This article is published in English.
Self-healing schema drift with Airflow, LangGraph, and MCP pull requests
When ETL fails on a column change, detect, classify drift, draft a code fix, open a GitHub PR, and Slack the team—without writing straight to production.
Introduction
Schema drift is a leading cause of ETL breakage: a harmless column add, type change, or rename in a source database can fail transformations, delay reporting, and burn data-engineering time. Classic responses are alerts plus manual fixes. Agentic systems can move from detection toward guided remediation.
This design combines:
- Apache Airflow for orchestration
- PostgreSQL as source and destination
- LangGraph for the agentic investigation workflow
- MCP hosts for tool integrations
- GitHub MCP for code edits and pull requests
- Slack MCP for team notifications
The loop is Detect → Analyze → Fix → Raise PR → Notify—without requiring an immediate human to start the investigation.
Architecture overview
+----------------------+
| 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
The problem: schema drift
Original schema
A customer table in source PostgreSQL:
customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)
A source team adds a column:
customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)
customer_segment VARCHAR(50)
The Airflow ETL still expects the old shape:
column "customer_segment" does not exist
or
INSERT has more expressions than target columns
The job fails; consumers feel it.
Step 1: Airflow detects failure
Airflow remains the orchestrator. On schema-mismatch failure:
def on_failure_callback(context):
send_to_langgraph(
dag_id=context["dag"].dag_id,
task_id=context["task_instance"].task_id,
error=context["exception"]
)
Instead of email-only alerts, Airflow forwards error context into a LangGraph agent. Payload:
{
"dag_id": "customer_load",
"task_id": "load_customer",
"error": "column customer_segment does not exist"
}
That payload starts autonomous remediation.
Step 2: LangGraph starts investigation
Agent 1: incident analyzer
Parses Airflow logs, names the impacted table, and proposes a probable cause. Output:
{
"table": "customer",
"issue": "schema_drift"
}
Step 3: Database MCP validates drift
The Database MCP host exposes tools against source and destination PostgreSQL. Example:
source_schema = get_source_schema("customer")
target_schema = get_target_schema("customer")
Source:
customer_id
customer_name
email
customer_segment
Target:
customer_id
customer_name
email
Comparison yields:
{
"new_column": "customer_segment",
"datatype": "VARCHAR(50)"
}
confirming a real drift event—not only a flaky network error.
Step 4: Root cause analysis agent
Classify drift:
Column added
customer_segment
Column removed
email
Datatype changed
customer_id BIGINT
Column renamed
customer_name
→
full_name
Classification drives remediation strategy.
Step 5: GitHub MCP generates a code fix
GitHub MCP scans the repo. Example ETL:
SELECT
customer_id,
customer_name,
email
FROM customer
Agent proposes:
SELECT
customer_id,
customer_name,
email,
customer_segment
FROM customer
or updates:
schema:
- customer_id
- customer_name
- email
- customer_segment
Step 6: Automated pull request
GitHub MCP creates a branch, applies edits, commits, and opens a PR:
Branch:
schema-fix/customer-segment
Commit:
Added customer_segment column support
PR:
Auto-generated schema drift remediation
PR description:
## 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
Changes stay auditable—no silent production edits.
Step 7: Slack MCP notifies the team
🚨 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
Engineers get a reviewable fix, not a bare stack trace.
LangGraph workflow design
Simplified flow:
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
Each node owns one responsibility with full traceability.
Benefits
Reduced MTTR
Traditional:
Failure
→ Alert
→ Investigation
→ Fix
→ PR
→ Deployment
hours to days. Agentic:
Failure
→ Detection
→ Analysis
→ PR Creation
→ Team Review
minutes to a PR.
Improved reliability
Continuous source/target validation cuts surprise production breaks.
Governance and auditability
Changes are versioned, PR-reviewed, Airflow-logged, and Slack-tracked. No direct prod writes by the agent.
Scalability
The same pattern extends to many DAGs and stores (Snowflake, BigQuery, Databricks, Kafka Schema Registry) with thin adapters.
Future enhancements
AI-generated test cases
Auto-build validation tests for proposed fixes.
Impact analysis
Map downstream jobs, dashboards, and reports touched by a drift.
Semantic drift detection
Catch business-meaning changes such as:
status
shifting from:
ACTIVE / INACTIVE
to:
A / I
even when column names stay put.
Autonomous merge
Low-risk fixes might auto-merge after checks—only with strict policy.
Conclusion
Schema drift is inevitable; manual-only remediation is not. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) turns detect-and-page into detect-analyze-PR-notify. As agentic platforms mature, data engineering moves from watching failures toward resolving them under governance.