← points de vue
14 novembre 2023

dbt sous Airflow : une tâche par modèle, pas une par projet

Un projet dbt lancé depuis un orchestrateur tient d’ordinaire dans une seule tâche. Le graphe affiché ne correspond alors pas au projet, et la reprise après erreur rejoue tout. Ce que change le fait de rendre le graphe à l’orchestrateur, et ce que cela coûte.

La façon ordinaire de faire tourner dbt depuis Airflow tient en une ligne : un opérateur bash, la commande dbt run, et voilà le projet dans une tâche. Cela fonctionne, et cela se paie sur trois postes qu’on ne voit qu’en incident. La reprise après erreur rejoue le projet entier alors qu’un seul modèle a cédé. Le graphe de dépendances que dbt connaît reste invisible à l’orchestrateur, qui affiche un carré là où il y a deux cents nœuds. Et les journaux arrivent en un seul bloc, où il faut chercher.

Cosmos, publié par Astronomer, lit le projet dbt et en construit des tâches Airflow : une par modèle, avec les dépendances que dbt a déjà calculées. L’intérêt n’est pas cosmétique : ce que l’orchestrateur voit, il sait le reprendre, le paralléliser et le raconter.

01Le graphe rendu à l’orchestrateur

Trois objets de configuration, et le DAG se déduit du projet. `ProjectConfig` dit où il se trouve, `ProfileConfig` comment s’y connecter, `ExecutionConfig` quel binaire dbt appeler : le dernier compte plus qu’il n’en a l’air, dbt et Airflow ne s’accordant presque jamais sur leurs dépendances.

python
from datetime import datetime

from cosmos import DbtDag, ExecutionConfig, ProfileConfig, ProjectConfig
from cosmos.profiles import AthenaAccessKeyProfileMapping

profile = ProfileConfig(
    profile_name="warehouse",
    target_name="prod",
    # The dbt profile is built from an Airflow connection of type aws:
    # no profiles.yml to deploy, no secret in the repository.
    profile_mapping=AthenaAccessKeyProfileMapping(
        conn_id="aws_warehouse",
        profile_args={"schema": "marts"},
    ),
)

dag = DbtDag(
    dag_id="warehouse_dbt",
    project_config=ProjectConfig("/usr/local/airflow/dags/dbt/warehouse"),
    profile_config=profile,
    # dbt in its own virtual environment: its dependencies and Airflow's do not
    # hold together, and that is the trap that costs you the day.
    execution_config=ExecutionConfig(
        dbt_executable_path="/usr/local/airflow/dbt_venv/bin/dbt",
    ),
    schedule_interval="0 5 * * *",
    start_date=datetime(2023, 11, 1),
    catchup=False,
    default_args={"retries": 2},
)

Le plus souvent le projet dbt n’est pas seul : il suit une ingestion et précède une publication. `DbtTaskGroup` produit alors le même graphe, mais posé au milieu d’un DAG existant, et les dépendances se déclarent comme entre deux tâches ordinaires.

python
from airflow.decorators import dag, task
from cosmos import DbtTaskGroup, ProjectConfig

@dag(schedule="0 5 * * *", start_date=datetime(2023, 11, 1), catchup=False)
def warehouse():
    @task
    def ingest() -> None:
        ...

    transform = DbtTaskGroup(
        group_id="transform",
        project_config=ProjectConfig("/usr/local/airflow/dags/dbt/warehouse"),
        profile_config=profile,
    )

    @task
    def publish() -> None:
        ...

    ingest() >> transform >> publish()

warehouse()
02Le profil, là où Athena complique

Athena n’est pas une base au sens habituel : le calcul est chez AWS, les données sont dans S3, et il faut un emplacement d’écriture pour les résultats intermédiaires. Un profil dbt-athena porte donc plus de champs qu’un profil Postgres, et deux d’entre eux, le répertoire de travail et le groupe de travail, n’ont pas d’équivalent ailleurs.

Le mapping de Cosmos lit ces champs dans le bloc `extra` d’une connexion Airflow de type `aws`, et prend les identifiants par le mécanisme habituel du fournisseur Amazon : ce qui veut dire qu’un rôle IAM fonctionne aussi bien qu’une clé, et vaut mieux.

json
{
  "region_name": "eu-west-1",
  "database": "awsdatacatalog",
  "schema": "marts",
  "s3_staging_dir": "s3://warehouse-athena/results/",
  "work_group": "dbt"
}

Deux avertissements, tirés de l’usage plutôt que de la documentation. Le premier : donner à dbt son propre groupe de travail Athena, comme ci-dessus, sépare sa consommation du reste et permet de lui poser une limite d’octets scannés. Sans quoi une jointure maladroite se découvre sur la facture. Le second : l’adaptateur s’installe sous le nom `dbt-athena-community`, et non `dbt-athena`, ce qui est la première demi-heure perdue par tout le monde.

03Ce que ça coûte

Le graphe doit être connu au moment où Airflow analyse le fichier, c’est-à-dire toutes les trente secondes par défaut. Laisser Cosmos appeler dbt à chaque analyse pour découvrir les modèles est la configuration qui se met en place toute seule, et la plus mauvaise : le planificateur ralentit à mesure que le projet grandit. La parade tient en une phrase : produire le manifeste au moment de la construction de l’image, et le donner à Cosmos plutôt que de le lui faire recalculer.

Le second coût est moins évident : deux cents tâches au lieu d’une, ce sont deux cents lignes dans la base de métadonnées à chaque exécution, et un tableau de bord qu’on ne lit plus d’un coup d’œil. Sur un projet de vingt modèles le gain est net ; au-delà de quelques centaines, il faut regrouper par dossier plutôt que par modèle, et Cosmos sait le faire.

04Ce que je laisse de côté

Je n’ai pas fait tourner ceci sur Kubernetes, où l’exécution se fait dans un conteneur par tâche et où l’arbitrage change complètement : le graphe détaillé coûte alors un démarrage de conteneur par modèle, et l’économie de la reprise fine peut s’inverser. Je ne dis rien non plus des tests dbt exposés en tâches distinctes, que Cosmos permet : je m’en méfie, parce qu’un test rouge doit arrêter la chaîne, et le rendre indépendant rend cette règle négociable.

Une dernière chose, datée. Le mapping Athena n’a qu’un mois au moment où j’écris : il est arrivé dans la version 1.2.0, le 13 octobre 2023, et cinq versions ont suivi en six semaines. C’est le rythme d’un projet qui n’a pas fini de bouger, et la raison pour laquelle le code ci-dessus est à lire comme un état plutôt que comme une recette.