Accueil / Articles / Dérive des schémas auto-régénératifs avec Airflow, LangGraph et les demandes de pull MCP

Dérive des schémas auto-régénératifs avec Airflow, LangGraph et les demandes de pull MCP

Lorsque le processus ETL échoue en raison d’une modification de colonne, détecter la dérive, la classifier, rédiger une correction de code, ouvrir un PR sur GitHub et informer l’équipe par Slack — sans jamais mettre à jour directement la production.

945 mots

Introduction

Le dérive de schéma est l’une des principales causes d’erreurs dans les processus ETL : l’ajout d’une colonne, un changement de type ou un renommage dans une base de données source peuvent faire échouer les transformations, retarder la génération des rapports et consommer beaucoup de temps dans le domaine du génie des données. Les solutions classiques consistent en des alertes accompagnées de corrections manuelles. Les systèmes agents peuvent passer de la simple détection à une correction guidée.

Cette conception combine :

  • Apache Airflow pour l’orchestration
  • PostgreSQL en tant que source et destination
  • LangGraph pour le flux de travail d’enquête automatisé
  • Des hôtes MCP pour l’intégration des outils
  • GitHub MCP pour les modifications de code et les demandes de fusion
  • Slack MCP pour les notifications d’équipe

Le cycle est Détection → Analyse → Correction → Soumission d’une demande de fusion → Notification—sans nécessiter qu’un humain intervienne immédiatement pour lancer l’enquête.

Aperçu de l’architecture

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

Le problème : la dérive de schéma

Schéma original

Une table de clients dans le PostgreSQL source :

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

Une équipe source ajoute une colonne :

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

L’ETL d’Airflow attend toujours la structure ancienne :

column "customer_segment" does not exist

ou

INSERT has more expressions than target columns

La tâche échoue ; les utilisateurs en subissent les conséquences.

Étape 1 : Airflow détecte l’échec

Airflow reste l’orchestrateur. En cas d’échec lié à une incohérence de schéma :

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

Au lieu d’alertes uniquement par e-mail, Airflow transmet le contexte de l’erreur à un agent LangGraph. Charge utile :

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

Cette charge utile déclenche une correction autonome.

Étape 2 : LangGraph commence l’enquête

Agent 1 : analyseur d’incidents

Il analyse les journaux d’Airflow, identifie la table affectée et propose une cause probable. Résultat :

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

Étape 3 : Database MCP valide les écarts

L’hôte Database MCP met à disposition des outils pour le PostgreSQL source et de destination. Exemple :

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

Source :

customer_id
customer_name
email
customer_segment

Cible :

customer_id
customer_name
email

La comparaison révèle :

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

ce qui confirme un véritable événement de dérive — et non simplement une erreur réseau intermittente.

Étape 4 : Agent d’analyse de la cause racine

Classer la dérive :

Colonne ajoutée

customer_segment

Colonne supprimée

email

Type de données modifié

customer_id BIGINT

Colonne renommée

customer_name
→
full_name

La classification guide la stratégie de correction.

Étape 5 : GitHub MCP génère une correction de code

GitHub MCP scanne le répertoire. Exemple d’ETL :

SELECT
 customer_id,
 customer_name,
 email
FROM customer

L’agent propose :

SELECT
 customer_id,
 customer_name,
 email,
 customer_segment
FROM customer

ou met à jour :

schema:
  - customer_id
  - customer_name
  - email
  - customer_segment

Étape 6 : Demande de fusion automatisée

GitHub MCP crée une branche, applique les modifications, effectue des commits et ouvre une demande de fusion :

Branch:
schema-fix/customer-segment

Commit:
Added customer_segment column support

PR:
Auto-generated schema drift remediation

Description de la demande de fusion :

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

Les modifications restent auditable — aucune modification silencieuse en production.

Étape 7 : Slack MCP informe l’équipe

🚨 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

Les ingénieurs reçoivent une solution à examiner, et non simplement un historique d’erreurs.

Conception du flux LangGraph

Flux simplifié :

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

Chaque nœud a une seule responsabilité, avec une traçabilité complète.

Avantages

Réduction du MTTR

Méthode traditionnelle :

Failure
→ Alert
→ Investigation
→ Fix
→ PR
→ Deployment

heures à jours. Méthode agente :

Failure
→ Detection
→ Analysis
→ PR Creation
→ Team Review

minutes jusqu’à la soumission d’un PR.

Fiabilité améliorée

La validation continue de la source et de la cible réduit les pannes inattendues en production.

Gouvernance et traçabilité

Les modifications sont versionnées, examinées via des PR, enregistrées dans Airflow et suivies sur Slack. Aucune écriture directe en production par l’agent.

Évolutivité

Même modèle applicable à de nombreux DAG et bases de données (Snowflake, BigQuery, Databricks, Kafka Schema Registry) grâce à des adaptateurs légers.

Améliorations futures

Cas de test générés par l’IA

Génération automatique de tests de validation pour les correctifs proposés.

Analyse d’impact

Cartographier les tâches, tableaux de bord et rapports affectés par un écart.

Détection des écarts sémantiques

Détecter les changements dans le sens métier tels que :

status

le passage de :

ACTIVE / INACTIVE

à :

A / I

même lorsque les noms de colonnes restent inchangés.

Fusion autonome

Les correctifs à faible risque peuvent être fusionnés automatiquement après vérification — uniquement selon des politiques strictes.

Conclusion

L’écart de schéma est inévitable ; la correction manuelle ne l’est pas. Airflow + PostgreSQL + LangGraph + MCP (DB, GitHub, Slack) transforme le processus de détection et d’affichage en un processus de détection, d’analyse, de création de PR et de notification. À mesure que les plateformes agnitives mûrissent, l’ingénierie des données passe de la simple surveillance des défaillances à leur résolution dans le cadre d’une gouvernance stricte.