Spark MCP(模型上下文协议)优化器
该项目实现了一个模型上下文协议(MCP)服务器和客户端,用于优化Apache Spark代码。该系统通过客户端-服务器架构提供智能代码优化建议和性能分析。
运作原理
代码优化工作流程
graph TB
subgraph Input
A[Input PySpark Code] --> |spark_code_input.py| B[run_client.py]
end
subgraph MCP Client
B --> |Async HTTP| C[SparkMCPClient]
C --> |Protocol Handler| D[Tools Interface]
end
subgraph MCP Server
E[run_server.py] --> F[SparkMCPServer]
F --> |Tool Registry| G[optimize_spark_code]
F --> |Tool Registry| H[analyze_performance]
F --> |Protocol Handler| I[Claude AI Integration]
end
subgraph Resources
I --> |Code Analysis| J[Claude AI Model]
J --> |Optimization| K[Optimized Code Generation]
K --> |Validation| L[PySpark Runtime]
end
subgraph Output
M[optimized_spark_code.py]
N[performance_analysis.md]
end
D --> |MCP Request| F
G --> |Generate| M
H --> |Generate| N
classDef client fill:#e1f5fe,stroke:#01579b
classDef server fill:#f3e5f5,stroke:#4a148c
classDef resource fill:#e8f5e9,stroke:#1b5e20
classDef output fill:#fff3e0,stroke:#e65100
class A,B,C,D client
class E,F,G,H,I server
class J,K,L resource
class M,N,O output组件详细信息
- 输入层
- spark_code_input.py:用于优化的PySpark源代码 - run_client.py:客户端启动和配置
- MCP客户端层
- 工具接口:符合协议的工具调用
- MCP服务器层
- run_server.py:服务器初始化 - 工具注册表:优化和分析工具 - 协议处理程序:MCP请求/响应管理
- 资源层
- Claude AI:代码分析和优化 - PySpark运行时:代码执行和验证
- 输出层
- optimized_spark_code.py:优化代码 - performance_analysis.md:详细分析
此工作流程说明:
- 输入PySpark代码提交
- MCP协议处理和路由
- Claude AI分析与优化
- 代码转换和验证
- 性能分析和报告
建筑
该项目遵循标准化AI模型交互的模型上下文协议架构:
┌──────────────────┐ ┌──────────────────┐ ┌──────────────────┐
│ │ │ MCP Server │ │ Resources │
│ MCP Client │ │ (SparkMCPServer)│ │ │
│ (SparkMCPClient) │ │ │ │ ┌──────────────┐ │
│ │ │ ┌─────────┐ │ │ │ Claude AI │ │
│ ┌─────────┐ │ │ │ Tools │ │ │ │ Model │ │
│ │ Tools │ │ │ │Registry │ │ │ └──────────────┘ │
│ │Interface│ │ │ └─────────┘ │ │ │
│ └─────────┘ │ │ ┌─────────┐ │ │ ┌──────────────┐ │
│ │ │ │Protocol │ │ │ │ PySpark │ │
│ │ │ │Handler │ │ │ │ Runtime │ │
│ │ │ └─────────┘ │ │ └──────────────┘ │
└──────────────────┘ └──────────────────┘ └──────────────────┘
│ │ │
│ │ │
v v v
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ Available │ │ Registered │ │ External │
│ Tools │ │ Tools │ │ Resources │
├──────────────┤ ├──────────────┤ ├──────────────┤
│optimize_code │ │optimize_code │ │ Claude API │
│analyze_perf │ │analyze_perf │ │ Spark Engine │
└──────────────┘ └──────────────┘ └──────────────┘组件
- MCP客户端
- 为代码优化提供工具界面 - 处理与服务器的异步通信 - 管理文件I/O以生成代码
- MCP服务器
- 实现MCP协议处理程序 - 管理工具注册表和执行 - 客户和资源之间的协调
- 资源
- Claude AI:提供代码优化智能 - PySpark运行时:执行和验证优化
协议流
- 客户端通过MCP协议发送优化请求
- 服务器验证请求并调用适当的工具
- 工具利用Claude AI进行优化
- 通过MCP响应返回优化代码
- 客户端保存并验证优化后的代码
端到端功能
sequenceDiagram
participant U as User
participant C as MCP Client
participant S as MCP Server
participant AI as Claude AI
participant P as PySpark Runtime
U->>C: Submit Spark Code
C->>S: Send Optimization Request
S->>AI: Analyze Code
AI-->>S: Optimization Suggestions
S->>C: Return Optimized Code
C->>P: Run Original Code
C->>P: Run Optimized Code
P-->>C: Execution Results
C->>C: Generate Analysis
C-->>U: Final Report- 代码提交
- 用户将PySpark代码放入 v1/input/spark_code_input.py - MCP客户端读取代码
- 优化过程
- MCP客户端通过标准化协议连接到服务器 - 服务器将代码转发给Claude AI进行分析 - AI建议基于最佳实践进行优化 - 服务器验证并处理建议
- 代码生成
- 优化后的代码保存到 v1/output/optimized_spark_code.py - 包括解释优化的详细注释 - 在提高性能的同时保持原始代码结构
- 性能分析
- 两个版本都在PySpark运行时执行 - 执行时间比较 - 验证结果的正确性 - 收集和分析的指标
- 结果生成
- 综合分析 v1/output/performance_analysis.md - 并排执行比较 - 性能改进统计 - 优化解释和基本原理
用法
需求
- Python 3.8+
- PySpark 3.2.0+
- Anthropic API密钥(用于Claude AI)
安装
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 # Original Spark code to optimize
├── output/
│ ├── optimized_spark_example.py # Generated optimized code
│ └── performance_analysis.md # Detailed performance comparison
├── spark_mcp/
│ ├── client.py # MCP client implementation
│ └── server.py # MCP server implementation
├── run_client.py # Client script to optimize code
├── run_server.py # Server startup script
└── run_optimized.py # Script to run and compare code versions为什么选择MCP?
模型上下文协议(MCP)为Spark代码优化提供了几个关键优势:
直接克劳德AI呼叫与MCP服务器
| 方面 | 直接Claude AI呼叫 | MCP服务器 |
|---|---|---|
| 整合 | •每个团队的定制集成 |
•手动响应处理 •重复实现|•预构建的客户端库 •自动化工作流程 •统一接口| | 基础设施 |•无内置验证 •无结果持久性 •手动跟踪|•自动验证 •结果持久性 •版本控制| | 上下文 |•基本代码建议 •无执行上下文 •优化范围有限|•上下文感知优化 •完整的执行历史 •全面改进| | 验证 |•需要手动测试 •无绩效指标 •不确定的结果|•自动化测试 •性能指标 •验证结果| | 工作流程 |•临时流程 •没有标准化 •需要人工干预|•结构化流程 •标准协议 •自动化管道|
主要区别:
1.人工智能集成
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | client = anthropic.Client(api_key) | |
response = client.messages.create(...) | •复杂的设置 |
•自定义错误处理 •紧密耦合| |MCP| client = SparkMCPClient() result = await client.optimize_spark_code(code) |•界面简单 •内置验证 •松耦合|
2.工具管理
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | class SparkOptimizer: |
def register_tool(self, name, func): self.tools[name] = func |•手动注册 •无验证 •复杂的维护| |MCP| @register_tool("optimize_spark_code") async def optimize_spark_code(code: str): |•自动发现 •类型检查 •易于扩展|
3.资源管理
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | def __init__(self): |
self.claude = init_claude() self.spark = init_spark() |•手动编排 •手动清理 •容易出错| |MCP| @requires_resources(["claude_ai", "spark"]) async def optimize_spark_code(code: str): |•自动协调 •生命周期管理 •错误处理|
4.通信协议
| 方法 | 代码示例 | 优点 |
|---|---|---|
| 传统 | {"type": "request", | |
"payload": {"code": code}} | •自定义格式 |
•手动验证 •自定义调试| |MCP| {"method": "tools/call", "params": {"name": "optimize_code"}} |•标准格式 •自动验证 •易于调试|
特性
- 智能代码优化:利用Claude AI分析和优化PySpark代码
- 性能分析:提供原始代码和优化代码之间性能差异的详细分析
- MCP架构:实现标准化AI模型交互的模型上下文协议
- 易于集成:用于代码优化请求的简单客户端界面
- 代码生成:自动将优化后的代码保存到单独的文件中
高级用法
您还可以通过编程方式使用客户端:
from spark_mcp.client import SparkMCPClient
async def main():
# Connect to the MCP server
client = SparkMCPClient()
await client.connect()
# Your Spark code to optimize
spark_code = '''
# Your PySpark code here
'''
# Get optimized code with performance analysis
optimized_code = await client.optimize_spark_code(
code=spark_code,
optimization_level="advanced",
save_to_file=True # Save to output/optimized_spark_example.py
)
# Analyze performance differences
analysis = await client.analyze_performance(
original_code=spark_code,
optimized_code=optimized_code,
save_to_file=True # Save to output/performance_analysis.md
)
# Run both versions and compare
# You can use the run_optimized.py script or implement your own comparison
await client.close()
# Analyze performance
performance = await client.analyze_performance(spark_code, optimized_code)
await client.close()输入和输出示例
存储库包括一个示例工作流:
- 输入代码 (
input/spark_code_input.py):
# Create DataFrames and join
emp_df = spark.createDataFrame(employees, ["id", "name", "age", "dept", "salary"])
dept_df = spark.createDataFrame(departments, ["dept", "location", "budget"])
# Join and analyze
result = emp_df.join(dept_df, "dept") \
.groupBy("dept", "location") \
.agg({"salary": "avg", "age": "avg", "id": "count"}) \
.orderBy("dept")- 优化代码 (
output/optimized_spark_example.py):
# Performance-optimized version with caching and improved configurations
spark = SparkSession.builder \
.appName("EmployeeAnalysis") \
.config("spark.sql.shuffle.partitions", 200) \
.getOrCreate()
# Create and cache DataFrames
emp_df = spark.createDataFrame(employees, ["id", "name", "age", "dept", "salary"]).cache()
dept_df = spark.createDataFrame(departments, ["dept", "location", "budget"]).cache()
# Optimized join and analysis
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):
## Execution Results Comparison
### Timing Comparison
- Original Code: 5.18 seconds
- Optimized Code: 0.65 seconds
- Performance Improvement: 87.4%
### Optimization Details
- Caching frequently used DataFrames
- Optimized shuffle partitions
- Improved column expressions
- Better memory management项目结构
ai-mcp/
├── spark_mcp/
│ ├── __init__.py
│ ├── client.py # MCP client implementation
│ └── server.py # MCP server implementation
├── examples/
│ ├── optimize_code.py # Example usage
│ └── optimized_spark_example.py # Generated optimized code
├── requirements.txt
└── run_server.py # Server startup script可用工具
- 优化公园代码
- 优化PySpark代码以获得更好的性能 - 支持基本和高级优化级别 - 自动将优化后的代码保存到examples/optimized_spark_example.py
- 分析性能
- 分析原始代码和优化代码之间的性能差异 - 提供以下见解: - 性能改进 - 资源利用率 - 可扩展性考虑因素 - 潜在的权衡
环境变量
ANTHROPIC_API_KEY:Claude AI的Anthropic API密钥
优化示例
该系统实现了各种PySpark优化,包括:
- 广播连接用于小型-大型表连接
- 高效使用窗口功能
- 战略性数据缓存
- 查询计划优化
- 以性能为导向的操作排序
贡献
请随时提交问题和增强请求!
许可证
MIT许可证
