Planification sensible aux partitions avec Asset Partitions

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Quand ds ne suffit pas

 

  • {{ ds }} lie chaque exécution à une date, utile pour les Dags planifiés dans le temps
  • Mais les Dags déclenchés par un actif ont ds à None : indisponible
  • Le Dag en aval ne sait pas quelle date vient d'être mise à jour
  • Asset Partitions propage une partition_key de l'amont vers l'aval

Exécution planifiée en amont déclenche une exécution déclenchée par actif en aval où ds vaut None

Créer des pipelines de données avec Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    ...

 

  • Chaque exécution planifiée reçoit automatiquement une partition_key
  • Les exécutions manuelles ne sont pas partitionnées sans clé fournie
Créer des pipelines de données avec Airflow

Accéder aux clés de partition

Dans les tâches Python :

@task
def process_data(dag_run=None):
    partition_key = dag_run.partition_key
    print(f"Processing: {partition_key}")

Dans les gabarits SQL :

DELETE FROM daily_summary
WHERE order_date = '{{ dag_run.partition_key[:10] }}';

 

  • Disponible dans chaque tâche d'une exécution partitionnée
Créer des pipelines de données avec Airflow

Événements d'actif partitionnés

from airflow.sdk import dag, task, Asset, CronPartitionTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    @task(outlets=[daily_sales])
    def load_data(**context):
        ...

 

  • Outlet de tâche + CronPartitionTimetable = événement d'actif partitionné
  • L'événement porte la même clé de partition que l'exécution du Dag
  • Indique quelle partition est prête, pas seulement qu'une donnée a changé
Créer des pipelines de données avec Airflow

PartitionedAssetTimetable

from airflow.sdk import dag, task, Asset, PartitionedAssetTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=PartitionedAssetTimetable(assets=daily_sales))
def sales_report():
    @task
    def generate_report(dag_run=None):
        partition_key = dag_run.partition_key
        print(f"Report for: {partition_key}")

 

  • Se déclenche uniquement sur des événements d'actif partitionnés
  • L'exécution en aval hérite de la clé de partition
Créer des pipelines de données avec Airflow

StartOfDayMapper

from airflow.sdk import (
    dag, PartitionedAssetTimetable,
    StartOfDayMapper,
)

@dag(schedule=PartitionedAssetTimetable(
    assets=daily_sales,
    partition_mapper_config={
        daily_sales: StartOfDayMapper()
    },
))
def sales_report():
    ...

 

Avant mappage : 2026-04-23T00:00:00

Après StartOfDayMapper : 2026-04-23

$$

  • Clés temporelles : StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Clés non temporelles : AllowedKeyMapper pour les régions, directions
Créer des pipelines de données avec Airflow

Le flux complet

Flux de partition complet

  • Amont : CronPartitionTimetable ajoute une clé de partition à chaque exécution
  • La tâche d'« outlet » émet un événement d'actif partitionné avec cette clé
  • StartOfDayMapper normalise l'horodatage en date
  • Aval : PartitionedAssetTimetable se déclenche et hérite de la clé mappée
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...