Tutoriel4 min de lecture

Introduction à Apache Airflow pour débutants

Ce qu’est Apache Airflow, d’où il vient, à quoi il sert, et les concepts clés d’Airflow 3 à maîtriser avant d’aller plus loin.

AirflowPythonOrchestration

Introduction

Cet article est une introduction à Apache Airflow, pour toute personne qui veut en comprendre les concepts de base. Il présente ce qu’est Airflow, d’où il vient, à quoi il sert, et les notions à maîtriser avant d’aller plus loin.

Apache Airflow est un outil open source pour écrire, planifier et superviser des workflows. Né chez Airbnb en 2014, il a rejoint l’Apache Software Foundation en 2016 et il est depuis développé par la communauté. La dernière version majeure est la 3.x, et cet article en décrit les concepts.

Airflow est un orchestrateur de workflows polyvalent, aux cas d’usage variés : de la gestion d’infrastructure (créer et détruire des ressources selon les besoins) aux pipelines de données (Extract Transform Load, ETL, ou Extract Load Transform, ELT). Il convient à tout processus fait d’étapes exécutées les unes après les autres ou en parallèle.

Un atout majeur d’Airflow est la définition des workflows en code. Les pratiques du développement logiciel, comme le versioning et les pipelines CI/CD, s’appliquent donc aussi aux workflows. Airflow fournit en plus une interface web pour déclencher des workflows manuellement et inspecter leurs exécutions et leurs logs lors du débogage.

Les concepts d’Airflow

Écrire et gérer des workflows est au cœur d’Apache Airflow. Un workflow est un ensemble de tâches exécutées dans un ordre précis. Airflow représente les workflows sous forme de graphes orientés acycliques (Directed Acyclic Graphs, DAG) : orientés parce que les tâches dépendent les unes des autres et s’exécutent dans un ordre donné, acycliques parce qu’aucune boucle n’est autorisée.

Qu’est-ce qu’un DAG ?

Un DAG est le modèle d’un workflow dans Airflow, avec toutes les métadonnées qui s’y rattachent :

  • les tâches du workflow ;
  • les dépendances entre les tâches (l’ordre d’exécution) ;
  • la planification des exécutions ;
  • les dates de début et de fin du workflow.

Un DAG décrit les tâches qui composent le workflow, mais ne fait pas le travail lui-même : c’est le rôle des tâches. Il y a trois façons de déclarer un DAG.

Méthode 1 : gestionnaire de contexte (recommandée)

from datetime import datetime

from airflow.sdk import DAG

with DAG(
    dag_id="my_dag",
    start_date=datetime(2024, 1, 1),
    schedule="@daily",
    catchup=False,
) as dag:
    # Tasks defined here are automatically assigned to this DAG
    pass

Méthode 2 : constructeur de la classe DAG

from datetime import datetime

from airflow.sdk import DAG

dag = DAG(
    dag_id="my_dag",
    start_date=datetime(2024, 1, 1),
    schedule="@daily",
    catchup=False,
)

# Tasks must reference the DAG explicitly
# task = SomeOperator(..., dag=dag)

Méthode 3 : décorateur TaskFlow

from datetime import datetime

from airflow.sdk import dag

@dag(
    dag_id="my_dag",
    start_date=datetime(2024, 1, 1),
    schedule="@daily",
    catchup=False,
)
def my_workflow():
    # Define tasks with the @task decorator here
    pass

my_workflow()

Tâches et opérateurs

Les tâches sont l’unité d’exécution d’un workflow : c’est là que le travail se fait. Une tâche est une instance d’un opérateur, la classe qui définit ce que fait la tâche. Airflow fournit des opérateurs comme BashOperator et PythonOperator, qui exécutent respectivement des commandes Bash et des fonctions Python. Dans Airflow 3, ils viennent du provider standard, installé avec Airflow.

BashOperator et PythonOperator montrent vite leurs limites pour des workflows qui s’intègrent à des services externes, comme des fournisseurs cloud ou des plateformes de données. C’est là que les providers deviennent utiles (voir plus bas).

Méthode 1 : gestionnaire de contexte

from datetime import datetime

from airflow.providers.standard.operators.bash import BashOperator
from airflow.providers.standard.operators.python import PythonOperator
from airflow.sdk import DAG

def my_python_function():
    print("Hello from Python!")

with DAG(
    dag_id="my_dag",
    start_date=datetime(2024, 1, 1),
    schedule="@daily",
) as dag:
    # Tasks are automatically assigned to the DAG
    task1 = BashOperator(
        task_id="bash_task",
        bash_command="echo 'Hello from Bash!'",
    )

    task2 = PythonOperator(
        task_id="python_task",
        python_callable=my_python_function,
    )

    # Define task dependencies
    task1 >> task2

Méthode 2 : constructeur de la classe DAG

from datetime import datetime

from airflow.providers.standard.operators.bash import BashOperator
from airflow.providers.standard.operators.python import PythonOperator
from airflow.sdk import DAG

def my_python_function():
    print("Hello from Python!")

dag = DAG(
    dag_id="my_dag",
    start_date=datetime(2024, 1, 1),
    schedule="@daily",
)

# Tasks must reference the DAG explicitly
task1 = BashOperator(
    task_id="bash_task",
    bash_command="echo 'Hello from Bash!'",
    dag=dag,
)

task2 = PythonOperator(
    task_id="python_task",
    python_callable=my_python_function,
    dag=dag,
)

# Define task dependencies
task1 >> task2

Méthode 3 : décorateur TaskFlow

from datetime import datetime

from airflow.sdk import dag, task

@dag(
    dag_id="my_dag",
    start_date=datetime(2024, 1, 1),
    schedule="@daily",
)
def my_workflow():
    @task
    def extract():
        return {"data": [1, 2, 3]}

    @task
    def transform(data: dict):
        return {"transformed": [x * 2 for x in data["data"]]}

    @task
    def load(data: dict):
        print(f"Loading: {data}")

    # TaskFlow infers dependencies from the function calls
    raw_data = extract()
    transformed_data = transform(raw_data)
    load(transformed_data)

my_workflow()

Providers

Les providers étendent les capacités de base d’Apache Airflow. Ils apportent des opérateurs, des sensors et d’autres briques qui facilitent l’intégration d’Airflow avec des systèmes externes.

Les providers s’installent séparément du cœur d’Airflow, sous forme de paquets Python. Airflow détecte ce qu’un provider apporte au redémarrage qui suit son installation.

L’architecture d’Airflow

Apache Airflow se compose de plusieurs composants. Il peut être déployé en système distribué (recommandé en production) ou sur une seule machine. Voici les composants qui le font fonctionner.

Scheduler

Le scheduler est le cerveau d’Airflow. Il lit les DAG dans la base de métadonnées, déclenche leurs exécutions à échéance et confie leurs tâches à l’executor, qui tourne dans le processus du scheduler.

API server

Sert l’interface web et l’API REST utilisées pour inspecter, déboguer et déclencher les DAG et les tâches.

DAG processor

Parcourt en continu le dossier des DAG, puis analyse les fichiers nouveaux ou modifiés et les sérialise dans la base de métadonnées.

Base de métadonnées

La base de métadonnées est la source de vérité de tout le système : elle stocke l’état des DAG et des tâches. Airflow prend en charge PostgreSQL, MySQL et SQLite (pas en production).

Conclusion

Apache Airflow est un orchestrateur de workflows polyvalent, aux nombreux cas d’usage. Les workflows sont définis en code, et une interface web aide à inspecter, déboguer et déclencher les DAG et les tâches.

Pour les sujets plus avancés, rendez-vous sur la documentation officielle.

Becko Junior Camara

Ingénieur DevOps (Azure & AWS)