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