Головна / Статті / Самовідновлюване відхилення схеми за допомогою 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.

Покращена надійність

Постійна перевірка джерела та цілі запобігає несподіваним збоям у роботі.

Керування та можливість аудиту

Зміни версіонуються, перевіряються через PR, фіксуються в Airflow та відстежуються у Slack. Агент не може безпосередньо змінювати дані в продакшені.

Масштабованість

Ця сама схема може застосовуватися до багатьох DAG та сховищ (Snowflake, BigQuery, Databricks, Kafka Schema Registry) за допомогою легких адаптерів.

Майбутні покращення

Тестові випадки, створені за допомогою ШІ

Автоматичне створення тестів на перевірку запропонованих виправлень.

Аналіз впливу

Відстеження завдань, панелей керування та звітів, які постраждають внаслідок змін.

Виявлення семантичних змін

Виявлення змін у бізнес-значеннях, таких як:

status

зміна з:

ACTIVE / INACTIVE

на:

A / I

навіть якщо назви стовпців залишаються незмінними.

Автономне об’єднання

Виправлення з низьким рівнем ризику можуть бути автоматично об’єднані після перевірок — лише за суворими правилами.

Висновок

Зміни схеми є неминучими; проте виправлення лише вручну — ні. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) перетворює процес виявлення та інформування на процес виявлення, аналізу, створення PR та сповіщення. У міру розвитку агентних платформ інженерія даних переходить від простого спостереження за збоями до їх вирішення в рамках чітких правил.