该项目实现了一个用于优化Apache Spark代码的模型上下文协议(MCP)服务器和客户端。该系统通过客户端-服务器架构提供智能代码优化建议和性能分析。
graph TB
subgraph 输入
A[输入PySpark代码] --> |spark_code_input.py| B[run_client.py]
end
subgraph MCP客户端
B --> |异步HTTP| C[SparkMCPClient]
C --> |协议处理器| D[工具接口]
end
subgraph MCP服务器
E[run_server.py] --> F[SparkMCPServer]
F --> |工具注册| G[优化Spark代码]
F --> |工具注册| H[分析性能]
F --> |协议处理器| I[Claude AI集成]
end
subgraph 资源
I --> |代码分析| J[Claude AI模型]
J --> |优化| K[生成优化代码]
K --> |验证| L[PySpark运行时]
end
subgraph 输出
M[optimized_spark_code.py]
N[performance_analysis.md]
end
D --> |MCP请求| F
G --> |生成| M
H --> |生成| N
classDef 客户端 fill:#e1f5fe,stroke:#01579b
classDef 服务器 fill:#f3e5f5,stroke:#4a148c
classDef 资源 fill:#e8f5e9,stroke:#1b5e20
classDef 输出 fill:#fff3e0,stroke:#e65100
class A,B,C,D 客户端
class E,F,G,H,I 服务器
class J,K,L 资源
class M,N,O 输出
输入层
spark_code_input.py: 待优化的源PySpark代码run_client.py: 客户端启动和配置MCP客户端层
MCP服务器层
run_server.py: 服务器初始化资源层
输出层
optimized_spark_code.py: 优化后的代码performance_analysis.md: 详细分析此工作流程说明了:
本项目遵循模型上下文协议架构以标准化AI模型交互:
┌──────────────────┐ ┌──────────────────┐ ┌──────────────────┐
│ │ │ MCP服务器 │ │ 资源 │
│ MCP客户端 │ │ (SparkMCPServer)│ │ │
│ (SparkMCPClient) │ │ │ │ ┌──────────────┐ │
│ │ │ ┌─────────┐ │ │ │ Claude AI │ │
│ ┌─────────┐ │ │ │ 工具 │ │ <──> │ │ 模型 │ │
│ │ 工具 │ │ │ │ 注册 │ │ │ └──────────────┘ │
│ │ 接口 │ │ <──> │ └─────────┘ │ │ │
│ └─────────┘ │ │ ┌─────────┐ │ │ ┌──────────────┐ │
│ │ │ │ 协议 │ │ │ │ PySpark │ │
│ │ │ │ 处理器 │ │ │ │ 运行时 │ │
│ │ │ └─────────┘ │ │ └──────────────┘ │
└──────────────────┘ └──────────────────┘ └──────────────────┘
│ │ │
│ │ │
v v v
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ 可用 │ │ 注册 │ │ 外部 │
│ 工具 │ │ 工具 │ │ 资源 │
├──────────────┤ ├──────────────┤ ├──────────────┤
│optimize_code │ │optimize_code │ │ Claude API │
│analyze_perf │ │analyze_perf │ │ Spark 引擎 │
└──────────────┘ └──────────────┘ └──────────────┘
MCP客户端
MCP服务器
资源
sequenceDiagram
participant U as 用户
participant C as MCP客户端
participant S as MCP服务器
participant AI as Claude AI
participant P as PySpark运行时
U->>C: 提交Spark代码
C->>S: 发送优化请求
S->>AI: 分析代码
AI-->>S: 优化建议
S->>C: 返回优化后的代码
C->>P: 运行原始代码
C->>P: 运行优化后的代码
P-->>C: 执行结果
C->>C: 生成分析
C-->>U: 最终报告
代码提交
v1/input/spark_code_input.py优化过程
代码生成
v1/output/optimized_spark_code.py性能分析
结果生成
v1/output/performance_analysis.md中进行全面分析pip install -r requirements.txt
将要优化的Spark代码添加到input/spark_code_input.py
启动MCP服务器:
python v1/run_server.py
python v1/run_client.py
这将生成两个文件:
output/optimized_spark_example.py: 带有详细优化注释的优化Spark代码output/performance_analysis.md: 全面的性能分析python v1/run_optimized.py
这将:
ai-mcp/
├── input/
│ └── spark_code_input.py # 待优化的原始Spark代码
├── output/
│ ├── optimized_spark_example.py # 生成的优化代码
│ └── performance_analysis.md # 详细的性能对比
├── spark_mcp/
│ ├── client.py # MCP客户端实现
│ └── server.py # MCP服务器实现
├── run_client.py # 客户端脚本以优化代码
├── run_server.py # 服务器启动脚本
└── run_optimized.py # 脚本以运行并比较代码版本
模型上下文协议(MCP)为Spark代码优化提供了几个关键优势:
| 方面 | 直接调用Claude AI | MCP服务器 |
|---|---|---|
| 集成 | • 每个团队自定义集成<br>• 手动响应处理<br>• 重复实现 | • 预构建客户端库<br>• 自动化工作流<br>• 统一接口 |
| 基础设施 | • 没有内置验证<br>• 没有结果持久化<br>• 手动跟踪 | • 自动验证<br>• 结果持久化<br>• 版本控制 |
| 上下文 | • 基础代码建议<br>• 没有执行上下文<br>• 有限的优化范围 | • 上下文感知优化<br>• 完整执行历史<br>• 综合改进 |
| 验证 | • 需要手动测试<br>• 没有性能指标<br>• 不确定的结果 | • 自动测试<br>• 性能指标<br>• 验证结果 |
| 工作流 | • 临时过程<br>• 没有标准化<br>• 需要人工干预 | • 结构化过程<br>• 标准协议<br>• 自动化流水线 |
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | client = anthropic.Client(api_key)<br>response = client.messages.create(...) | • 复杂设置<br>• 自定义错误处理<br>• 紧耦合 |
| MCP | client = SparkMCPClient()<br>result = await client.optimize_spark_code(code) | • 简单接口<br>• 内置验证<br>• 松耦合 |
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | class SparkOptimizer:<br> def register_tool(self, name, func):<br> self.tools[name] = func | • 手动注册<br>• 没有验证<br>• 复杂维护 |
| MCP | @register_tool("optimize_spark_code")<br>async def optimize_spark_code(code: str): | • 自动发现<br>• 类型检查<br>• 易于扩展 |
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | def __init__(self):<br> self.claude = init_claude()<br> self.spark = init_spark() | • 手动编排<br>• 手动清理<br>• 错误多 |
| MCP | @requires_resources(["claude_ai", "spark"])<br>async def optimize_spark_code(code: str): | • 自动协调<br>• 生命周期管理<br>• 错误处理 |
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | {"type": "request",<br> "payload": {"code": code}} | • 自定义格式<br>• 手动验证<br>• 自定义调试 |
| MCP | {"method": "tools/call",<br> "params": {"name": "optimize_code"}} | • 标准格式<br>• 自动验证<br>• 易于调试 |
您也可以编程使用客户端:
from spark_mcp.client import SparkMCPClient
async def main():
# 连接到MCP服务器
client = SparkMCPClient()
await client.connect()
# 您的Spark代码以优化
spark_code = '''
# 您的PySpark代码在这里
'''
# 获取优化后的代码和性能分析
optimized_code = await client.optimize_spark_code(
code=spark_code,
optimization_level="advanced",
save_to_file=True # 保存到output/optimized_spark_example.py
)
# 分析性能差异
analysis = await client.analyze_performance(
original_code=spark_code,
optimized_code=optimized_code,
save_to_file=True # 保存到output/performance_analysis.md
)
# 运行并比较两个版本
# 您可以使用run_optimized.py脚本或实现自己的比较
await client.close()
# 分析性能
performance = await client.analyze_performance(spark_code, optimized_code)
await client.close()
仓库包括一个示例工作流程:
input/spark_code_input.py):# 创建数据帧并连接
emp_df = spark.createDataFrame(employees, ["id", "name", "age", "dept", "salary"])
dept_df = spark.createDataFrame(departments, ["dept", "location", "budget"])
# 连接并分析
result = emp_df.join(dept_df, "dept") \
.groupBy("dept", "location") \
.agg({"salary": "avg", "age": "avg", "id": "count"}) \
.orderBy("dept")
output/optimized_spark_example.py):# 性能优化版本,带有缓存和改进的配置
spark = SparkSession.builder \
.appName("EmployeeAnalysis") \
.config("spark.sql.shuffle.partitions", 200) \
.getOrCreate()
# 创建并缓存数据帧
emp_df = spark.createDataFrame(employees, ["id", "name", "age", "dept", "salary"]).cache()
dept_df = spark.createDataFrame(departments, ["dept", "location", "budget"]).cache()
# 优化连接和分析
result = emp_df.join(dept_df, "dept") \
.groupBy("dept", "location") \
.agg(
avg("salary").alias("avg_salary"),
avg("age").alias("avg_age"),
count("id").alias("employee_count")
) \
.orderBy("dept")
output/performance_analysis.md):## 执行结果对比
### 时间对比
- 原始代码:5.18秒
- 优化后的代码:0.65秒
- 性能提升:87.4%
### 优化细节
- 缓存频繁使用的数据帧
- 优化shuffle分区
- 改进列表达式
- 更好的内存管理
ai-mcp/
├── spark_mcp/
│ ├── __init__.py
│ ├── client.py # MCP客户端实现
│ └── server.py # MCP服务器实现
├── examples/
│ ├── optimize_code.py # 示例用法
│ └── optimized_spark_example.py # 生成的优化代码
├── requirements.txt
└── run_server.py # 服务器启动脚本
optimize_spark_code
analyze_performance
ANTHROPIC_API_KEY: 您的Anthropic API密钥,用于Claude AI系统实现了各种PySpark优化,包括:
欢迎提交问题和增强请求!
MIT许可