Asset Partitions के साथ partition-aware scheduling

Airflow के साथ Data Pipelines बनाना

Volker Janz

Senior Developer Advocate at Astronomer

जब ds काफ़ी नहीं होता

 

  • {{ ds }} हर रन को एक तारीख से जोड़ता है, जो time-scheduled Dags के लिए ठीक काम करता है
  • लेकिन asset-triggered Dags में ds None होता है, यानी उपलब्ध नहीं
  • डाउनस्ट्रीम Dag को पता नहीं चलता कि अभी कौन-सी तारीख अपडेट हुई है
  • Asset Partitions अपस्ट्रीम से डाउनस्ट्रीम तक partition_key पास करते हैं

अपस्ट्रीम scheduled रन डाउनस्ट्रीम asset-scheduled रन को ट्रिगर करता है जहाँ ds None है

Airflow के साथ Data Pipelines बनाना

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • हर scheduled run को partition_key अपने-आप मिलता है
  • मैनुअल रन partitioned नहीं होते जब तक key न दी जाए
Airflow के साथ Data Pipelines बनाना

Partition keys एक्सेस करना

Python tasks में:

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

SQL templates में:

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

 

  • Partitioned रन में हर task के अंदर उपलब्ध
Airflow के साथ Data Pipelines बनाना

Partitioned 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 = partitioned asset event
  • इवेंट Dag रन वाला same partition key लेकर चलता है
  • यह बताता है कि कौन-सा partition तैयार है, न कि बस data अपडेट हुआ
Airflow के साथ Data Pipelines बनाना

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

 

  • ट्रिगर सिर्फ़ partitioned asset events पर होता है
  • डाउनस्ट्रीम रन partition key इनहेरिट करता है
Airflow के साथ Data Pipelines बनाना

StartOfDayMapper

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

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

 

मैपिंग से पहले: 2026-04-23T00:00:00

StartOfDayMapper के बाद: 2026-04-23

$$

  • Temporal keys: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Non-temporal keys: क्षेत्रों, विभागों के लिए AllowedKeyMapper
Airflow के साथ Data Pipelines बनाना

पूरा फ्लो

पूरी partition फ्लो

  • अपस्ट्रीम: CronPartitionTimetable हर रन से partition key जोड़ता है
  • Outlet task उसी key के साथ partitioned asset event भेजता है
  • StartOfDayMapper timestamp को date में normalize करता है
  • डाउनस्ट्रीम: PartitionedAssetTimetable ट्रिगर होता है और mapped key इनहेरिट करता है
Airflow के साथ Data Pipelines बनाना

अभ्यास करते हैं!

Airflow के साथ Data Pipelines बनाना

Preparing Video For Download...