Галоўная / Артыкулы / Самастабілізуючыся дрейф схемы з Airflow, LangGraph і запросамі MCP

Самастабілізуючыся дрейф схемы з Airflow, LangGraph і запросамі MCP

Калі ETL не выконваецца пад час змены столбца, трэба выявіць і класыфікаціяваць адхылэнне, напісаць праект выправлення коду, створыць PR на GitHub і паведаміць команду через Slack — без прымусовага выканання змян у продакшэне.

945 слоў

Введэнне

Змяны схемы ў базе дадзеных яўляюцца адной з галоўных прычын парадкі ETL: дадаць безнэпакойную столбец, змяніць його тип або перайменаваць у выхіднай базе дадзеных можа спрычыніць неудачу трансфармацый, запазджэння у падготовцы звятынаў і витрачанне часу на роботу з дадзеннямі. Класычныя способы рашэння — аўтаматычныя апавешчэнні і ручныя практыкі. Системы з агентамі могу перайсці ад выяўлення проблемы да кераванага ўсунення яе.

У гэтым дизайне выкарыстоўваюцца:

  • Apache Airflow для оркестрацыі
  • PostgreSQL як выхідная і канечная база дадзеных
  • LangGraph для рабочага процесу агентскага расследавання
  • MCP-хосты для інтеграціі інструментаў
  • GitHub MCP для рэдагавання коду і падачы запрошэнняў
  • Slack MCP для апавешчэнняя каманды

Цикл складаецца з Выяўленне → Аналіз → Усуненне → Падача запрошэнняя → Апавешчэнне — без неабходнасці негайнага втручання чалавека для пачатку расследавання.

Аптака архітэктуры

+----------------------+
| 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: Database MCP пераканальвае наявнасць адхылення

Хост Database 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 стварае галузь, застосоўвае змены, фіксуе іх і ачынае PR:

Branch:
schema-fix/customer-segment

Commit:
Added customer_segment column support

PR:
Auto-generated schema drift remediation

Апісанне PR:

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

мінуты да прайняття змян.

Падышанне надзеямоўнасці

Постацэйная перакананне выхадных і канечных данных скарочвае неспакоўныя перерывы ў роботе.

Карацель і можлівасць аудыту

Змяны версіююцца, пераглядаюцца ў рамках PR, фіксуюцца в Airflow і стежыцца ў Slack. Агент не выканае безпасэчных змян у прыемным сервісе.

Масштабаванне

Той жа патэрн можна застосавіць да багоў DAG і хранальняў (Snowflake, BigQuery, Databricks, Kafka Schema Registry) за дапамогою лёгкіх адаптараў.

Будучыя падборкі

Тэстовыя кейсы, створаныя за дапамогою AI

Автаматычна ствароўка тэстаў на перакананне для запропаных пасправек.

Аналіз адзейнасці

Апісваць задачы, панелі керування і атласы, якія паўтараюць уплыв.

Адказванне на семантычны уплыв

Зафіксаваць змены ў бізнес-значэннях, такія як:

status

змена з:

ACTIVE / INACTIVE

на:

A / I

нават калі назвы столбцоў застаюцца тымі ж.

Автонамны з’еднанне

Пасправкі з низкім рызыкам могу быць автаматычна з’еднаныя пасля перакананняў — толькі за строгай політыкай.

Вывык

Уплыв схемы ўнэварцыйны; толькі ручная корэкцыя — ней. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) ператварае процес адказвання з простага выяўлення на выяўленне, аналіз, ствароўку PR і паведамленне. Калі платформы-агенты будуць развівацца, інжынерыя дадзеных пераходзіць ад простага стэрагавання проблем да ўсунення іх пад кераваннем.