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.
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
passMé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 >> task2Mé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 >> task2Mé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.