पाठ 7 / 25

निर्भरताएँ तय करना

क्रम, fan-out और fan-in व्यक्त करने के लिए >>, lists और helper functions उपयोग करें।

कोड में तीर

a >> b का मतलब है "a सफल होने के बाद b चलाओ"। List fan-out या fan-in करती है: extract >> [clean_a, clean_b] >> load extract के बाद दोनों clean tasks समानांतर चलाता है, और load दोनों का इंतज़ार करता है। chain(...) और cross_downstream(...) बड़े graphs में मदद करते हैं। DAG acyclic होना चाहिए: कोई task सीधे या परोक्ष रूप से अपने ऊपर निर्भर नहीं हो सकता, और कोशिश करने पर Airflow import error बताता है।

निष्पादन क्रम, चलाकर

मैंने graph extract → [transform_a, transform_b] → load → notify के लिए छोटा dependency resolver (Airflow ख़ुद नहीं) चलाया। extract ख़त्म होते ही दोनों transforms साथ तैयार होते हैं, और load दोनों का इंतज़ार करता है।

from collections import deque

class Task:
    def __init__(self, name): self.name, self.down = name, []
    def __rshift__(self, other):
        for o in (other if isinstance(other, list) else [other]):
            self.down.append(o)
        return other

extract, t1, t2, load, notify = (Task(n) for n in ("extract", "transform_a", "transform_b", "load", "notify"))
extract >> [t1, t2]; t1 >> load; t2 >> load; load >> notify

def topo(tasks):
    indeg = {t.name: 0 for t in tasks}
    for t in tasks:
        for d in t.down: indeg[d.name] += 1
    q = deque(t for t in tasks if indeg[t.name] == 0); out = []
    while q:
        t = q.popleft(); out.append(t.name)
        for d in t.down:
            indeg[d.name] -= 1
            if indeg[d.name] == 0: q.append(d)
    return out

print(topo([extract, t1, t2, load, notify]))

Output:

['extract', 'transform_a', 'transform_b', 'load', 'notify']

Graphs पठनीय रखें

DAG में दर्जनों tasks हों तो उन्हें TaskGroup से समूहित करें ताकि UI समझने योग्य रहे, या pipeline को assets या triggers से जुड़े कई DAGs में बाँटें।

त्वरित जाँच: `extract >> [clean_a, clean_b] >> load` का क्या मतलब है?

  • दोनों clean tasks extract के बाद चलते हैं, और load दोनों का इंतज़ार करता है
  • load पहले चलता है
  • tasks सिर्फ़ list के क्रम में एक के बाद एक चलते हैं
  • यह syntax error है
Answer

दोनों clean tasks extract के बाद चलते हैं, और load दोनों का इंतज़ार करता है — >> के दाईं ओर list fan-out करती है; बाईं ओर list fan-in करती है।