टास्क के बीच डेटा साझा करना

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

Volker Janz

Senior Developer Advocate at Astronomer

XCom क्या है?

XCom push pull flow

 

  • XCom = टास्कों के बीच क्रॉस-कम्युनिकेशन
  • वैल्यूज़ metadata database में स्टोर होती हैं
  • क्लासिक तरीका: स्पष्ट xcom_push() और xcom_pull()
Airflow के साथ Data Pipelines बनाना

XComArg: TaskFlow का तरीका

 

  • किसी @task का रिटर्न वैल्यू XComArg बन जाता है
  • XComArg एक lazy reference है, असली डेटा नहीं
  • इसे दूसरी फंक्शन में पास करने से dependency बनती है
  • रनटाइम पर Airflow इसे XCom से पुल करके resolve करता है

XComArgs in TaskFlow

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

कोड में XComArg

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # dependency + data
Airflow के साथ Data Pipelines बनाना

क्लासिक और TaskFlow को मिलाएँ

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Airflow के साथ Data Pipelines बनाना

XCom की सीमाएँ

XCom size limits by backend

 

  • अधिकतम साइज आपके database backend पर निर्भर है
  • वैल्यू JSON serializable होनी चाहिए (dicts, lists, strings, numbers)
  • हर वैल्यू डेटाबेस पर लोड बढ़ाती है
Airflow के साथ Data Pipelines बनाना

कस्टम XCom बैकएंड्स

Custom XCom backend routing

  • वैल्यूज़ को किसी object storage पर रूट करें
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: सिर्फ़ threshold पार होने पर ही XComs को object storage में स्टोर करें
Airflow के साथ Data Pipelines बनाना

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

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

Preparing Video For Download...