RSS →
ARTICLES · AIRFLOW

Airflow + DBT + Snowflake : orchestrer la trinité sans devenir fou

J'ai hérité d'une pipeline Airflow qui orchestrait un projet DBT de 380 modèles sur Snowflake. Le DAG faisait 2 400 lignes, un seul BashOperator qui lançait dbt run en bloc, et un timeout de 45 minutes qui explosait à chaque pic de charge. Quand un modèle plantait, tout le pipeline redémarrait depuis le début. Voici comment je l'ai restructuré — et pourquoi orchestrer DBT avec Airflow n'est pas juste "mettre un BashOperator dans un DAG".

Orchestration Airflow + DBT + Snowflake

Le constat : pourquoi tout le monde se plante au début

Le pattern naïf, tout le monde l'a écrit :

# ⛔ Ce que tout le monde fait au début
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta

with DAG("dbt_pipeline", schedule="0 6 * * *") as dag:
    BashOperator(
        task_id="dbt_run",
        bash_command="cd /opt/dbt && dbt run",
        execution_timeout=timedelta(minutes=45),
    )

Ça marche pour 10 modèles. À 380, c'est l'enfer :

  • Un modèle échoue → tout le DAG est en erreur → pas de retry granulaire
  • Pas de visibilité sur quel modèle a échoué
  • Le timeout de 45 min explose dès qu'un warehouse Snowflake est sous chargé
  • Pas de notion de dépendances entre modèles — Airflow ne sait pas que marts dépend de staging

Le problème : Airflow ne connaît pas la structure de ton projet DBT. Il lance une commande opaque et prie pour que ça passe.

Trois approches pour orchestrer DBT avec Airflow

1. BashOperator — le minimum vital

Pour un petit projet (< 20 modèles), un BashOperator peut suffire. Mais il faut au minimum découper par étape :

# ✅ Découper par étape DBT
dbt_seed = BashOperator(task_id="dbt_seed", bash_command="cd /opt/dbt && dbt seed")
dbt_run_staging = BashOperator(
    task_id="dbt_run_staging",
    bash_command="cd /opt/dbt && dbt run --select staging.*",
)
dbt_run_marts = BashOperator(
    task_id="dbt_run_marts",
    bash_command="cd /opt/dbt && dbt run --select marts.*",
)
dbt_test = BashOperator(task_id="dbt_test", bash_command="cd /opt/dbt && dbt test")

dbt_seed >> dbt_run_staging >> dbt_run_marts >> dbt_test

Limites : tu gères les dépendances à la main. Si un nouveau modèle arrive dans staging, tu n'as rien à changer — mais si quelqu'un ajoute une couche intermediate, il faut modifier le DAG.

2. DbtTaskGroup (Cosmos) — le pattern que j'utilise en production

Cosmos est un framework open-source d'Astronomer qui convertit ton projet DBT en TaskGroup Airflow. Chaque modèle DBT devient une tâche Airflow, avec les dépendances automatiquement résolues depuis le manifest DBT.

# ✅ Cosmos : le DAG se construit tout seul depuis DBT
from cosmos import DbtTaskGroup, ProjectConfig, ProfileConfig, ExecutionConfig
from cosmos.profiles import SnowflakeUserPasswordProfileMapping
from airflow import DAG
from datetime import datetime

profile_config = ProfileConfig(
    profile_name="default",
    target_name="prod",
    profile_mapping=SnowflakeUserPasswordProfileMapping(
        conn_id="snowflake_default",
        profile_args={"database": "ANALYTICS", "warehouse": "COMPUTE_WH"},
    ),
)

project_config = ProjectConfig(
    dbt_project_path="/opt/dbt",
    manifest_path="/opt/dbt/target/manifest.json",
)

with DAG("dbt_snowflake_pipeline", schedule="0 6 * * *",
         start_date=datetime(2026, 1, 1)) as dag:

    dbt_task_group = DbtTaskGroup(
        project_config=project_config,
        profile_config=profile_config,
        execution_config=ExecutionConfig(
            execution_timeout=timedelta(minutes=10),
        ),
        operator_args={
            "install_deps": True,
        },
        default_args={"retries": 2},
    )

Ce que ça change :

  • Chaque modèle DBT est une tâche Airflow individuelle avec son propre retry
  • Les dépendances (ref()) deviennent des dépendances Airflow automatiquement
  • Si stg_orders échoue, seuls les modèles qui en dépendent sont marqués upstream_failed — les autres continuent
  • La visibilité dans l'UI Airflow est parfaite : on voit exactement quel modèle a échoué

3. DbtDag — quand DBT est le seul occupant du DAG

Si ton DAG ne fait que lancer DBT (pas d'autres tâches upstream/downstream), DbtDag est encore plus simple :

from cosmos import DbtDag

dag = DbtDag(
    project_config=project_config,
    profile_config=profile_config,
    schedule="0 6 * * *",
    start_date=datetime(2026, 1, 1),
    dag_id="dbt_snowflake_pipeline",
)

Le piège du timeout : découper ou mourir

Le problème numéro 1 en production : dbt run qui tourne 45 minutes sur un warehouse Snowflake sous-dimensionné.

Solution 1 : Sélectionner seulement ce qui a changé

# ✅ Ne runner que les modèles modifiés depuis le dernier run
dbt_run = BashOperator(
    task_id="dbt_run_modified",
    bash_command="cd /opt/dbt && dbt run --select state:modified+",
)

state:modified+ compare le manifest actuel avec le précédent et ne reconstruit que les modèles modifiés et leurs dépendances. Sur un projet de 380 modèles, ça passe de 47 minutes à 3-5 minutes en routine.

Solution 2 : Cosmos avec timeout par modèle

Avec Cosmos, chaque modèle a son propre timeout. Un staging simple timeout à 2 minutes, un mart lourd à 10 minutes :

execution_config = ExecutionConfig(
    execution_timeout=timedelta(minutes=5),  # timeout par défaut
)

Si un modèle dépasse le timeout, seul celui-ci échoue — les autres ne sont pas affectés.

Gestion des environnements : dev → recette → prod

Le pattern que j'applique systématiquement : un profil DBT par environnement, un DAG Airflow par projet, une variable d'environnement pour le target.

# profiles.yml
default:
  outputs:
    dev:
      type: snowflake
      account: "{{ env_var('SNOWFLAKE_ACCOUNT') }}"
      user: "{{ env_var('SNOWFLAKE_USER') }}"
      password: "{{ env_var('SNOWFLAKE_PASSWORD') }}"
      database: ANALYTICS_DEV
      warehouse: COMPUTE_DEV
      schema: "{{ env_var('DBT_SCHEMA', 'public') }}"
    prod:
      type: snowflake
      account: "{{ env_var('SNOWFLAKE_ACCOUNT') }}"
      user: "{{ env_var('SNOWFLAKE_USER') }}"
      password: "{{ env_var('SNOWFLAKE_PASSWORD') }}"
      database: ANALYTICS_PROD
      warehouse: COMPUTE_PROD
      schema: "{{ env_var('DBT_SCHEMA', 'public') }}"
  target: "{{ env_var('DBT_TARGET', 'dev') }}"

Côté Airflow, le target est injecté via une Variable :

# ✅ L'environnement est piloté par Airflow, pas par DBT
target = Variable.get("dbt_target", default_var="dev")

profile_config = ProfileConfig(
    profile_name="default",
    target_name=target,
    profile_mapping=SnowflakeUserPasswordProfileMapping(
        conn_id=f"snowflake_{target}",
        profile_args={"database": f"ANALYTICS_{target.upper()}"},
    ),
)

La règle d'or : le code DBT est identique en dev et en prod. Seuls le profil et les variables changent.

Le pattern complet que j'utilise en production

from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.models import Variable
from cosmos import DbtTaskGroup, ProjectConfig, ProfileConfig, ExecutionConfig
from cosmos.profiles import SnowflakeUserPasswordProfileMapping
from datetime import datetime, timedelta

# Variables d'environnement
target = Variable.get("dbt_target", default_var="dev")
snowflake_conn = f"snowflake_{target}"

# Configs Cosmos
profile_config = ProfileConfig(
    profile_name="default",
    target_name=target,
    profile_mapping=SnowflakeUserPasswordProfileMapping(
        conn_id=snowflake_conn,
        profile_args={
            "database": f"ANALYTICS_{target.upper()}",
            "warehouse": f"COMPUTE_{target.upper()}",
        },
    ),
)

project_config = ProjectConfig(
    dbt_project_path="/opt/dbt",
    manifest_path="/opt/dbt/target/manifest.json",
)

with DAG(
    "dbt_snowflake_pipeline",
    schedule="0 6 * * *",
    start_date=datetime(2026, 1, 1),
    default_args={"retries": 2, "retry_delay": timedelta(minutes=5)},
    catchup=False,
) as dag:

    # 1. Pré-flight : vérifier les sources Snowflake
    pre_flight = BashOperator(
        task_id="pre_flight_check",
        bash_command="cd /opt/dbt && dbt debug",
    )

    # 2. DBT : seed + run + test en un TaskGroup Cosmos
    dbt_group = DbtTaskGroup(
        group_id="dbt_transformations",
        project_config=project_config,
        profile_config=profile_config,
        execution_config=ExecutionConfig(
            execution_timeout=timedelta(minutes=10),
        ),
        operator_args={"install_deps": True},
    )

    # 3. Post-flight : notifier Slack si succès
    notify_success = BashOperator(
        task_id="notify_success",
        bash_command='echo "Pipeline OK - {{ ti.xcom_pull(task_ids=\'dbt_transformations\') }}"',
        trigger_rule="all_success",
    )

    pre_flight >> dbt_group >> notify_success

Le résultat

Métrique Avant (BashOperator mono) Après (Cosmos)
Lignes de code DAG 2 400 60
Visibilité par modèle Aucune (un seul task) Chaque modèle = une tâche
Retry granularity Tout ou rien Par modèle
Temps pipeline (routine) 47 min 5 min (state:modified+)
Ajout d'un modèle DBT Modifier le DAG Rien (Cosmos lit le manifest)
Timeout qui explose Tout le DAG en erreur Seul le modèle concerné

Orchestrer DBT avec Airflow n'est pas mettre un BashOperator dans un DAG. C'est laisser DBT décrire les dépendances, et Airflow exécuter, monitorer, et retry. Si ton DAG Airflow connaît la structure de ton projet DBT, tu as déjà gagné.

M

Mikael Paulhiout

Lead Tech / Data Architect. Écrit sur Snowflake, DBT, Airflow et l'IA appliquée.

RSS LinkedIn