利用 Airflow、LangGraph 和 MCP 提交 Pull Request 实现自愈式模式漂移检测
当列变更导致ETL失败时,能够检测并分类数据漂移,起草代码修复方案,创建GitHub PR,并通过Slack通知团队——且无需直接将修改内容部署到生产环境。
简介
模式漂移是导致ETL流程失败的主要原因:源数据库中看似无害的列添加、类型变更或重命名都可能使转换失败、延迟报表生成,并耗费大量数据工程时间。传统的应对方式是发送警报并手动修复问题。而代理系统则能从检测阶段迈向引导式修复。
该设计整合了以下组件:
- 用于流程协调的Apache Airflow
- 作为源数据和目标数据的PostgreSQL
- 用于代理式问题排查工作流的LangGraph
- 用于工具集成的MCP主机
- 用于代码编辑和提交Pull Request的GitHub MCP
- 用于团队通知的Slack MCP
整个流程为检测 → 分析 → 修复 → 提交PR → 通知——无需人工立即介入即可启动问题排查。
架构概览
+----------------------+
| Source PostgreSQL |
+----------+-----------+
|
v
+----------------------+
| Airflow ETL Pipeline |
+----------+-----------+
|
v
+----------------------+
| Destination Postgres |
+----------+-----------+
|
|
Failure/Event
|
v
+----------------------+
| LangGraph Agent |
+----------+-----------+
|
+-------+--------+
| |
v v
Database MCP GitHub MCP
| |
Schema Diff Code Fix
Validation PR Creation
|
+-------+
|
v
Slack MCP
|
v
Data Engineering Team
问题:模式漂移
原始架构
源端 PostgreSQL 中的客户表:
customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)
源端团队添加了一列:
customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)
customer_segment VARCHAR(50)
但 Airflow ETL 仍期望旧的表结构:
column "customer_segment" does not exist
或者
INSERT has more expressions than target columns
任务将失败,使用者会感受到影响。
步骤 1:Airflow 检测到故障
Airflow 依然是任务协调器。在架构不匹配导致故障时:
def on_failure_callback(context):
send_to_langgraph(
dag_id=context["dag"].dag_id,
task_id=context["task_instance"].task_id,
error=context["exception"]
)
Airflow 不仅会发送邮件警报,还会将错误信息转发给 LangGraph 代理。传输的数据包括:
{
"dag_id": "customer_load",
"task_id": "load_customer",
"error": "column customer_segment does not exist"
}
这些数据会触发自动修复流程。
步骤 2:LangGraph 开始调查
代理 1:故障分析器
它会解析 Airflow 日志,确定受影响的表,并提出可能的故障原因。输出结果为:
{
"table": "customer",
"issue": "schema_drift"
}
步骤 3:数据库 MCP 验证数据偏差
数据库 MCP 主机提供了用于操作源端和目标端 PostgreSQL 的工具。示例:
source_schema = get_source_schema("customer")
target_schema = get_target_schema("customer")
来源:
customer_id
customer_name
email
customer_segment
目标:
customer_id
customer_name
email
对比结果为:
{
"new_column": "customer_segment",
"datatype": "VARCHAR(50)"
}
确认是真正的数据漂移事件,而非偶发的网络错误。
第4步:根本原因分析模块
对数据漂移进行分类:
新增列
customer_segment
删除列
email
数据类型变更
customer_id BIGINT
列重命名
customer_name
→
full_name
分类结果将决定相应的修复策略。
第5步:GitHub MCP生成代码修复方案
GitHub MCP会扫描代码仓库。例如ETL流程:
SELECT
customer_id,
customer_name,
email
FROM customer
该模块会提出:
SELECT
customer_id,
customer_name,
email,
customer_segment
FROM customer
或进行更新:
schema:
- customer_id
- customer_name
- email
- customer_segment
第6步:自动生成拉取请求
GitHub MCP会创建一个分支,应用修改内容并提交代码,随后打开拉取请求:
Branch:
schema-fix/customer-segment
Commit:
Added customer_segment column support
PR:
Auto-generated schema drift remediation
拉取请求的描述内容为:
## Schema Drift Detected
Table: customer
Change:
Added customer_segment VARCHAR(50)
Actions Performed:
- Updated ETL query
- Updated schema definition
- Added validation test
Generated by LangGraph Agent
所有更改均可被追溯,不会存在隐秘的生产环境修改。
第7步:Slack MCP通知团队
🚨 Schema Drift Detected
Table: customer
Drift Type:
New Column Added
Action Taken:
✔ Code Updated
✔ Pull Request Created
PR:
https://github.com/org/repo/pull/123
Review Required:
Data Engineering Team
工程师会收到可供审查的修复方案,而非单纯的堆栈跟踪信息。
LangGraph工作流设计
简化后的流程:
Start
|
v
Incident Analyzer
|
v
Schema Validation Agent
|
v
Drift Classification Agent
|
v
Code Generation Agent
|
v
GitHub MCP Agent
|
v
Slack Notification Agent
|
v
End
每个节点仅承担一项职责,且具备完整的可追溯性。
优势
缩短平均修复时间
传统方式:
Failure
→ Alert
→ Investigation
→ Fix
→ PR
→ Deployment
需要数小时到数天。自动化方式:
Failure
→ Detection
→ Analysis
→ PR Creation
→ Team Review
仅需几分钟即可提交修复请求。
提升可靠性
持续的源代码/目标数据验证可避免意外的生产环境故障。
可管控性与审计性
所有变更都会被版本控制、经过修复请求审核、记录在Airflow日志中,并通过Slack进行追踪。代理程序不会直接写入生产环境。
可扩展性
通过轻量级的适配器,同一模式可应用于多种DAG及数据存储系统(如Snowflake、BigQuery、Databricks、Kafka Schema Registry)。
未来改进方向
AI生成的测试用例
自动为拟议的修复方案生成验证测试。
影响分析
展示因数据偏移而受影响的下游任务、仪表板及报告。
语义偏移检测
能够捕捉业务含义上的变化,例如:
status
从:
ACTIVE / INACTIVE
变为:
A / I
即便列名保持不变也是如此。
自动合并
低风险修复方案在经过检查后可能会自动合并——但仅限于严格遵循相关策略的情况下。
结论
模式偏移是不可避免的,但仅靠人工修复并非如此。Airflow + PostgreSQL + LangGraph + MCP(数据库、GitHub、Slack)将“检测并反馈”转变为“检测-分析-提交修复请求-通知”。随着智能平台的成熟,数据工程正从单纯监控故障转向在规范管理下解决故障。