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

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
martsdépend destaging
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é.