Startseite / Artikel / Selbstheilendes Schema-Drift mit Airflow, LangGraph und MCP-Pull-Requests

Selbstheilendes Schema-Drift mit Airflow, LangGraph und MCP-Pull-Requests

Wenn bei einer Spaltenänderung der ETL-Prozess fehlschlägt, soll die Abweichung erkannt und klassifiziert werden, ein Code-Fix entworfen, eine GitHub-Pull-Request erstellt und das Team per Slack informiert werden – ohne direkt in die Produktion zu schreiben.

945 Wörter

Einführung

Schema-Drift ist eine der Hauptursachen für Störungen bei ETL-Prozessen: Die hinzufügung einer harmlosen Spalte, eine Typänderung oder ein Umbenennen in einer Quelldatenbank können Transformationen zum Scheitern bringen, die Berichterstellung verzögern und wertvolle Zeit im Bereich Datenengineering verschwenden. Klassische Gegenmaßnahmen sind Warnungen in Kombination mit manuellen Korrekturen. Agentenbasierte Systeme können von der Erkennung hin zu einer geführten Behebung übergehen.

Dieses Design kombiniert:

  • Apache Airflow für die Orchestrierung
  • PostgreSQL als Quelle und Ziel
  • LangGraph für den Workflow der agentenbasierten Untersuchung
  • MCP-Hosts für die Integration von Tools
  • GitHub MCP für Codeänderungen und Pull-Requests
  • Slack MCP für Teambenachrichtigungen

Der Ablauf ist Erkennen → Analysieren → Beheben → Pull-Request erstellen → Benachrichtigen – ohne dass sofort ein Mensch die Untersuchung starten muss.

Architekturübersicht

+----------------------+
| 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

Das Problem: Schema-Drift

Ursprüngliches Schema

Eine Kunden-Tabelle im Quell-PostgreSQL:

customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)

Ein Quell-Team fügt eine Spalte hinzu:

customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)
customer_segment VARCHAR(50)

Airflow ETL erwartet weiterhin die alte Struktur:

column "customer_segment" does not exist

oder

INSERT has more expressions than target columns

Die Aufgabe fehlschlägt; die Nutzer spüren dies.

Schritt 1: Airflow erkennt den Fehler

Airflow bleibt der Orchesterer. Bei Fehlern aufgrund von Schema-Unterschieden:

def on_failure_callback(context):
    send_to_langgraph(
        dag_id=context["dag"].dag_id,
        task_id=context["task_instance"].task_id,
        error=context["exception"]
    )

Anstelle von nur E-Mail-Benachrichtigungen leitet Airflow den Fehlerkontext an einen LangGraph-Agenten weiter. Payload:

{
  "dag_id": "customer_load",
  "task_id": "load_customer",
  "error": "column customer_segment does not exist"
}

Dieser Payload startet eine automatische Behebung.

Schritt 2: LangGraph beginnt mit der Untersuchung

Agent 1: Incident-Analysator

Er analysiert die Airflow-Logs, benennt die betroffene Tabelle und schlägt eine wahrscheinliche Ursache vor. Ausgabe:

{
  "table": "customer",
  "issue": "schema_drift"
}

Schritt 3: Database MCP überprüft Abweichungen

Der Database MCP-Host stellt Tools für den Quell- und Ziel-PostgreSQL zur Verfügung. Beispiel:

source_schema = get_source_schema("customer")
target_schema = get_target_schema("customer")

Quelle:

customer_id
customer_name
email
customer_segment

Ziel:

customer_id
customer_name
email

Der Vergleich ergibt:

{
  "new_column": "customer_segment",
  "datatype": "VARCHAR(50)"
}

eine echte Abweichung – nicht nur einen vorübergehenden Netzwerkfehler.

Schritt 4: Analyse des Ursachenfaktors

Klassifizierung der Abweichung:

Spalte hinzugefügt

customer_segment

Spalte entfernt

email

Datentyp geändert

customer_id BIGINT

Spalte umbenannt

customer_name
→
full_name

Die Klassifizierung bestimmt die Korrekturmaßnahmen.

Schritt 5: GitHub MCP erzeugt eine Codekorrektur

GitHub MCP durchsucht das Repository. Beispiel für ETL:

SELECT
 customer_id,
 customer_name,
 email
FROM customer

Der Agent schlägt vor:

SELECT
 customer_id,
 customer_name,
 email,
 customer_segment
FROM customer

oder aktualisiert:

schema:
  - customer_id
  - customer_name
  - email
  - customer_segment

Schritt 6: Automatisierte Pull-Request-Erstellung

GitHub MCP erstellt einen Branch, wendet Änderungen an, macht Commits und öffnet einen PR:

Branch:
schema-fix/customer-segment

Commit:
Added customer_segment column support

PR:
Auto-generated schema drift remediation

Beschreibung des PRs:

## 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

Die Änderungen bleiben nachvollziehbar – es finden keine stillen Änderungen in der Produktion statt.

Schritt 7: Slack MCP benachrichtigt das 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

Die Entwickler erhalten eine überprüfbare Lösung, nicht nur einen reinen Stack-Trace.

Design des LangGraph-Ablaufs

Einfacher Ablauf:

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

Jeder Knoten hat eine einzige Verantwortung mit vollständiger Nachvollziehbarkeit.

Vorteile

Kürzere MTTR-Zeiten

Traditionell:

Failure
→ Alert
→ Investigation
→ Fix
→ PR
→ Deployment

Stunden bis Tage. Mit Agenten:

Failure
→ Detection
→ Analysis
→ PR Creation
→ Team Review

Minuten bis zur Freigabe eines PRs.

Bessere Zuverlässigkeit

Durch kontinuierliche Validierung von Quelle und Ziel werden unerwartete Ausfälle in der Produktion vermieden.

Steuerung und Prüfbarkeit

Änderungen werden versioniert, im PR geprüft, in Airflow protokolliert und in Slack nachverfolgt. Der Agent führt keine direkten Schreibvorgänge in der Produktion durch.

Skalierbarkeit

Das gleiche Muster lässt sich auf viele DAGs sowie Speicherplattformen (Snowflake, BigQuery, Databricks, Kafka Schema Registry) mit leichten Adaptern anwenden.

Zukünftige Erweiterungen

Von KI generierte Testfälle

Automaatische Erstellung von Validierungstests für vorgeschlagene Korrekturen.

Auswirkungsanalyse

Karten der nachgelagerten Aufgaben, Dashboards und Berichte, die durch Abweichungen betroffen sind.

Semantische Abweichungserkennung

Aufnahme von Veränderungen im geschäftlichen Kontext wie zum Beispiel:

status

von:

ACTIVE / INACTIVE

nach:

A / I

selbst wenn die Spaltennamen unverändert bleiben.

Autonome Zusammenführung

Korrekturen mit geringem Risiko können nach Überprüfungen automatisch zusammengeführt werden – allerdings nur unter strengen Richtlinien.

Fazit

Schema-Abweichungen sind unvermeidlich; eine rein manuelle Behebung hingegen nicht. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) verwandeln das Verfahren „erkennen und melden“ in „erkennen, analysieren, Pull Request erstellen und benachrichtigen“. Mit der Weiterentwicklung agiler Plattformen rückt die Datenverarbeitung stärker von der bloßen Fehlerbeobachtung zur lösungsorientierten Bearbeitung unter Governance-Regeln.