Asset Partitions ile bölüme duyarlı zamanlama

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

ds yeterli olmadığında

 

  • {{ ds }} her çalıştırmayı bir tarihe bağlar; bu, zamana göre zamanlanmış Dag'ler için işe yarar
  • Ama varlık tetiklemeli Dag'lerde ds None olur, yani yoktur
  • Aşağı akıştaki Dag az önce hangi tarihin güncellendiğini bilemez
  • Asset Partitions, yukarı akıştan aşağı akışa bir partition_key iletir

Yukarı akış zamanlanmış çalıştırma, ds'in None olduğu aşağı akış varlık-zamanlamalı çalıştırmayı tetikler

Airflow ile Veri İş Hatları Oluşturma

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Her zamanlanmış çalıştırma otomatik olarak bir partition_key alır
  • Manuel çalıştırmalar, anahtar verilmedikçe bölümlenmez
Airflow ile Veri İş Hatları Oluşturma

Partition key'lere erişim

Python görevlerinde:

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

SQL şablonlarında:

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

 

  • Bölümlü bir çalıştırmadaki her görevde kullanılabilir
Airflow ile Veri İş Hatları Oluşturma

Bölümlü varlık olayları

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

 

  • Görev outlet + CronPartitionTimetable = bölümlü varlık olayı
  • Olay, Dag çalıştırmasıyla aynı partition key'i taşır
  • Sadece veri güncellendiğini değil, hangi bölümün hazır olduğunu belirtir
Airflow ile Veri İş Hatları Oluşturma

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

 

  • Yalnızca bölümlü varlık olaylarında tetiklenir
  • Aşağı akış çalıştırması partition key'i devralır
Airflow ile Veri İş Hatları Oluşturma

StartOfDayMapper

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

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

 

Eşlemeden önce: 2026-04-23T00:00:00

StartOfDayMapper sonrası: 2026-04-23

$$

  • Zamansal anahtarlar: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Zamansal olmayan anahtarlar: bölgeler, departmanlar için AllowedKeyMapper
Airflow ile Veri İş Hatları Oluşturma

Tam akış

Tam bölüm akışı

  • Yukarı akış: CronPartitionTimetable her çalıştırmaya partition key ekler
  • Outlet görevi bu anahtarla bölümlü varlık olayı yayar
  • StartOfDayMapper zaman damgasını bir tarihe normalleştirir
  • Aşağı akış: PartitionedAssetTimetable tetiklenir ve eşlenen anahtarı devralır
Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...