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 位于repo根目录的文件。
跑
通过控制台脚本安装并运行:
pip install -e .
flink-mcpMCP客户端应使用以下命令通过stdio启动服务器: flink-mcp.
确保 SQL_GATEWAY_API_BASE_URL 在您的环境中设置或 .env.
工具(v0.2.5)
flink_info(资源):从返回集群信息/v3/info.open_new_session(properties?: dict)->{ sessionHandle, ... }.get_config(sessionHandle: str):返回会话配置。configure_session(sessionHandle: str, statement: str):应用会话范围的DDL/config(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):execute,在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 '',等待DESCRIBE作业状态为非运行;回报{ jobID, status, jobGone, jobStatus }.
备注
- 工具是无状态的;客户端显式地管理和传递会话/操作句柄。
run_query_stream_start返回两者jobID和operationHandle;使用fetch_result_page以流式传输结果。
cancel_job使用DESCRIBE作业发出STOP并等待;close_operation在适当的情况下在内部调用。
- 端点以SQL网关v3样式的路径为目标。
