Partitionsbewusstes Scheduling mit Asset Partitions

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Wenn ds nicht reicht

 

  • {{ ds }} verknüpft jeden Run mit einem Datum – passt für zeitlich geplante Dags
  • Bei asset-getriggerten Dags ist ds None, also nicht verfügbar
  • Der Downstream-Dag weiß nicht, welches Datum gerade aktualisiert wurde
  • Asset Partitions propagieren einen partition_key von Upstream zu Downstream

Upstream-geplanter Run triggert Downstream-Asset-Run, bei dem ds None ist

Data-Pipelines mit Airflow aufbauen

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Jeder geplante Run erhält automatisch einen partition_key
  • Manuelle Runs sind nicht partitioniert, außer ein Key wird angegeben
Data-Pipelines mit Airflow aufbauen

Auf Partition-Keys zugreifen

In Python-Tasks:

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

In SQL-Templates:

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

 

  • In jedem Task eines partitionierten Runs verfügbar
Data-Pipelines mit Airflow aufbauen

Partitionierte Asset-Events

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 = partitioniertes Asset-Event
  • Das Event trägt denselben Partition-Key wie der Dag-Run
  • Signalisiert, welche Partition bereit ist, nicht nur dass Daten aktualisiert wurden
Data-Pipelines mit Airflow aufbauen

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}")

 

  • Triggert nur auf partitionierte Asset-Events
  • Downstream-Run erbt den Partition-Key
Data-Pipelines mit Airflow aufbauen

StartOfDayMapper

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

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

 

Vor dem Mapping: 2026-04-23T00:00:00

Nach StartOfDayMapper: 2026-04-23

$$

  • Zeitliche Keys: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Nicht-zeitliche Keys: AllowedKeyMapper für Regionen, Abteilungen
Data-Pipelines mit Airflow aufbauen

Der komplette Flow

Kompletter Partition-Flow

  • Upstream: CronPartitionTimetable hängt jeden Run einen Partition-Key an
  • Outlet-Task emittiert ein partitioniertes Asset-Event mit diesem Key
  • StartOfDayMapper normalisiert den Zeitstempel auf ein Datum
  • Downstream: PartitionedAssetTimetable triggert und erbt den gemappten Key
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...