Plánování s ohledem na Asset Partitions

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Kdy ds nestačí

 

  • {{ ds }} váže každý běh na datum – funguje pro časově naplánované DAGy
  • U assetem spouštěných DAGů je ds nastaveno na None – prostě není k dispozici
  • Downstream DAG nemá jak zjistit, které datum bylo právě aktualizováno
  • Asset Partitions předávají partition_key z upstream do downstream

Naplánovaný upstream běh spouští downstream asset-scheduled běh, kde ds je None

Tvorba datových pipeline s Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Každý naplánovaný běh dostane partition_key automaticky
  • Manuální běhy nejsou rozděleny do partition, pokud klíč není zadán
Tvorba datových pipeline s Airflow

Přístup k partition keys

V Python tasku:

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

V SQL šablonách:

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

 

  • Dostupné v každém tasku v rámci rozděleného běhu
Tvorba datových pipeline s Airflow

Události rozděleného assetu

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 tasku + CronPartitionTimetable = událost rozděleného assetu
  • Událost nese stejný partition key jako běh DAGu
  • Signalizuje, která partition je připravena, nejen že data byla aktualizována
Tvorba datových pipeline s 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}")

 

  • Spouští se pouze při událostech rozděleného assetu
  • Downstream běh zdědí partition key
Tvorba datových pipeline s Airflow

StartOfDayMapper

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

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

 

Před mapováním: 2026-04-23T00:00:00

Po StartOfDayMapper: 2026-04-23

$$

  • Časové klíče: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Nečasové klíče: AllowedKeyMapper pro regiony, oddělení
Tvorba datových pipeline s Airflow

Celý průběh

Celý průběh partition flow

  • Upstream: CronPartitionTimetable přiřadí partition key každému běhu
  • Outlet task vygeneruje událost rozděleného assetu s tímto klíčem
  • StartOfDayMapper normalizuje časové razítko na datum
  • Downstream: PartitionedAssetTimetable se spustí a zdědí mapovaný klíč
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...