Главная / Статьи / Самовосстанавливающееся управление отклонениями схем с использованием Airflow, LangGraph и петиций MCP

Самовосстанавливающееся управление отклонениями схем с использованием Airflow, LangGraph и петиций MCP

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

945 слов

Введение

Изменения схемы являются одной из основных причин сбоев в процессах 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) превращают процесс обнаружения и отслеживания в процесс обнаружения, анализа, создания поправок и уведомления. По мере совершенствования агентных платформ инжиниринг данных переходит от простого наблюдения за сбоями к их устранению в рамках четких правил управления.