一个用于管理跨不同运行器(Flink、Spark、Dataflow 和 Direct)的 Apache Beam 数据管道的 Model Context Protocol (MCP) 服务器。
Apache Beam MCP 服务器提供了一个标准化的 API,用于管理跨不同运行器的 Apache Beam 数据管道。它设计用于:
# 克隆仓库
git clone https://github.com/yourusername/beam-mcp-server.git
cd beam-mcp-server
# 创建虚拟环境
python -m venv beam-mcp-venv
source beam-mcp-venv/bin/activate # 在 Windows 上:beam-mcp-venv\Scripts\activate
# 安装依赖
pip install -r requirements.txt
# 使用 Direct 运行器(无需外部依赖)
python main.py --debug --port 8888
# 使用 Flink 运行器(如果已安装 Flink)
CONFIG_PATH=config/flink_config.yaml python main.py --debug --port 8888
# 创建测试输入
echo "这是 Apache Beam WordCount 示例的测试文件" > /tmp/input.txt
# 提交作业
curl -X POST http://localhost:8888/api/v1/jobs \
-H "Content-Type: application/json" \
-d '{
"job_name": "test-wordcount",
"runner_type": "direct",
"job_type": "BATCH",
"code_path": "examples/pipelines/wordcount.py",
"pipeline_options": {
"input_file": "/tmp/input.txt",
"output_path": "/tmp/output"
}
}'
预构建的 Docker 镜像可以在 GitHub Container Registry 中找到:
# 拉取最新镜像
docker pull ghcr.io/yourusername/beam-mcp-server:latest
# 运行容器
docker run -p 8888:8888 \
-v $(pwd)/config:/app/config \
-e GCP_PROJECT_ID=your-gcp-project \
-e GCP_REGION=us-central1 \
ghcr.io/yourusername/beam-mcp-server:latest
# 构建镜像
./scripts/build_and_push_images.sh
# 构建并推送到注册表
./scripts/build_and_push_images.sh --registry your-registry --push --latest
用于本地开发,包含多个服务(Flink、Spark、Prometheus、Grafana):
docker-compose -f docker-compose.dev.yaml up -d
该仓库包含了 Kubernetes 清单文件,用于将 Beam MCP 服务器部署到 Kubernetes:
# 使用 kubectl 部署
kubectl apply -k kubernetes/
# 使用 Helm 部署
helm install beam-mcp ./helm/beam-mcp-server \
--namespace beam-mcp \
--create-namespace
详细的部署说明,请参阅 Kubernetes 部署指南。
Beam MCP 服务器实现了所有标准的 Model Context Protocol (MCP) 端点,提供了全面的框架来管理 AI 控制的数据管道:
/tools 端点管理用于管道处理的 AI 代理和模型:
# 注册情感分析工具
curl -X POST "http://localhost:8888/api/v1/tools/" \
-H "Content-Type: application/json" \
-d '{
"name": "sentiment-analyzer",
"description": "分析文本中的情感",
"type": "transformation",
"parameters": {
"text_column": {
"type": "string",
"description": "包含要分析的文本的列"
}
}
}'
/resources 端点管理数据集和其他管道资源:
# 注册数据集
curl -X POST "http://localhost:8888/api/v1/resources/" \
-H "Content-Type: application/json" \
-d '{
"name": "客户交易",
"description": "每日客户交易数据",
"resource_type": "dataset",
"location": "gs://analytics-data/transactions/*.csv"
}'
/contexts 端点定义管道执行环境:
# 创建 Dataflow 执行环境
curl -X POST "http://localhost:8888/api/v1/contexts/" \
-H "Content-Type: application/json" \
-d '{
"name": "Dataflow 生产",
"description": "生产 Dataflow 环境",
"context_type": "dataflow",
"parameters": {
"region": "us-central1",
"project": "beam-analytics-prod"
}
}'
这些 MCP 标准端点与 Beam 的核心功能无缝集成,提供了一整套解决方案来管理数据管道。详细示例和用例,请参阅 MCP 协议合规性。
import requests
# 获取可用运行器
headers = {"MCP-Session-ID": "my-session-123"}
runners = requests.get("http://localhost:8888/api/v1/runners", headers=headers).json()
# 创建作业
job = requests.post(
"http://localhost:8888/api/v1/jobs",
headers=headers,
json={
"job_name": "wordcount-example",
"runner_type": "flink",
"job_type": "BATCH",
"code_path": "examples/pipelines/wordcount.py",
"pipeline_options": {
"parallelism": 2,
"input_file": "/tmp/input.txt",
"output_path": "/tmp/output"
}
}
).json()
# 监控作业状态
job_id = job["data"]["job_id"]
status = requests.get(f"http://localhost:8888/api/v1/jobs/{job_id}", headers=headers).json()
该仓库包含一个 GitHub Actions 工作流,用于持续集成和部署:
Beam MCP 服务器内置了对监控和可观测性的支持:
/metrics 端点暴露指标/health 端点提供健康检查我们欢迎贡献!详情请参阅我们的 贡献指南。
要运行测试:
# 运行回归测试
./scripts/run_regression_tests.sh
此项目根据 Apache 许可证 2.0 授权。
MCP(Model Context Protocol)的实现分为几个阶段:
当构建客户端与 MCP 服务器交互时,必须遵循 Model Context Protocol。详情请参阅 MCP 协议合规性。