Partition-bewuste planning met Asset Partitions

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Als ds niet genoeg is

 

  • {{ ds }} koppelt elke run aan een datum, handig voor tijdgebaseerde, geplande Dags
  • Maar bij asset-getriggerde Dags is ds None en dus niet beschikbaar
  • De downstream Dag weet niet welke datum zojuist is bijgewerkt
  • Asset Partitions geven een partition_key door van upstream naar downstream

Upstream geplande run triggert downstream asset-geplande run waar ds None is

Data-pijplijnen bouwen met Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Elke geplande run krijgt automatisch een partition_key
  • Handmatige runs zijn niet gepartitioneerd tenzij je een key opgeeft
Data-pijplijnen bouwen met Airflow

Partition keys gebruiken

In Python-taken:

@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] }}';

 

  • Beschikbaar in elke taak binnen een gepartitioneerde run
Data-pijplijnen bouwen met Airflow

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

 

  • Taak-outlet + CronPartitionTimetable = gepartitioneerd asset-event
  • Het event draagt dezelfde partition key als de Dag-run
  • Geeft aan welke partition klaar is, niet alleen dat data is bijgewerkt
Data-pijplijnen bouwen met 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}")

 

  • Triggert alleen op gepartitioneerde asset-events
  • Downstream-run erft de partition key
Data-pijplijnen bouwen met Airflow

StartOfDayMapper

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

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

 

Voor mapping: 2026-04-23T00:00:00

Na StartOfDayMapper: 2026-04-23

$$

  • Temporele keys: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Niet-temporele keys: AllowedKeyMapper voor regio's, afdelingen
Data-pijplijnen bouwen met Airflow

De volledige flow

Volledige partition-flow

  • Upstream: CronPartitionTimetable koppelt een partition key aan elke run
  • Outlet-taak zendt een gepartitioneerd asset-event uit met die key
  • StartOfDayMapper normaliseert de timestamp naar een datum
  • Downstream: PartitionedAssetTimetable triggert en erft de gemapte key
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...