पाठ 12 / 25
Branching
@task.branch से run के समय एक रास्ता चुनें और उचित trigger rule से दोबारा जुड़ें।
सिर्फ़ एक रास्ता चलता है
Branch task उस task_id (या IDs की सूची) को लौटाता है जिसका अनुसरण करना है; बाक़ी सीधे downstream रास्ते छोड़े जाते हैं। इसे "क्या यह कार्यदिवस है?" या "क्या नया डेटा आया?" जैसे निर्णयों के लिए उपयोग करें। रास्तों को फिर जोड़ने वाले task को none_failed_min_one_success जैसा trigger rule चाहिए, क्योंकि डिफ़ॉल्ट all_success में उसके upstream रास्तों में से एक छोड़ा जाने पर वह भी छोड़ दिया जाता।
तय करें, इंतज़ार करें, फैलाएँ
Branching, trigger rules, sensors और dynamic mapping DAG को डेटा और बाहरी दुनिया के अनुसार ढलने देते हैं।
कार्यदिवस branch
मैंने यह Airflow 3.0.6 पर airflow dags test report_branch 2026-10-01 से चलाया। 1 अक्टूबर 2026 गुरुवार है, इसलिए weekday_report चला, weekend_summary skipped हुआ, और Airflow ने एक downstream task skipped दर्ज किया। logical_date task context में उपलब्ध है।
from datetime import datetime
from airflow.sdk import dag, task
from airflow.providers.standard.operators.empty import EmptyOperator
@dag(schedule="@daily", start_date=datetime(2026, 9, 1), catchup=False)
def report_branch():
@task.branch
def choose(logical_date=None):
return "weekday_report" if logical_date.weekday() < 5 else "weekend_summary"
weekday = EmptyOperator(task_id="weekday_report")
weekend = EmptyOperator(task_id="weekend_summary")
done = EmptyOperator(task_id="done", trigger_rule="none_failed_min_one_success")
choose() >> [weekday, weekend] >> done
report_branch()स्पष्ट join task रखें
Branches के बाद हमेशा join task रखें और उसे सही trigger rule दें। यह भूलना branch के बाद सब कुछ रहस्यमय ढंग से skipped होने का पुराना कारण है।
त्वरित जाँच: Branch task द्वारा न चुने गए रास्तों का क्या होता है?
- वे फिर भी चलते हैं
- वे skipped होते हैं
- वे विफल होते हैं
- वे DAG से हटाए जाते हैं
Answer
वे skipped होते हैं — न चुनी गई branches उस run के लिए skipped स्थिति पाती हैं।