Apache Flink MCP 服务器
一个针对Apache Flink的模型上下文协议(MCP)服务器实现,它使得人工智能助手和大型语言模型能够通过自然语言接口与Flink集群进行交互。该服务器提供了全面的工具,用于监控、管理和分析Apache Flink流处理应用程序。
概述
Apache Flink MCP服务器通过提供标准化的MCP接口,弥合了人工智能助手与Apache Flink集群之间的鸿沟。它使用户能够通过对话式人工智能执行复杂的Flink操作,从而让流处理管理更加便捷且直观。
特点/特性
🎯 核心能力
- 集群监控获取实时集群信息,包括作业、槽位和TaskManagers
- 作业管理列出、监控和分析Flink作业的详细信息和指标
- 异常追踪检索并分析作业异常以进行调试
- 资源管理监控TaskManager资源和JAR文件的部署
- 指标收集访问全面的职位和集群指标
🔧 可用工具:
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– (发送电子邮件通知)
______________________________________________________________________
🚀 好处/益处
- 自然语言接口使用对话式AI与Flink进行交互
- 实时监控即时了解集群和作业状态
- 调试支持轻松访问异常日志和指标
- 资源优化监控所有TaskManager的资源使用情况
- 开发人员生产力减少在Flink Web UI中导航所花费的时间
安装
先决条件
- Apache Flink 集群(正在运行且可访问)
- Python 3.8 或更高版本
- 与MCP兼容的客户端(如Claude Desktop、Continue等)
客户端配置
Continue.dev(可译为“继续开发平台”或根据具体语境简化为“继续.dev”)
添加到您的继续配置(Flink-mcp-server.yaml):
name: Sample MCP
version: 0.0.1
schema: v1
mcpServers:
- name: Flink MCP Server
type: streamable-http
url: http://127.0.0.1:9090/mcp/ 使用示例
基本集群监控
Human: What's the status of my Flink cluster?
AI: I'll check your Flink cluster status for you.
[Uses get_cluster_info tool to fetch cluster overview]工作分析
Human: Show me all running jobs and their performance metrics
AI: Let me get the current jobs and their metrics.
[Uses list_jobs and get_job_metrics tools]故障排除
Human: My job with ID abc123 is failing. Can you help me debug it?
AI: I'll check the job details and any exceptions for job abc123.
[Uses get_job_details and get_job_exceptions tools]资源管理
Human: How are my TaskManager resources being utilized?
AI: Let me check your TaskManager status and resource allocation.
[Uses list_taskmanagers tool]API 参考文档
可用的MCP工具
get_cluster_info
描述获取Flink集群的概览,包括作业、槽位和TaskManager。 参数无 回报集群概览及资源信息
list_jobs
描述列出所有当前和最近的Flink作业及其状态。 参数无 退货包含状态、开始时间和持续时间的工作列表
get_job_details
描述获取特定Flink作业的详细信息。 参数:
job_id(字符串,必填):Flink 作业的唯一标识符
list_taskmanagers
描述列出集群中所有已注册的 TaskManager。 参数无 退货包含资源信息的TaskManager列表
get_job_exceptions
描述获取在指定作业中发生的异常。 参数:
job_id(字符串,必填):Flink 作业的唯一标识符
list_jar_files
描述列出Flink集群中所有已上传的JAR文件。 参数无 回报可用的JAR文件列表
get_job_metrics
描述获取运行中Flink作业的选定有用指标。 参数:
job_id(字符串,必填):Flink 作业的唯一标识符
______________________________________________________________________
send_mail
描述发送电子邮件通知,例如来自Flink MCP服务器的警报、状态更新或报告。
故障排除
常见问题
连接失败
- 验证Flink集群是否正在运行且可访问
- 确保网络连接到 Flink JobManager
权限错误
- 验证Flink REST API是否已启用
- 检查您的Flink设置是否需要认证
调试模式
启用详细日志记录:
export LOG_LEVEL=DEBUG
python mcp_server.py贡献
我们欢迎投稿!请按照以下步骤操作:
- 克隆该仓库
- 创建一个特性分支:
git checkout -b feature-name - 进行你的更改并添加测试
- 提交您的更改:
git commit -m "Add feature" - 推送到你的叉(分支):
git push origin feature-name - 创建一个拉取请求
开发指南
- 遵循PEP 8风格指南
- 根据需要更新文档
- 确保向后兼容性
许可证
这个项目遵循以下许可协议 麻省理工学院许可证(MIT License).
相关项目
- 模型上下文协议 - MCP规范
- Apache Flink - Apache Flink 流处理框架
- MCP Kafka - Confluent/Kafka 的 MCP 服务器
- MCP 容器 - 用于Container的MCP服务器
支持
- 问题:
- 文档: 项目维基
- 讨论: 项目讨论
致谢
- Apache Flink社区,卓越的流处理框架
- 标准化接口的模型上下文协议团队
- 该项目的贡献者和用户
______________________________________________________________________
注此MCP服务器默认仅提供对Flink集群信息的只读访问权限。对于写操作,可能需要额外的配置和安全考虑。
