Strona główna / Artykuły / Samonaprawiające się odchylenia schematu przy użyciu Airflow, LangGraph i prośb o integrację MCP

Samonaprawiające się odchylenia schematu przy użyciu Airflow, LangGraph i prośb o integrację MCP

Gdy proces ETL zawodzi podczas zmiany kolumny, należy zidentyfikować i sklasyfikować odchylenia, przygotować poprawkę kodu, otworzyć PR na GitHubie oraz poinformować zespół przez Slack — bez konieczności wprowadzania zmian bezpośrednio do środowiska produkcyjnego.

945 słów

Wprowadzenie

Drift schematu jest jedną z głównych przyczyn awarii procesów ETL: dodanie bezpiecznej kolumny, zmiana typu lub przemianowanie nazwy w bazie danych źródłowej może spowodować niepowodzenie transformacji, opóźnienie raportowania oraz stratę czasu poświęconego inżynierii danych. Klasycznymi rozwiązaniami są alerty oraz ręczne korekty. Systemy agentowe mogą przechodzić od wykrywania problemów do kierowanego ich naprawiania.

Ta architektura łączy w sobie:

  • Apache Airflow do orkiestracji
  • PostgreSQL jako baza źródłowa i docelowa
  • LangGraph do przeprowadzania zbadania przez systemy agentowe
  • Hosty MCP do integracji narzędzi
  • GitHub MCP do edycji kodu i zgłaszania pull requestów
  • Slack MCP do powiadamiania zespołu

Pętla działania to Wykryć → Przeanalizować → Naprawić → Zgłosić PR → Poinformować — bez konieczności natychmiastowego interwencji człowieka w celu rozpoczęcia dochodzenia.

Przegląd architektury

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

Problem: drift schematu

Oryginalna struktura

Tabela klientów w oryginalnym PostgreSQL:

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

Zespół odpowiedzialny za źródło dodaje kolumnę:

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

Airflow ETL nadal oczekuje starej struktury:

column "customer_segment" does not exist

lub

INSERT has more expressions than target columns

Zadanie zawodzi; użytkownicy tego odczuwają.

Krok 1: Airflow wykrywa awarię

Airflow pozostaje orkiestratorem. W przypadku błędu związанego z niezgodnością struktury:

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

Zamiast tylko powiadamień e-mailem, Airflow przekazuje kontekst błędu do agenta LangGraph. Treść przesyłana:

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

Ta treść uruchamia autonomiczną korektę.

Krok 2: LangGraph rozpoczyna dochodzenie

Agent 1: analityk incydentów

Analizuje logi Airflow, określa nazwę dotkniętej tabeli i proponuje prawdopodobną przyczynę. Wynik:

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

Krok 3: Database MCP weryfikuje odchylenia

Host Database MCP udostępnia narzędzia do pracy z oryginalnym i docelowym PostgreSQL. Przykład:

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

Źródło:

customer_id
customer_name
email
customer_segment

Cel:

customer_id
customer_name
email

Porównanie pokazuje:

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

potwierdzające rzeczywisty incydent odchylenia — a nie tylko przelotny błąd sieciowy.

Krok 4: Agent do analizy przyczyn podstawowych

Klasyfikuj odchylenie:

Kolumna dodana

customer_segment

Kolumna usunięta

email

Zmieniony typ danych

customer_id BIGINT

Kolumna przemianowana

customer_name
→
full_name

Klasyfikacja określa strategię naprawczą.

Krok 5: GitHub MCP generuje poprawkę kodu

GitHub MCP skanuje repozytorium. Przykład ETL:

SELECT
 customer_id,
 customer_name,
 email
FROM customer

Agent proponuje:

SELECT
 customer_id,
 customer_name,
 email,
 customer_segment
FROM customer

lub aktualizuje:

schema:
  - customer_id
  - customer_name
  - email
  - customer_segment

Krok 6: Automatyczna prośba o połączenie

GitHub MCP tworzy gałąź, aplikuje zmiany, dokonuje komitowania i otwiera prośbę o połączenie:

Branch:
schema-fix/customer-segment

Commit:
Added customer_segment column support

PR:
Auto-generated schema drift remediation

Opis prośby o połączenie:

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

Zmiany pozostają poddawane audytowi — bez ukrytych edycji w produkcji.

Krok 7: Slack MCP powiadamia zespół

🚨 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

Inżynierowie otrzymują gotową do sprawdzenia poprawkę, a nie tylko surowy zapis ścieżki błędu.

Projekt przepływu LangGraph

Społeczny przepływ:

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

Każdy węzeł odpowiada za jedną funkcję przy pełnej możliwości śledzenia.

Korzyści

Zmniejszony czas reakcji

W tradycyjnym podejściu:

Failure
→ Alert
→ Investigation
→ Fix
→ PR
→ Deployment

godziny do dni. W podejściu opartym na agentach:

Failure
→ Detection
→ Analysis
→ PR Creation
→ Team Review

minuty do złożenia PR.

Lepsza niezawodność

Ciągła weryfikacja źródła i celu zapobiega nieoczekiwanym awariom w produkcji.

Zarządzanie i możliwość audytu

Zmiany są wersjonowane, sprawdzane pod kątem PR, rejestrowane przez Airflow i śledzone w Slacku. Agent nie może bezpośrednio zapisywać danych do produkcji.

Możliwości skalowania

To samo rozwiązanie można zastosować do wielu DAG-ów oraz baz danych (Snowflake, BigQuery, Databricks, Kafka Schema Registry) za pomocą prostych adapterów.

Budujące się funkcje w przyszłości

Przypadki testowe generowane przez AI

Automaticzne tworzenie testów walidacyjnych dla proponowanych poprawek.

Analiza wpływu

Mapowanie zadań, paneli kontrolnych i raportów dotkniętych odchyleniem.

Detekcja semantycznego odchylenia

Zauważanie zmian o znaczeniu biznesowym, takich jak:

status

przejście od:

ACTIVE / INACTIVE

do:

A / I

nawet gdy nazwy kolumn pozostają bez zmian.

Autonomiczne łączenie

Poprawki o niskim ryzyku mogą być automatycznie łączone po sprawdzeniach — tylko przy ścisłej polityce.

Wniosek

Odchylenia schematu są nieuniknione; jedynie ręczne naprawy nie. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) przekształcają proces wykrywania i raportowania w proces wykrywania, analizy, tworzenia PR oraz powiadamiania. W miarę dojrzewania platform agentowych inżynieria danych przechodzi od obserwowania awarii do ich rozwiązywania w ramach określonych zasad.

Pozycje pokrewne

  • Enterprise MCP na Bedrock AgentCore z Okta SSO i Trino — silnik Lambda OAuth, przekaźnik PKCE, brama i środowisko działania AgentCore oraz połączenia Trino dla każdego użytkownika bez wycieku wspólnych tokenów.