Inicio / Artículos / Deriva de esquemas autoreparables con Airflow, LangGraph y solicitudes de pull de MCP

Deriva de esquemas autoreparables con Airflow, LangGraph y solicitudes de pull de MCP

Cuando el proceso ETL falla al modificar una columna, se debe detectar y clasificar la desviación, redactar una solución en código, abrir un PR en GitHub y notificar al equipo por Slack, todo ello sin enviar cambios directamente a producción.

945 palabras

Introducción

La deriva de esquemas es una de las principales causas de fallos en los procesos ETL: la adición de columnas, el cambio de tipo o el renombramiento en una base de datos de origen pueden hacer que fallen las transformaciones, retrasar la generación de informes y consumir tiempo valioso en ingeniería de datos. Las respuestas habituales son las alertas junto con correcciones manuales. Los sistemas basados en agentes pueden pasar de la detección a una solución guiada.

Este diseño combina:

  • Apache Airflow para la orquestación
  • PostgreSQL como origen y destino
  • LangGraph para el flujo de trabajo de investigación mediante agentes
  • Anfitriones MCP para la integración de herramientas
  • GitHub MCP para ediciones de código y solicitudes de pull request
  • Slack MCP para notificaciones al equipo

El ciclo es Detectar → Analizar → Corregir → Crear PR → Notificar—sin necesidad de que una persona intervenga de inmediato para iniciar la investigación.

Visión general de la arquitectura

+----------------------+
| 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

El problema: la deriva de esquemas

Estructura original

Una tabla de clientes en el PostgreSQL de origen:

customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)

Un equipo del origen agrega una columna:

customer_id INT
customer_name VARCHAR(100)
email VARCHAR(255)
customer_segment VARCHAR(50)

Airflow ETL sigue esperando la estructura antigua:

column "customer_segment" does not exist

o

INSERT has more expressions than target columns

La tarea falla; los usuarios lo perciben.

Paso 1: Airflow detecta el fallo

Airflow sigue siendo el orquestador. En caso de fallo por discrepancia en la estructura:

def on_failure_callback(context):
    send_to_langgraph(
        dag_id=context["dag"].dag_id,
        task_id=context["task_instance"].task_id,
        error=context["exception"]
    )

En lugar de solo alertas por correo electrónico, Airflow envía el contexto del error a un agente de LangGraph. Carga útil:

{
  "dag_id": "customer_load",
  "task_id": "load_customer",
  "error": "column customer_segment does not exist"
}

Esa carga útil inicia la corrección autónoma.

Paso 2: LangGraph inicia la investigación

Agente 1: analizador de incidentes

Analiza los registros de Airflow, identifica la tabla afectada y propone una causa probable. Salida:

{
  "table": "customer",
  "issue": "schema_drift"
}

Paso 3: Database MCP valida la desviación

El host de Database MCP ofrece herramientas para el PostgreSQL de origen y destino. Ejemplo:

source_schema = get_source_schema("customer")
target_schema = get_target_schema("customer")

Fuente:

customer_id
customer_name
email
customer_segment

Destino:

customer_id
customer_name
email

El comparativo muestra:

{
  "new_column": "customer_segment",
  "datatype": "VARCHAR(50)"
}

lo que confirma un verdadero evento de desviación, y no solo un error de red ocasional.

Paso 4: Agente de análisis de la causa raíz

Clasificar la desviación:

Columna añadida

customer_segment

Columna eliminada

email

Cambio en el tipo de dato

customer_id BIGINT

Columna renombrada

customer_name
→
full_name

La clasificación determina la estrategia de corrección.

Paso 5: GitHub MCP genera una solución en código

GitHub MCP escanea el repositorio. Ejemplo de ETL:

SELECT
 customer_id,
 customer_name,
 email
FROM customer

El agente propone:

SELECT
 customer_id,
 customer_name,
 email,
 customer_segment
FROM customer

o actualiza:

schema:
  - customer_id
  - customer_name
  - email
  - customer_segment

Paso 6: Solicitud de integración automatizada

GitHub MCP crea una rama, aplica los cambios, realiza commits y abre una solicitud de integración:

Branch:
schema-fix/customer-segment

Commit:
Added customer_segment column support

PR:
Auto-generated schema drift remediation

Descripción de la solicitud de integración:

## 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

Los cambios permanecen auditables; no hay ediciones silenciosas en producción.

Paso 7: Slack MCP notifica al equipo

🚨 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

Los ingenieros reciben una solución revisable, no solo un rastro de errores básico.

Diseño del flujo de trabajo de LangGraph

Flujo simplificado:

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

Cada nodo tiene una responsabilidad única con total trazabilidad.

Beneficios

Reducción del MTTR

En el método tradicional:

Failure
→ Alert
→ Investigation
→ Fix
→ PR
→ Deployment

de horas a días. Con el enfoque basado en agentes:

Failure
→ Detection
→ Analysis
→ PR Creation
→ Team Review

de minutos a la presentación de una PR.

Mayor fiabilidad

La validación continua de fuentes y destinos evita interrupciones inesperadas en producción.

Gobernanza y auditabilidad

Los cambios se versionan, se revisan en las PR, se registran en Airflow y se rastrean en Slack. El agente no realiza escrituras directas en producción.

Escalabilidad

El mismo patrón se puede aplicar a numerosos DAG y almacenamientos (Snowflake, BigQuery, Databricks, Kafka Schema Registry) mediante adaptadores simples.

Mejoras futuras

Casos de prueba generados por IA

Generación automática de pruebas de validación para las correcciones propuestas.

Análisis de impacto

Mapear las tareas posteriores, los paneles de control y los informes afectados por una desviación.

Detección de desviaciones semánticas

Detener cambios en el significado empresarial como, por ejemplo:

status

el paso de:

ACTIVE / INACTIVE

a:

A / I

incluso cuando los nombres de las columnas permanecen sin cambios.

Fusión autónoma

Las correcciones de bajo riesgo pueden fusionarse automáticamente tras las verificaciones, pero solo bajo políticas estrictas.

Conclusión

La desviación del esquema es inevitable; la solución exclusivamente manual no lo es. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) convierte el proceso de detección y notificación en uno que incluye análisis, creación de PRs y notificaciones. A medida que las plataformas agentes maduran, la ingeniería de datos pasa de simplemente observar los fallos a resolverlos bajo un marco de gobernanza.

Lecturas relacionadas

  • Enterprise MCP en Bedrock AgentCore con Okta SSO y Trino — Cerebro OAuth de Lambda, relé PKCE, gateway y entorno de ejecución de AgentCore, además de conexiones de Trino por usuario sin fugas de tokens compartidos.