Harmonogramowanie z uwzględnieniem partycji

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Kiedy ds nie wystarcza

 

  • {{ ds }} wiąże każde uruchomienie z datą — działa to w DAG-ach harmonogramowanych czasowo
  • W DAG-ach wyzwalanych przez zasoby ds ma wartość None — po prostu nie jest dostępny
  • Downstream DAG nie wie, która data właśnie została zaktualizowana
  • Partycje zasobów przekazują partition_key z upstream do downstream

Zaplanowane uruchomienie upstream wyzwala downstream DAG wyzwalany przez zasób, gdzie ds ma wartość None

Budowanie potoków danych z Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • Każde zaplanowane uruchomienie automatycznie otrzymuje partition_key
  • Uruchomienia ręczne nie są partycjonowane, chyba że podasz klucz
Budowanie potoków danych z Airflow

Dostęp do kluczy partycji

W zadaniach Python:

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

W szablonach SQL:

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

 

  • Dostępny w każdym zadaniu w ramach partycjonowanego uruchomienia
Budowanie potoków danych z Airflow

Zdarzenia partycjonowanych zasobów

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 zadania + CronPartitionTimetable = zdarzenie partycjonowanego zasobu
  • Zdarzenie niesie ten sam klucz partycji co uruchomienie DAG-a
  • Sygnalizuje, która partycja jest gotowa, a nie tylko że dane zostały zaktualizowane
Budowanie potoków danych z 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}")

 

  • Wyzwala się wyłącznie na zdarzeniach partycjonowanych zasobów
  • Uruchomienie downstream dziedziczy klucz partycji
Budowanie potoków danych z Airflow

StartOfDayMapper

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

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

 

Przed mapowaniem: 2026-04-23T00:00:00

Po zastosowaniu StartOfDayMapper: 2026-04-23

$$

  • Klucze czasowe: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Klucze nieczasowe: AllowedKeyMapper dla regionów i działów
Budowanie potoków danych z Airflow

Pełny przepływ

Pełny przepływ partycji

  • Upstream: CronPartitionTimetable przypisuje klucz partycji do każdego uruchomienia
  • Zadanie outlet emituje zdarzenie partycjonowanego zasobu z tym kluczem
  • StartOfDayMapper normalizuje znacznik czasu do daty
  • Downstream: PartitionedAssetTimetable wyzwala się i dziedziczy zmapowany klucz
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...