Planification sensible aux données avec les Assets

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Le problème de synchronisation

 

  • Le pipeline A charge des données brutes depuis une API
  • Le pipeline B crée un tableau de bord à partir de ces données
  • Si les deux tournent via cron, B peut s'exécuter avant la fin de A
  • Ajouter une marge de temps est fragile

Le problème de synchronisation

Créer des pipelines de données avec Airflow

Qu'est-ce qu'un Asset ?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Une référence à des données, identifiée par un nom unique
  • Une URI peut être associée à l'asset lorsqu'il représente des données concrètes
  • Peut être un fichier, une table de base de données ou toute donnée
Créer des pipelines de données avec Airflow

Producteurs et consommateurs

  • Les tâches peuvent mettre à jour des assets, ce qui crée un événement d'asset
  • Les DAGs peuvent être déclenchés lors des mises à jour via un horaire d'asset

 

Flux d'assets

Créer des pipelines de données avec Airflow

Producteur : signaler avec outlets

sales_data = Asset("s3://bucket/sales/daily.parquet")

@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...

 

  • En cas de succès, Airflow enregistre la mise à jour de l'asset
  • Le producteur n'a pas à connaître les consommateurs
Créer des pipelines de données avec Airflow

Consommateur : planifier selon un asset

sales_data = Asset("s3://bucket/sales/daily.parquet")

@dag(schedule=[sales_data]) def build_dashboard(): ...

 

  • Aucune expression cron, seulement la référence d'asset
  • Airflow déclenche le consommateur automatiquement après la mise à jour
Créer des pipelines de données avec Airflow

Planification conditionnelle

sales = Asset("s3://bucket/sales.parquet")
inventory = Asset("s3://bucket/inventory.parquet")

# Attendre la mise à jour des DEUX assets @dag(schedule=(sales & inventory)) def full_report(): ...
# Déclencher dès qu'AU MOINS UN asset est mis à jour @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Créer des pipelines de données avec Airflow

Vérifier les assets avec la CLI

$ airflow assets list
name                                  | group    | uri
s3://data-lake/sales/daily.csv        | asset    | s3://data-lake/sales/daily.csv
$ airflow assets details --name "s3://data-lake/sales/daily.csv"
property_name    | property_value
=================+=====================                                         
name             | s3://data-lake/sales/daily.csv
                ...
updated_at       | 2026-04-14T10:00:50.820041Z
  • airflow assets list affiche tous les assets enregistrés
  • airflow assets details affiche les métadonnées, dont updated_at
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...