Самовосстанавливающееся управление отклонениями схем с использованием Airflow, LangGraph и петиций MCP
Когда процесс ETL сбивается при изменении столбца, необходимо обнаружить и классифицировать отклонения, подготовить код для исправления, создать петицию в GitHub и уведомить команду через Slack — при этом не выпуская изменения непосредственно в производство.
Введение
Изменения схемы являются одной из основных причин сбоев в процессах ETL: добавление столбца, изменение типа или переименование в исходной базе данных могут нарушить преобразования, замедлить формирование отчетов и потратить много времени на работу с данными. Традиционные решения включают использование алертов и ручное устранение проблем. Системы с агентами позволяют перейти от обнаружения проблем к их направленному устранению.
В данной архитектуре используются:
- Apache Airflow для оркестрации задач
- PostgreSQL в качестве исходной и целевой базы данных
- LangGraph для рабочего процесса расследования с использованием агентов
- Хосты MCP для интеграции инструментов
- GitHub MCP для редактирования кода и отправки pull-запросов
- Slack MCP для уведомлений команды
Цикл действий — Обнаружить → Анализировать → Исправить → Отправить PR → Уведомить — не требует немедленного вмешательства человека для начала расследования.
Обзор архитектуры
+----------------------+
| 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) превращают процесс обнаружения и отслеживания в процесс обнаружения, анализа, создания поправок и уведомления. По мере совершенствования агентных платформ инжиниринг данных переходит от простого наблюдения за сбоями к их устранению в рамках четких правил управления.