Самовідновлюване відхилення схеми за допомогою 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.
Покращена надійність
Постійна перевірка джерела та цілі запобігає несподіваним збоям у роботі.
Керування та можливість аудиту
Зміни версіонуються, перевіряються через PR, фіксуються в Airflow та відстежуються у Slack. Агент не може безпосередньо змінювати дані в продакшені.
Масштабованість
Ця сама схема може застосовуватися до багатьох DAG та сховищ (Snowflake, BigQuery, Databricks, Kafka Schema Registry) за допомогою легких адаптерів.
Майбутні покращення
Тестові випадки, створені за допомогою ШІ
Автоматичне створення тестів на перевірку запропонованих виправлень.
Аналіз впливу
Відстеження завдань, панелей керування та звітів, які постраждають внаслідок змін.
Виявлення семантичних змін
Виявлення змін у бізнес-значеннях, таких як:
status
зміна з:
ACTIVE / INACTIVE
на:
A / I
навіть якщо назви стовпців залишаються незмінними.
Автономне об’єднання
Виправлення з низьким рівнем ризику можуть бути автоматично об’єднані після перевірок — лише за суворими правилами.
Висновок
Зміни схеми є неминучими; проте виправлення лише вручну — ні. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) перетворює процес виявлення та інформування на процес виявлення, аналізу, створення PR та сповіщення. У міру розвитку агентних платформ інженерія даних переходить від простого спостереження за збоями до їх вирішення в рамках чітких правил.