Partitionsmedveten schemaläggning med Asset Partitions

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

När ds inte räcker

 

  • {{ ds }} kopplar varje körning till ett datum, vilket fungerar för tidsschemalagda DAG:ar
  • Men asset-triggade DAG:ar har ds satt till None – värdet är helt enkelt inte tillgängligt
  • Det nedströms DAG:et har inget sätt att veta vilket datum som just uppdaterades
  • Asset Partitions sprider en partition_key från uppströms till nedströms

Uppströms schemalagd körning triggar nedströms asset-schemalagd körning där ds är None

Bygg datapipelines med Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Varje schemalagd körning får automatiskt en partition_key
  • Manuella körningar är inte partitionerade om ingen nyckel anges
Bygg datapipelines med Airflow

Åtkomst till partitionsnycklar

I Python-tasks:

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

I SQL-mallar:

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

 

  • Tillgänglig i varje task inom en partitionerad körning
Bygg datapipelines med Airflow

Partitionerade asset-händelser

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 = partitionerad asset-händelse
  • Händelsen bär samma partitionsnyckel som DAG-körningen
  • Signalerar vilken partition som är klar, inte bara att data uppdaterats
Bygg datapipelines med 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}")

 

  • Triggas bara vid partitionerade asset-händelser
  • Nedströms körning ärver partitionsnyckeln
Bygg datapipelines med Airflow

StartOfDayMapper

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

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

 

Före mappning: 2026-04-23T00:00:00

Efter StartOfDayMapper: 2026-04-23

$$

  • Temporala nycklar: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Icke-temporala nycklar: AllowedKeyMapper för regioner, avdelningar
Bygg datapipelines med Airflow

Det fullständiga flödet

Fullständigt partitionsflöde

  • Uppströms: CronPartitionTimetable kopplar en partitionsnyckel till varje körning
  • Outlet-task sänder ut en partitionerad asset-händelse med den nyckeln
  • StartOfDayMapper normaliserar tidsstämpeln till ett datum
  • Nedströms: PartitionedAssetTimetable triggas och ärver den mappade nyckeln
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...