该项目提供了一个连接到 Apache Flink SQL 网关的 MCP 服务器。
运行中的 Apache Flink 集群和 SQL 网关
./bin/start-cluster.sh./bin/sql-gateway.sh start -Dsql-gateway.endpoint.rest.address=localhostcurl 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。
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 返回 jobID 和 operationHandle;使用 fetch_result_page 流式传输结果。
cancel_job 发出 STOP 并使用 DESCRIBE JOB 等待;在适当情况下内部调用 close_operation。
端点针对 SQL 网关 v3 样式的路径。