Pianificazione con partizioni usando Asset Partitions

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Quando ds non basta

 

  • {{ ds }} lega ogni esecuzione a una data: va bene per Dag programmati nel tempo
  • Ma nei Dag avviati da asset ds è None: non è disponibile
  • Il Dag a valle non sa quale data è stata aggiornata
  • Le Asset Partitions propagano un partition_key da monte a valle

L'esecuzione programmata a monte attiva un'esecuzione asset-programmed a valle dove ds è None

Creare data pipeline con Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Ogni esecuzione programmata riceve automaticamente un partition_key
  • Le esecuzioni manuali non sono partizionate se non fornisci una chiave
Creare data pipeline con Airflow

Accedere alle chiavi di partizione

Nei task Python:

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

Nei template SQL:

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

 

  • Disponibile in ogni task dentro un'esecuzione partizionata
Creare data pipeline con Airflow

Eventi asset partizionati

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):
        ...

 

  • Task outlet + CronPartitionTimetable = evento asset partizionato
  • L'evento porta la stessa chiave di partizione dell'esecuzione del Dag
  • Indica quale partizione è pronta, non solo che i dati sono aggiornati
Creare data pipeline con 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}")

 

  • Si attiva solo su eventi asset partizionati
  • L'esecuzione a valle eredita la chiave di partizione
Creare data pipeline con Airflow

StartOfDayMapper

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

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

 

Prima del mapping: 2026-04-23T00:00:00

Dopo StartOfDayMapper: 2026-04-23

$$

  • Chiavi temporali: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Chiavi non temporali: AllowedKeyMapper per regioni, reparti
Creare data pipeline con Airflow

Il flusso completo

Flusso completo delle partizioni

  • A monte: CronPartitionTimetable associa una chiave di partizione a ogni esecuzione
  • Il task outlet emette un evento asset partizionato con quella chiave
  • StartOfDayMapper normalizza il timestamp a una data
  • A valle: PartitionedAssetTimetable si attiva e eredita la chiave mappata
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...