एसेट्स के साथ डेटा-अवेयर शेड्यूलिंग

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

Volker Janz

Senior Developer Advocate at Astronomer

टाइमिंग समस्या

 

  • पाइपलाइन A एक API से कच्चा डेटा लोड करती है
  • पाइपलाइन B उसी डेटा से एक डैशबोर्ड बनाती है
  • अगर दोनों cron पर चलें, तो B, A के खत्म होने से पहले चल सकती है
  • समय बफ़र जोड़ना नाज़ुक होता है

टाइमिंग समस्या

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

Asset क्या है?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • डेटा के किसी हिस्से का एक रेफरेंस, जिसे एक यूनिक नाम से पहचाना जाता है
  • जब यह कोई ठोस डेटा एंटीटी हो, तो एसेट पर एक URI जोड़ा जा सकता है
  • यह एक फ़ाइल, डेटाबेस टेबल, या कोई भी डेटा हो सकता है
Airflow के साथ Data Pipelines बनाना

प्रोड्यूसर और कंज़्यूमर

  • टास्क एसेट्स को अपडेट कर सकते हैं, जिससे एक एसेट इवेंट बनता है
  • DAGs को एसेट शेड्यूल के साथ एसेट अपडेट पर ट्रिगर किया जा सकता है

 

एसेट फ्लो

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

Producer: outlets से संकेत देना

sales_data = Asset("s3://bucket/sales/daily.parquet")

@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...

 

  • सफलता पर, Airflow रिकॉर्ड करता है कि एसेट अपडेट हुआ
  • प्रोड्यूसर को किसी भी कंज़्यूमर के बारे में जानने की ज़रूरत नहीं होती
Airflow के साथ Data Pipelines बनाना

Consumer: किसी एसेट पर शेड्यूल करना

sales_data = Asset("s3://bucket/sales/daily.parquet")

@dag(schedule=[sales_data]) def build_dashboard(): ...

 

  • कोई cron एक्सप्रेशन नहीं, सिर्फ एसेट रेफरेंस
  • अपडेट के बाद Airflow कंज़्यूमर को अपने-आप ट्रिगर करता है
Airflow के साथ Data Pipelines बनाना

कंडीशनल शेड्यूलिंग

sales = Asset("s3://bucket/sales.parquet")
inventory = Asset("s3://bucket/inventory.parquet")

# दोनों एसेट्स के अपडेट होने तक प्रतीक्षा @dag(schedule=(sales & inventory)) def full_report(): ...
# जब कोई भी एसेट अपडेट हो, ट्रिगर करें @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Airflow के साथ Data Pipelines बनाना

CLI से एसेट्स वेरिफाई करना

$ airflow assets list
name                                  | group    | uri
s3://data-lake/sales/daily.csv        | asset    | s3://data-lake/sales/daily.csv
$ airflow assets details --name "s3://data-lake/sales/daily.csv"
property_name    | property_value
=================+=====================                                         
name             | s3://data-lake/sales/daily.csv
                ...
updated_at       | 2026-04-14T10:00:50.820041Z
  • airflow assets list सभी रजिस्टर्ड एसेट्स दिखाता है
  • airflow assets details में updated_at सहित मेटाडेटा दिखता है
Airflow के साथ Data Pipelines बनाना

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

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

Preparing Video For Download...