返回市场
光束-MCP服务器

光束-MCP服务器

作者:souravch4 星标更新:2025-03-22

项目介绍

Apache Beam MCP 服务器

一个用于管理跨不同运行器(Flink、Spark、Dataflow 和 Direct)的 Apache Beam 数据管道的 Model Context Protocol (MCP) 服务器。

Python 3.9+ MCP 版本 Apache Beam Docker Kubernetes

这是什么?

Apache Beam MCP 服务器提供了一个标准化的 API,用于管理跨不同运行器的 Apache Beam 数据管道。它设计用于:

  • 数据工程师:使用一致的 API 管理管道,无论运行器如何
  • AI/LLM 开发者:通过 MCP 标准启用 AI 控制的数据管道
  • DevOps 团队:简化管道操作和监控

主要特性

  • 多运行器支持:一个 API 支持 Flink、Spark、Dataflow 和 Direct 运行器
  • 符合 MCP:遵循 Model Context Protocol 以实现 AI 集成
  • 管道管理:创建、监控和控制数据管道
  • 易于扩展:添加新的运行器或自定义功能
  • 生产就绪:包括 Docker/Kubernetes 部署、监控和扩展

快速开始

安装

# 克隆仓库
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 支持

使用预构建镜像

预构建的 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

Docker Compose

用于本地开发,包含多个服务(Flink、Spark、Prometheus、Grafana):

docker-compose -f docker-compose.dev.yaml up -d

Kubernetes 部署

该仓库包含了 Kubernetes 清单文件,用于将 Beam MCP 服务器部署到 Kubernetes:

# 使用 kubectl 部署
kubectl apply -k kubernetes/

# 使用 Helm 部署
helm install beam-mcp ./helm/beam-mcp-server \
  --namespace beam-mcp \
  --create-namespace

详细的部署说明,请参阅 Kubernetes 部署指南

MCP 标准端点

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 协议合规性

文档

Python 客户端示例

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()

CI/CD 管道

该仓库包含一个 GitHub Actions 工作流,用于持续集成和部署:

  • CI:在每个拉取请求上运行测试、代码检查和类型检查
  • CD:在每次向主分支推送时构建并推送 Docker 镜像
  • 部署:自动部署到开发和生产环境

监控和可观测性

Beam MCP 服务器内置了对监控和可观测性的支持:

  • Prometheus 指标:在 /metrics 端点暴露指标
  • Grafana 仪表盘:预配置的监控仪表盘
  • 健康检查:在 /health 端点提供健康检查
  • 日志记录:结构化的 JSON 日志记录,便于与日志聚合系统集成

贡献

我们欢迎贡献!详情请参阅我们的 贡献指南

要运行测试:

# 运行回归测试
./scripts/run_regression_tests.sh

许可证

此项目根据 Apache 许可证 2.0 授权。

MCP 实现状态

MCP(Model Context Protocol)的实现分为几个阶段:

第一阶段:核心连接生命周期(已完成)

  • ✅ 连接初始化
  • ✅ 连接状态管理
  • ✅ 基础能力协商
  • ✅ 带有 SSE 的 HTTP 传输
  • ✅ JSON-RPC 消息处理
  • ✅ 错误处理

第二阶段:完整的功能协商(已完成)

  • ✅ 增强的功能兼容性检查
  • ✅ 功能的语义版本兼容性
  • ✅ 功能的支持级别(必需、首选、可选、实验)
  • ✅ 功能属性验证
  • ✅ 基于功能的 API 端点控制
  • ✅ 与 FastAPI 的功能路由器集成

第三阶段:高级消息处理(已完成)

  • ✅ 结构化消息类型
  • ✅ 消息验证
  • ✅ 改进的错误处理
  • ✅ 批量消息处理

第四阶段:生产优化(待完成)

  • ⬜ 性能优化
  • ⬜ 监控和指标
  • ⬜ 高级安全功能
  • ⬜ 高可用性支持

当构建客户端与 MCP 服务器交互时,必须遵循 Model Context Protocol。详情请参阅 MCP 协议合规性