ETL代理X MCP-智能奖章管道
基于LangGraph和copula的人工智能驱动的多代理ETL编排系统Claude Sonnet 4.5
通过Bronze转换原始数据→ 银→ 具有智能代理的黄金层,可以自动规划、编码、审查和执行转换。每一层都由具有上下文感知代码生成功能的GitHub PR支持。
______________________________________________________________________
🎯 快速开始
先决条件
- Python 3.12+
- 配备无服务器SQL仓库的copula工作空间
- 带有PAT令牌的GitHub存储库
- Rancher个人访问令牌
安装
# Clone and install
git clone
cd ETLAgenticXMCP
# Configure environment
cp .env.example .env
# Edit .env with your credentials
# Install dependencies
pip install -r requirements.txt运行管道
from graph_workflow import MEDALLION_PIPELINE
# Execute full Bronze → Silver → Gold transformation
result = MEDALLION_PIPELINE.invoke({
"user_query": "Clean weather data and create daily city aggregations",
"source_table": "samples.accuweather.forecast_daily_calendar_imperial"
})或者使用MCP服务器进行VS Code集成:
python server.py______________________________________________________________________
🏗️ 架构概述
flowchart TB
subgraph user_interface["User Interface"]
vscode_client["VS Code + Copilot"]
prompt_entry["User Prompt Entry"]
end
subgraph mcp_layer["MCP Protocol & Orchestration"]
mcp_server["FastMCP Server (Python)"]
langgraph_orch["LangGraph Workflow"]
etl_state["ETL State Management"]
end
subgraph agent_layer["Agent Layer (Databricks Claude 4.5)"]
planner_agent["PlannerAgent (Claude)"]
codegen_agent["CodeGenAgent (Claude)"]
reviewer_agent["ReviewerAgent (Claude)"]
pr_creator_agent["PRCreatorAgent (GitHub API)"]
executor_agent["ExecutorAgent (Databricks SDK)"]
context_agent["ContextEnrichmentAgent"]
summary_agent["SummaryAgent (Claude)"]
end
subgraph tools_integration["Python Integration & External APIs"]
databricks_tools["Databricks Tools (SDK, SQL)"]
github_tools["GitHub Tools (REST, GitPython)"]
medallion_tools["Medallion Utils (SQL Generation)"]
context_tools["Context Tools (Schema, Analysis)"]
end
subgraph databricks_infra["Databricks Workspace"]
workspace["Workspace / Serverless Infra"]
foundation_model_api["Foundation Model API
Claude Sonnet 4.5 Endpoint"]
sql_warehouse["SQL Warehouse (Serverless)"]
delta_tables["Delta Lake Tables (Bronze/Silver/Gold)"]
unity_catalog["Unity Catalog"]
end
subgraph github_vcs["GitHub Version Control"]
github_repo["Repository"]
github_pr["Pull Requests"]
github_ci["Security & Approval"]
end
subgraph storage["Medallion Data Layers"]
bronze_table["Bronze Table (Raw)"]
silver_table["Silver Table (Clean)"]
gold_table["Gold Table (Aggregate)"]
end
subgraph config["Configuration"]
env_vars[".env File"]
rules_txt["rules.txt"]
end
%% USER TO INTERFACE
vscode_client --> prompt_entry
prompt_entry --> mcp_server
%% CONFIGURATION
env_vars --> mcp_server
rules_txt --> mcp_server
%% MCP to ORCHESTRATION
mcp_server --> langgraph_orch
mcp_server --> etl_state
etl_state langgraph_orch
%% Orchestrator calls agents in sequence
langgraph_orch --> planner_agent
planner_agent --> codegen_agent
codegen_agent --> reviewer_agent
reviewer_agent -.-> pr_creator_agent
pr_creator_agent --> executor_agent
executor_agent --> context_agent
context_agent -.-> planner_agent
context_agent --> summary_agent
summary_agent --> mcp_server
%% AGENTS to TOOLS
planner_agent --> databricks_tools
codegen_agent --> medallion_tools
reviewer_agent --> databricks_tools
pr_creator_agent --> github_tools
executor_agent --> databricks_tools
context_agent --> context_tools
%% DATABRICKS: Model API and SQL Warehouse
planner_agent --> foundation_model_api
codegen_agent --> foundation_model_api
reviewer_agent --> foundation_model_api
summary_agent --> foundation_model_api
executor_agent --> sql_warehouse
databricks_tools --> sql_warehouse
medallion_tools --> sql_warehouse
context_tools --> sql_warehouse
sql_warehouse --> delta_tables
delta_tables --> unity_catalog
%% DATA FLOW Bronze → Silver → Gold
bronze_table --> silver_table
silver_table --> gold_table
%% Delta tables backed by physical tables
bronze_table --> delta_tables
silver_table --> delta_tables
gold_table --> delta_tables
%% CONTEXT ENRICHMENT
bronze_table -.-> context_agent
silver_table -.-> context_agent
gold_table -.-> context_agent
%% PR Creator → GitHub
pr_creator_agent --> github_pr
github_pr --> github_repo
github_repo --> github_ci
github_pr -.-> executor_agent
%% Final summary returns to user
summary_agent --> vscode_client______________________________________________________________________
📊 管道工作流
单层转换
对于每个奖章层(青铜、银、金):
1. 🎯 PLANNER
├─ Analyzes user query & transformation rules
├─ Reviews context from previous layer (if available)
└─ Creates detailed transformation plan
2. 💻 CODE GENERATOR
├─ Generates PySpark/SQL code
├─ Includes error handling & data quality checks
└─ Follows medallion best practices
3. 🔍 REVIEWER
├─ Validates syntax & logic
├─ Checks against rules.txt compliance
├─ Generates quality score
└─ Approves or requests revision
4. 🔗 PR CREATOR (if approved)
├─ Creates GitHub PR with generated code
├─ Adds documentation & test plans
└─ Waits for approval & merge
5. ⚡ EXECUTOR (on PR merge)
├─ Detects PR merge event
├─ Executes transformation on Databricks
└─ Monitors execution & collects metrics
6. 📈 CONTEXT ENRICHMENT
├─ Queries output table schema
├─ Analyzes data quality metrics
├─ Creates summary for next layer
└─ Loops back to PLANNER for next layer
7. 📝 SUMMARY GENERATOR (all layers complete)
├─ Aggregates execution metrics
├─ Generates executive report
└─ Returns to user多层智能
每一层接收 真实、新鲜的背景 从上一层开始:
Bronze Layer
↓ [Context: Raw data schema, row counts, data types]
Silver Layer (uses Bronze context)
↓ [Context: Cleaned data schema, quality metrics, duplicates removed]
Gold Layer (uses Silver context)
↓
Executive Summary______________________________________________________________________
🔧 核心组件
代理商(agents/)
| 代理 | 角色 |
|---|---|
| 规划师代理人 | 分析请求,创建转换策略,从先前的上下文中丰富内容 |
| CodeGenAgent 的 | 生成带有奖章模式、验证和错误处理的PySpark/SQL代码 |
| 审阅者代理 | 验证代码质量、语法、是否符合规则,生成置信度评分 |
| PRCreatorAgent | 使用代码、文档、测试计划创建GitHub PR,跟踪合并状态 |
| 执行人代理人 | 在PR合并后在copula上执行代码,监控作业,收集指标 |
| ContextEnrichmentAgent | 分析层输出,提取模式/指标,为下一层准备上下文 |
| 摘要代理 | 生成包含公关链接、指标和业务见解的执行摘要 |
工具(tools/)
| 模块 | 目的 |
|---|---|
| databricks_tools.py | Rancher SDK包装器-SQL执行、仓库查询、作业监控 |
| github工具.py | GitHub API和GitPython-PR创建、合并检测、回购操作 |
| 奖章_工具.py | SQL模板生成器-带转换的青铜/银/金模式 |
| context_tools.py | 模式提取、数据质量分析、指标聚合 |
国家管理(state.py)
集中式ETL状态跟踪:
- 图层进度(已完成/剩余)
- 具有合并状态的PR历史记录
- 数据质量指标
- 错误记录
- 跨层上下文累积
编排(graph_workflow.py)
LangGraph工作流程包括:
- 基于审批状态的有条件路由
- 顺序多层处理
- 层间上下文丰富
- 错误处理和回退路径
MCP服务器(server.py)
FastMCP集成展示:
run_full_medallion_pipeline()-执行完整铜牌→银→Goldcheck_transformation_status()-监控PR/执行状态view_transformation_rules()-显示/管理转换规则- 事件驱动的PR合并检测
______________________________________________________________________
📋 转换规则
规则定义见 rules.txt 并应用于所有层:
青铜层规则
- ✅ 按原样摄取原始数据
- ✅ 添加审核元数据(摄入时间戳、源文件、行id)
- ✅ 使用Change Data Feed在三角洲湖存储
- ✅ 未应用转换
银层规则
- 🧹 去除重复项
- 🔤 标准化数据类型和文本大小写
- ❌ 删除空临界值
- 📊 添加质量标志(is_valid_record、data_quality_score)
- 🚫 筛选得分\<0.7的记录
- 📅 按日期划分
黄金层规则
- 📈 按城市和国家汇总
- 🧮 计算导出的指标(comfort_index等)
- 🗂️ 创建维度和事实表
- 🔍 过滤列上的Z顺序簇
- 📊 生成趋势视图(每周、每月)
数据质量检查
- 温度:-50°C至60°C
- 湿度:0-100%
- 坐标:纬度±90,经度±180
- 风/降水:≥0
- 异常值:标志超过3σ
______________________________________________________________________
🚀 主要特点
🤖 人工智能驱动决策
- 克劳德·索内特4.5为所有代理人提供权力
- 自然语言转换请求
- 具有错误处理功能的智能代码生成
- 自适应上下文丰富
🔄 上下文感知管道
- 每一层读取前一层的输出
- 新的模式和指标为下一个规划阶段提供信息
- 没有硬编码的假设-真正自适应
🔗 GitHub原生工作流
- 每层转换一个PR
- 自动合并检测触发执行
- Git历史记录中的完整审计跟踪
- 内置安全审查
📊 内置质量保证
- PR前的自动代码审查
- 每层跟踪的数据质量指标
- 强制执行质量分数阈值
- 全面的错误记录
⚡ 无服务器且可扩展
- Rancher无服务器SQL仓库
- 自动缩放执行器
- 三角洲湖变化数据馈送
- Unity治理目录
🎯 企业就绪
- 多层奖章建筑
- 概念转换
- 全面的日志记录和监控
- 合规就绪的数据沿袭
______________________________________________________________________
📦 环境设置
创建 .env 文件:
# Databricks
DATABRICKS_HOST=https://your-workspace.databricks.com
DATABRICKS_TOKEN=dapi2xxxxx
DATABRICKS_WAREHOUSE_ID=xxxxx
DATABRICKS_MODEL_ENDPOINT=databricks-claude-sonnet-4-5
# Model Configuration
MODEL_TEMPERATURE=0.2
MODEL_MAX_TOKENS=4096
MODEL_TOP_P=0.95
# GitHub
GITHUB_TOKEN=ghp_xxxxx
GITHUB_REPO_OWNER=your-username
GITHUB_REPO_NAME=your-repo
GIT_LOCAL_PATH=/path/to/local/repo
# Databricks Defaults
DEFAULT_CATALOG=samples
DEFAULT_SCHEMA=accuweather
# Logging
LOG_LEVEL=INFO
DEBUG_API_CALLS=false______________________________________________________________________
💻 用法示例
示例1:完整管道执行
from graph_workflow import MEDALLION_PIPELINE
state = MEDALLION_PIPELINE.invoke({
"user_query": "Clean weather data and create daily city aggregations",
"source_table": "samples.accuweather.forecast_daily_calendar_imperial"
})
print(state["executive_summary"])示例2:自定义转换规则
编辑 rules.txt 然后运行:
python server.py
# Access via VS Code MCP tools示例3:监控转换状态
from server import check_transformation_status
status = await check_transformation_status(pr_number=42)
print(status)示例4:视图转换规则
from server import view_transformation_rules
rules = await view_transformation_rules()
print(rules)______________________________________________________________________
📈 管道执行流程
User Query
↓
[MCP Server receives request]
↓
[Initialize ETLState]
↓
┌─ BRONZE LAYER ────────────────────────────────┐
│ Planner → CodeGen → Reviewer → PR → Executor │
│ [Create bronze_table from source] │
│ [Enrich context with schema/metrics] │
└───────────────────────────────────────────────┘
↓ [Context: Bronze table schema, row_count]
┌─ SILVER LAYER ────────────────────────────────┐
│ Planner (with Bronze context) → ... → Executor│
│ [Clean & deduplicate bronze_table] │
│ [Create silver_table with quality flags] │
│ [Enrich context with quality metrics] │
└───────────────────────────────────────────────┘
↓ [Context: Silver table schema, quality score]
┌─ GOLD LAYER ──────────────────────────────────┐
│ Planner (with Silver context) → ... → Executor│
│ [Aggregate silver_table by dimensions] │
│ [Create gold_table with derived metrics] │
└───────────────────────────────────────────────┘
↓
[Generate Executive Summary]
↓
Return to User with:
- 3 PR links (Bronze, Silver, Gold)
- Execution metrics & timing
- Data quality scores
- Error logs (if any)
- Next steps & recommendations______________________________________________________________________
🔍 监控与调试
启用调试日志
# In .env
LOG_LEVEL=DEBUG
DEBUG_API_CALLS=true检查执行日志
# View server logs
tail -f medallion_etl.log
# Check Git PR history
git log --oneline | grep "medallion\|etl"监控copula作业
- 打开copula工作区
- 导航到工作流→ 查看已执行的笔记本
- 检查作业运行的性能指标
______________________________________________________________________
📚 文件结构
ETLAgenticXMCP/
├── agents/ # AI-powered agents
│ ├── planner_agent.py # Transformation planning
│ ├── codegen_agent.py # SQL/PySpark generation
│ ├── reviewer_agent.py # Code quality review
│ ├── pr_creator_agent.py # GitHub PR management
│ ├── executor_agent.py # Databricks execution
│ ├── context_enrichment_agent.py # Context extraction
│ └── summary_agent.py # Report generation
├── tools/ # External integrations
│ ├── databricks_tools.py # Databricks SDK wrapper
│ ├── github_tools.py # GitHub API & GitPython
│ ├── medallion_tools.py # SQL generators
│ └── context_tools.py # Schema & metrics
├── graph_workflow.py # LangGraph orchestration
├── state.py # ETL state definition
├── server.py # FastMCP server
├── main.py # CLI entry point
├── rules.txt # Transformation rules
├── requirements.txt # Dependencies
└── .env.example # Environment template______________________________________________________________________
🛠️ 依赖项
- 兰格拉夫 -多代理编排
- 语言链 -LLM框架
- 数据块语言链 -copula集成
- sdk数据库 -Databricks API客户端
- fastmcp -MCP协议服务器
- gitpython -Git操作
- 请求: -HTTP客户端
- 皮丹提克 -数据验证
______________________________________________________________________
🤝 贡献
- 分叉存储库
- 创建要素分支:
git checkout -b feature/your-feature - 提交更改:
git commit -am 'Add feature' - 推送到分支:
git push origin feature/your-feature - 提交拉取请求
______________________________________________________________________
📄 许可证
MIT许可证-有关详细信息,请参阅许可证文件
______________________________________________________________________
📞 支持
对于问题、疑问或功能请求:
- 打开GitHub问题
- 检查中的现有文档
rules.txt - 查看代理日志以进行调试
______________________________________________________________________
内置于❤️ 数据工程团队\ *由Rancher Claude Sonnet 4.5和LangGraph提供技术支持*
