Planificarea conștientă de partiții cu Asset Partitions

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Când ds nu este suficient

 

  • {{ ds }} leagă fiecare rulare de o dată, util pentru Dag-urile programate
  • Dar Dag-urile declanșate de assets au ds setat la None – valoarea nu e disponibilă
  • Dag-ul din aval nu știe ce dată tocmai a fost actualizată
  • Asset Partitions propagă un partition_key din amonte în aval

Rularea programată upstream declanșează rularea downstream programată de asset, unde ds este None

Construirea pipeline-urilor de date cu Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Fiecare rulare programată primește automat un partition_key
  • Rulările manuale nu sunt partiționare dacă nu se furnizează cheia
Construirea pipeline-urilor de date cu Airflow

Accesarea cheilor de partiție

În taskuri Python:

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

În șabloane SQL:

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

 

  • Disponibil în fiecare task dintr-o rulare partiționată
Construirea pipeline-urilor de date cu Airflow

Evenimente de asset partiționat

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 task + CronPartitionTimetable = eveniment de asset partiționat
  • Evenimentul poartă aceeași cheie de partiție ca rularea Dag-ului
  • Semnalează ce partiție e gata, nu doar că datele s-au actualizat
Construirea pipeline-urilor de date cu 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 declanșează doar la evenimente de asset partiționat
  • Rularea din aval moștenește cheia de partiție
Construirea pipeline-urilor de date cu Airflow

StartOfDayMapper

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

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

 

Înainte de mapare: 2026-04-23T00:00:00

După StartOfDayMapper: 2026-04-23

$$

  • Chei temporale: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Chei non-temporale: AllowedKeyMapper pentru regiuni, departamente
Construirea pipeline-urilor de date cu Airflow

Fluxul complet

Fluxul complet al partiției

  • Upstream: CronPartitionTimetable atașează cheia de partiție fiecărei rulări
  • Taskul outlet emite un eveniment de asset partiționat cu acea cheie
  • StartOfDayMapper normalizează timestamp-ul la o dată
  • Downstream: PartitionedAssetTimetable se declanșează și moștenește cheia mapată
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...