Mappage dynamique des tâches

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Le problème des tâches codées en dur

@task
def process_file_1():
    transform("/data/file_1.csv")
@task
def process_file_2():
    transform("/data/file_2.csv")
@task
def process_file_3():
    transform("/data/file_3.csv")

Exemple de mappage dynamique des tâches

Créer des pipelines de données avec Airflow

.expand()

@task
def fetch_files():
    return ["/data/file_1.csv", "/data/file_2.csv", "/data/file_3.csv"]

@task(max_active_tis_per_dagrun=2) # control max parallelism def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

Monter en charge à l'exécution

Créer des pipelines de données avec Airflow

Tâches mappées dans l'interface

 

  • Les instances mappées apparaissent comme des instances indexées sous une seule tâche
  • Cliquez dans chacune pour voir les journaux et l'état individuels
  • Le DAG reste lisible même avec des centaines d'instances

Tâches mappées dans l'interface

Créer des pipelines de données avec Airflow

.partial()

@task
def process_file(path: str, output_dir: str):
    transform(path, output_dir)

files = get_files() process_file.partial(output_dir="/out").expand(path=files)

partial et expand

Créer des pipelines de données avec Airflow

Le problème du produit croisé

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# Produit cartésien : 3 x 3 = 9 instances (pas ce qu'on veut!) process.expand(path=files, dest=destinations)

Zipper des listes

Créer des pipelines de données avec Airflow

.zip()

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# Apparié : 3 instances (a->a, b->b, c->c) process.expand_kwargs(files.zip(destinations))

 

  • zip() associe les éléments un à un par position
  • expand_kwargs décompresse chaque paire en arguments nommés
Créer des pipelines de données avec Airflow

Quand utiliser quoi

 

Modèle Cas d'utilisation
expand() Mapper une seule liste
partial() Fixer des paramètres partagés entre instances
zip() Apparier plusieurs listes par position

 

  • Aide à créer des DAG évolutifs
  • Quand les données changent, le nombre de tâches s'ajuste automatiquement à l'exécution
Créer des pipelines de données avec Airflow

Passons à la pratique !

Créer des pipelines de données avec Airflow

Preparing Video For Download...