首页 / 文章 / 利用 Airflow、LangGraph 和 MCP 提交 Pull Request 实现自愈式模式漂移检测

利用 Airflow、LangGraph 和 MCP 提交 Pull Request 实现自愈式模式漂移检测

当列变更导致ETL失败时,能够检测并分类数据漂移,起草代码修复方案,创建GitHub PR,并通过Slack通知团队——且无需直接将修改内容部署到生产环境。

945 词

简介

模式漂移是导致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)将“检测并反馈”转变为“检测-分析-提交修复请求-通知”。随着智能平台的成熟,数据工程正从单纯监控故障转向在规范管理下解决故障。