返回市场
流处理-mcp

流处理-mcp

作者:Cledar8 星标更新:2025-09-04

项目介绍

flink-mcp — Flink MCP 服务器

该项目提供了一个连接到 Apache Flink SQL 网关的 MCP 服务器。

先决条件

  • 运行中的 Apache Flink 集群和 SQL 网关

    • 启动集群:./bin/start-cluster.sh
    • 启动网关:./bin/sql-gateway.sh start -Dsql-gateway.endpoint.rest.address=localhost
    • 验证:curl http://localhost:8083/v3/info
  • 配置环境:

    • 设置 SQL_GATEWAY_API_BASE_URL(默认 http://localhost:8083)。您可以在仓库根目录使用一个 .env 文件。

运行

通过控制台脚本安装并运行:

pip install -e .
flink-mcp

MCP 客户端应通过命令 flink-mcp 在标准 I/O 上启动服务器。

确保在您的环境中或 .env 文件中设置了 SQL_GATEWAY_API_BASE_URL

工具 (v0.2.5)

  • flink_info (资源):从 /v3/info 返回集群信息。
  • open_new_session(properties?: dict) -> { sessionHandle, ... }
  • get_config(sessionHandle: str):返回会话配置。
  • configure_session(sessionHandle: str, statement: str):应用会话范围内的 DDL/配置(CREATE/USE/SET/RESET/LOAD/UNLOAD/ADD JAR)。
  • run_query_collect_and_stop(sessionHandle: str, query: str, max_rows: int=5, max_seconds: float=15.0):执行查询,在 T 秒内获取最多 N 行数据,如果存在 jobID 则停止作业;关闭操作。
  • run_query_stream_start(sessionHandle: str, query: str):执行流式查询并返回 { jobID, operationHandle };作业保持运行状态。
  • fetch_result_page(sessionHandle: str, operationHandle: str, token: int):获取单页结果;返回 { page, nextToken, isEnd }
  • cancel_job(sessionHandle: str, jobId: str):发出 STOP JOB '<jobId>' 命令,等待直到 DESCRIBE JOB 状态不再是 RUNNING;返回 { jobID, status, jobGone, jobStatus }

注意事项

  • 工具是无状态的;客户端需明确管理并传递会话/操作句柄。

  • run_query_stream_start 返回 jobIDoperationHandle;使用 fetch_result_page 流式传输结果。

  • cancel_job 发出 STOP 并使用 DESCRIBE JOB 等待;在适当情况下内部调用 close_operation

  • 端点针对 SQL 网关 v3 样式的路径。