IAMaîtriser l'IA générative Plan du corpus

02 — Pipelines et orchestration

Un pipeline, ce n'est pas du code qui marche. C'est du code qui redémarre proprement quand il casse à 3 h du matin.

Temps de lecture : 16 min | Niveau : Intermédiaire

Ce que tu sauras faire après

  • Faire concevoir un DAG Airflow à partir d'une description métier, sans code inventé.
  • Générer du code d'ingestion Python solide et pas seulement « qui tourne une fois ».
  • Faire écrire la gestion des erreurs, des réessais et des reprises.
  • Documenter un pipeline existant que tu viens de reprendre.
  • Analyser un échec d'exécution à partir des logs.

1. Ce que l'IA fait bien, et ce qu'elle fait mal, sur un pipeline

L'IA est forte pourL'IA est faible pour
Écrire la structure d'un DAGConnaître tes conventions internes de nommage
Générer du code de chargement répétitifSavoir combien de temps tient ta fenêtre d'API
Proposer une stratégie de réessaiChoisir si une reprise doit être idempotente ou non
Lire une stack trace de 200 lignesSavoir qui prévenir quand ça casse
Écrire la documentation d'un DAG existantDeviner une dépendance non déclarée dans le code

Conclusion pratique : tu délègues la forme, tu gardes la décision.


2. Concevoir un DAG : donner le vocabulaire exact

Airflow propose trois façons de déclarer un DAG, toutes documentées. Si tu ne dis pas laquelle tu veux, le modèle en choisit une au hasard et ton dépôt devient incohérent.

# 1. Gestionnaire de contexte (with)
with DAG(
    dag_id="my_dag_name",
    start_date=datetime.datetime(2021, 1, 1),
    schedule="@daily",
):
    EmptyOperator(task_id="task")

# 2. Constructeur standard
my_dag = DAG(
    dag_id="my_dag_name",
    start_date=datetime.datetime(2021, 1, 1),
    schedule="@daily",
)
EmptyOperator(task_id="task", dag=my_dag)

# 3. Décorateur @dag
@dag(start_date=datetime.datetime(2021, 1, 1), schedule="@daily")
def generate_dag():
    EmptyOperator(task_id="task")

generate_dag()

Les paramètres que tu dois toujours imposer dans ton prompt :

  • dag_id : le nom, selon ta convention.
  • start_date : la date de début.
  • schedule : la fréquence ("@daily", "0 0 * * *", …).
  • default_args : les réglages hérités par toutes les tâches.

La doc Airflow formule cet héritage ainsi :

« you can instead pass default_args to the Dag when you create it, and it will auto-apply them to any operator tied to it. »

Concrètement, si tu écris default_args={"retries": 2} sur le DAG, chaque opérateur du DAG hérite de retries=2. Un opérateur peut toujours redéfinir la valeur pour lui seul.

Les dépendances s'écrivent avec les opérateurs >> et << :

first_task >> [second_task, third_task]

La TaskFlow API

Si ton pipeline est surtout du Python, la doc Airflow recommande la TaskFlow API :

« If you write most of your Dags using plain Python code rather than Operators, then the TaskFlow API will make it much easier to author clean Dags without extra boilerplate, all using the @task decorator. »

@task
def get_ip():
    return my_ip_service.get_main_ip()

@task(multiple_outputs=True)
def compose_email(external_ip):
    return {
        'subject': f'Server connected from {external_ip}',
        'body': f'Your server executing Airflow is connected from the external IP {external_ip}<br>'
    }

email_info = compose_email(get_ip())

La dernière ligne est la clé. Quand tu écris compose_email(get_ip()), Airflow passe la sortie via XCom et enregistre la dépendance tout seul. Tu n'écris pas de >>.

Attention pédagogique : XCom sert à passer de petites valeurs (un chemin de fichier, une date, un identifiant). Pas un DataFrame de 2 Go. C'est une erreur que les modèles font souvent si tu ne le précises pas.

AVANT / APRÈS

AVANT :

Écris-moi un DAG Airflow qui charge les commandes tous les jours.

Tu obtiens du code générique, avec un opérateur inventé, sans start_date, sans gestion d'erreur.

APRÈS :

Écris un DAG Airflow 3.

Contexte :
- Source : API REST paginée https://api.interne/commandes, auth par token dans
  la variable d'environnement API_TOKEN.
- Cible : table BigQuery analytics_raw.commandes, partitionnée sur date_commande.
- Fréquence : quotidienne à 4 h UTC.
- Volumétrie : environ 50 000 commandes par jour.

Contraintes :
- Style : TaskFlow API avec le décorateur @task.
- dag_id : "ingestion_commandes_daily".
- default_args avec retries=3.
- Via XCom, ne fais transiter que des chemins de fichiers, jamais des données.
- Le chargement doit être idempotent : relancer la même date doit donner
  le même résultat, pas des doublons.

Produis : le code, puis une liste <points_a_valider> de tout ce que tu as
supposé et que je dois confirmer.

Pourquoi c'est mieux : le modèle connaît la source, la cible, le volume, le style de code et la convention de nommage. Surtout, il connaît la contrainte d'idempotence. C'est elle qui casse le plus de pipelines en production.


3. Générer du code d'ingestion

Un code d'ingestion sérieux répond à cinq questions. Mets-les dans ton prompt, sinon tu n'auras que la première.

  1. Comment je lis ? Pagination, taille de lot, limite de débit de l'API.
  2. Comment je gère une coupure ? Réessai avec attente croissante, sur quelles erreurs seulement.
  3. Comment je rejoue ? Écriture idempotente : effacer la partition du jour puis réécrire, plutôt qu'ajouter.
  4. Comment je sais que ça a marché ? Nombre de lignes lues, nombre de lignes écrites, et comparaison des deux.
  5. Qu'est-ce que je laisse comme trace ? Journal structuré, avec la date de traitement et le nombre de lignes.
import logging
import time

import requests

LOG = logging.getLogger(__name__)

def lire_page(url: str, token: str, page: int, tentatives: int = 3) -> dict:
    """Lit une page de l'API. Réessaie uniquement sur erreur réseau ou 5xx."""
    for essai in range(1, tentatives + 1):
        try:
            reponse = requests.get(
                url,
                headers={"Authorization": f"Bearer {token}"},
                params={"page": page},
                timeout=30,
            )
            if reponse.status_code >= 500:
                raise requests.HTTPError(f"HTTP {reponse.status_code}")
            reponse.raise_for_status()
            return reponse.json()
        except (requests.ConnectionError, requests.Timeout, requests.HTTPError) as err:
            if essai == tentatives:
                raise
            attente = 2 ** essai
            LOG.warning("Page %s: echec %s, nouvel essai dans %ss", page, err, attente)
            time.sleep(attente)
    raise RuntimeError("inatteignable")

Note bien le détail qui fait la différence : on ne réessaie pas sur une erreur 4xx. Un 401 ou un 422 ne se répareront pas tout seuls en réessayant trois fois ; ils vont juste retarder l'alerte. C'est exactement le genre de nuance que le modèle n'ajoute que si tu la demandes.


4. Erreurs et reprises : le vocabulaire à imposer

Quatre mots à utiliser dans tes prompts. Ils déclenchent le bon comportement.

  • Idempotent : relancer deux fois donne le même état final. Pour une table partitionnée, c'est « je supprime la partition puis je réécris », pas « j'ajoute ».
  • Réessai avec attente croissante (exponential backoff) : 2 s, 4 s, 8 s… Empêche de matraquer une API déjà en difficulté.
  • Échec partiel : 49 000 lignes sur 50 000 sont passées. Que fait-on ? On rejoue tout, ou on reprend au dernier point ? Décide-le toi, et écris-le dans le prompt.
  • File d'attente des rejets (dead letter) : les lignes non traitables vont dans une table à part au lieu de faire tomber le pipeline.

La doc Anthropic sur le pilotage des agents insiste aussi sur la réversibilité. Elle propose explicitement de faire demander confirmation avant les actions destructrices :

« for actions that are hard to reverse, affect shared systems, or could be destructive, ask the user before proceeding. »

Applique-le : quand tu demandes du code qui fait un DROP, un TRUNCATE ou un DELETE, exige que ce soit isolé, commenté, et jamais exécuté automatiquement sans garde-fou.


5. Analyser un échec d'exécution

C'est un des meilleurs usages quotidiens. Le log est dans le prompt, donc le modèle raisonne sur du réel.

Trois règles pour que ça marche :

  1. Colle le log entier, pas juste la dernière ligne. La cause est souvent 80 lignes plus haut.
  2. Ajoute le code de la tâche qui a échoué. Sans lui, le modèle devine.
  3. Ajoute ce qui a changé récemment. « Le DAG tournait hier, on a déployé un changement de schéma ce matin » divise le temps de diagnostic par cinq.

AVANT

J'ai cette erreur Airflow : KeyError: 'montant_ttc'. Pourquoi ?

APRÈS

<log>
[COLLE LE LOG COMPLET DE LA TÂCHE]
</log>

<code_tache>
[COLLE LE CODE DE LA FONCTION QUI A ÉCHOUÉ]
</code_tache>

<contexte>
- DAG : ingestion_commandes_daily, tâche : transform_commandes.
- Ça tournait hier sans erreur.
- Changement depuis : la source a renommé "montant" en "montant_ttc".
</contexte>

Donne-moi, dans cet ordre :
1. <cause_probable> : la cause la plus probable, et pourquoi tu le penses,
   en citant la ligne exacte du log qui l'appuie.
2. <autres_pistes> : 2 autres hypothèses, classées par probabilité.
3. <correctif> : le patch minimal.
4. <prevention> : le test ou le contrôle qui aurait attrapé ça avant la prod.

Ne suppose rien qui ne soit pas dans le log ou le code fournis.

Le « en citant la ligne exacte du log » n'est pas de la décoration. C'est la technique d'ancrage par citation recommandée par Anthropic pour réduire les hallucinations :

« Verify with citations: Make Claude's response auditable by having it cite quotes and sources for each of its claims. »


6. Documenter un pipeline existant

Tu reprends un DAG de 400 lignes. Personne ne sait ce qu'il fait. Le code est la source de vérité, colle-le en entier.

Demande une fiche en six blocs : objectif métier, sources, cibles, planification, dépendances entre tâches, points de fragilité. Le dernier bloc est celui qui a le plus de valeur pour toi : il liste les endroits où ça peut casser.


🧰 Templates à copier-coller

Template A — Concevoir un DAG

Tu es ingénieur data senior, expert Apache Airflow.

Contexte du pipeline :
- Source : [TYPE_SOURCE_API_FICHIER_BASE] — [DETAIL_URL_OU_CHEMIN]
- Authentification : [METHODE_ET_VARIABLE_ENV]
- Cible : [TABLE_CIBLE] sur [ENTREPOT], partitionnée sur [COLONNE_PARTITION]
- Fréquence : [FREQUENCE] à [HEURE] [FUSEAU]
- Volumétrie estimée : [NB_LIGNES_PAR_EXECUTION]
- Fenêtre de traitement acceptable : [DUREE_MAX]

Contraintes techniques :
- Version d'Airflow : [VERSION]
- Style de déclaration : [TaskFlow @task | context manager with DAG | constructeur]
- dag_id : [NOM_DU_DAG]
- default_args : retries=[NB], et [AUTRES_DEFAULT_ARGS]
- XCom : uniquement des métadonnées légères (chemins, dates, compteurs).
- Le chargement doit être idempotent : rejouer une date ne crée pas de doublons.

Produis :
1. <dag> : le code complet.
2. <points_a_valider> : tout ce que tu as supposé et que je dois confirmer.
3. <risques> : les 3 façons les plus probables dont ce DAG cassera en production.

Si une information te manque, demande-la au lieu de l'inventer.

Template B — Code d'ingestion robuste

Écris une fonction Python d'ingestion.

Source : [DESCRIPTION_SOURCE]
Format de réponse : [JSON_CSV_PARQUET] — exemple d'une ligne :
[COLLE_UNE_LIGNE_REELLE_ANONYMISEE]
Cible : [TABLE_OU_FICHIER_CIBLE]

Exigences non négociables :
- Pagination : [DECRIS_LA_PAGINATION]
- Réessai avec attente croissante UNIQUEMENT sur erreurs réseau et 5xx.
  Aucun réessai sur 4xx : lève l'erreur immédiatement.
- Écriture idempotente : [DECRIS_LA_STRATEGIE_EX_SUPPRIMER_LA_PARTITION_PUIS_REECRIRE]
- Journalisation : nombre de lignes lues, nombre de lignes écrites, durée.
- Contrôle final : si lignes_ecrites != lignes_lues, lève une exception explicite.
- Lignes invalides : envoie-les dans [TABLE_DE_REJETS] au lieu de faire tomber le job.

Type-hint toutes les signatures. Pas de dépendance exotique :
utilise seulement [LISTE_DES_LIBRAIRIES_AUTORISEES].

Template C — Analyser un échec

<log>
[COLLE_LE_LOG_COMPLET]
</log>

<code_tache>
[COLLE_LE_CODE_DE_LA_TACHE_EN_ECHEC]
</code_tache>

<contexte>
- DAG : [NOM_DAG] / tâche : [NOM_TACHE]
- Dernière exécution réussie : [DATE]
- Changements depuis : [CE_QUI_A_CHANGE_CODE_SCHEMA_VOLUME_DROITS]
- Le problème est-il reproductible : [OUI_NON_PARFOIS]
</contexte>

Rends-moi :
1. <cause_probable> : la cause la plus probable, en citant la ligne exacte du log.
2. <autres_pistes> : 2 hypothèses alternatives classées par probabilité.
3. <correctif> : le patch minimal, sans refactorisation.
4. <prevention> : le test ou le contrôle qui aurait détecté ça avant la prod.

N'affirme rien qui ne soit pas visible dans le log ou le code fournis.
Si le log est insuffisant, dis-moi quelle information supplémentaire aller chercher.

Template D — Documenter un pipeline repris

<code_dag>
[COLLE_LE_CODE_COMPLET_DU_DAG]
</code_dag>

Écris la fiche de ce pipeline, en français, en six blocs :
1. Objectif métier (3 phrases max, sans jargon).
2. Sources lues : table/API, granularité, fréquence de rafraîchissement.
3. Cibles écrites : table, mode d'écriture (ajout ou remplacement), partitionnement.
4. Planification : fréquence, heure, fuseau, fenêtre de rattrapage.
5. Graphe des tâches : un tableau | tâche | dépend de | ce qu'elle fait |.
6. Points de fragilité : où ça casse le plus probablement, et pourquoi.

N'invente aucune dépendance qui ne soit pas visible dans le code. Si une
dépendance semble implicite (un fichier déposé par un autre job), signale-la
comme "dépendance implicite à confirmer".

A retenir

  • Impose le style de déclaration du DAG (@dag, with DAG, constructeur) : sinon ton dépôt devient incohérent.
  • default_args avec retries est hérité par toutes les tâches, sauf si elles le redéfinissent.
  • Avec la TaskFlow API, les dépendances se déduisent des appels ; ne fais transiter par XCom que des métadonnées légères.
  • Le mot « idempotent » dans ton prompt change radicalement le code produit : exige-le.
  • Pour un échec : log complet + code de la tâche + ce qui a changé, et demande une citation du log à l'appui du diagnostic.

Sources

Corpus personnel de formation · genere le 26/09/2026 · source : 02-pipelines-et-orchestration.md