Самастабілізуючыся дрейф схемы з Airflow, LangGraph і запросамі MCP
Калі ETL не выконваецца пад час змены столбца, трэба выявіць і класыфікаціяваць адхылэнне, напісаць праект выправлення коду, створыць PR на GitHub і паведаміць команду через Slack — без прымусовага выканання змян у продакшэне.
Введэнне
Змяны схемы ў базе дадзеных яўляюцца адной з галоўных прычын парадкі 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 і паведамленне. Калі платформы-агенты будуць развівацца, інжынерыя дадзеных пераходзіць ад простага стэрагавання проблем да ўсунення іх пад кераваннем.