这是一个针对Apache Flink的模型上下文协议(MCP)服务器实现,它使AI助手和大型语言模型能够通过自然语言接口与Flink集群进行交互。此服务器提供了全面的工具,用于监控、管理和分析Apache Flink流处理应用程序。

Apache Flink MCP Server通过提供标准化的MCP接口,弥合了AI助手与Apache Flink集群之间的差距。它允许用户通过对话式AI执行复杂的Flink操作,使得流处理管理更加便捷和直观。
initialize_flink_connection – 连接到Flink REST APIget_connection_status – 检查连接状态get_cluster_info – 获取Flink集群概览list_jobs – 列出所有Flink作业及其状态get_job_details – 通过ID获取详细作业信息get_job_exceptions – 获取作业级别的异常get_job_metrics – 获取作业指标list_taskmanagers – 列出带有资源的TaskManagerslist_jar_files – 列出已上传的JAR文件send mail – (发送电子邮件通知)添加到您的Continue配置文件(Flink-mcp-server.yaml):
name: 示例MCP
version: 0.0.1
schema: v1
mcpServers:
- name: Flink MCP Server
type: streamable-http
url: http://127.0.0.1:9090/mcp/
Human: 我的Flink集群状态如何?
AI: 我会为您检查Flink集群的状态。
[使用get_cluster_info工具获取集群概览]
Human: 显示所有正在运行的作业及其性能指标
AI: 让我获取当前作业及其指标。
[使用list_jobs和get_job_metrics工具]
Human: 我的作业ID abc123失败了。你能帮我调试吗?
AI: 我会检查作业abc123的详细信息和任何异常。
[使用get_job_details和get_job_exceptions工具]
Human: 我的TaskManager资源是如何被利用的?
AI: 让我检查您的TaskManager状态和资源分配。
[使用list_taskmanagers工具]
get_cluster_info描述:获取Flink集群概览,包括作业、插槽和TaskManagers。 参数:无 返回值:包含资源信息的集群概览
list_jobs描述:列出所有当前和最近的Flink作业及其状态。 参数:无 返回值:带有状态、开始时间和持续时间的作业列表
get_job_details描述:获取特定Flink作业的详细信息。 参数:
job_id (字符串,必需):Flink作业的唯一标识符list_taskmanagers描述:列出集群中注册的所有TaskManagers。 参数:无 返回值:带有资源信息的TaskManagers列表
get_job_exceptions描述:获取指定作业中发生的异常。 参数:
job_id (字符串,必需):Flink作业的唯一标识符list_jar_files描述:列出Flink集群中上传的所有JAR文件。 参数:无 返回值:可用JAR文件列表
get_job_metrics描述:获取正在运行的Flink作业的选定有用指标。 参数:
job_id (字符串,必需):Flink作业的唯一标识符send_mail描述:从Flink MCP服务器发送电子邮件通知,例如警报、状态更新或报告。
连接失败
权限错误
启用详细日志记录:
export LOG_LEVEL=DEBUG
python mcp_server.py
我们欢迎贡献!请遵循以下步骤:
git checkout -b feature-namegit commit -m "添加功能"git push origin feature-name本项目根据MIT许可证授权。
注意:此MCP服务器默认提供对Flink集群信息的只读访问。对于写入操作,可能需要额外的配置和安全考虑。