Human-in-the-loop वर्कफ़्लो

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

Volker Janz

Senior Developer Advocate at Astronomer

HITL उपयोग के मामले

$$

$$

  • AI द्वारा जनरेटेड उत्तर को अप्रूव करना
  • सपोर्ट टिकट को सही टीम तक रूट करना
  • पब्लिश करने से पहले डेटा क्वालिटी रिपोर्ट रिव्यू करना

लूप में मानव का दृश्य

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

HITL ऑपरेटर परिवार

$$

ऑपरेटर उपयोग का मामला
ApprovalOperator बाइनरी approve/reject
HITLBranchOperator अगला कौन-से टास्क चलें चुनें
HITLOperator कस्टम options + form
HITLEntryOperator केवल form input, कोई options नहीं

$$

  • सभी HITL ऑपरेटर deferrable operators के रूप में इम्प्लीमेंट किए गए हैं
Airflow के साथ Data Pipelines बनाना

HITLEntryOperator

from airflow.providers.standard.operators.hitl import HITLEntryOperator
from airflow.sdk import Param

review = HITLEntryOperator( task_id="human_review", subject="Review AI draft response", body="{{ ti.xcom_pull(task_ids='generate_draft') }}", params={ "feedback": Param("", type="string"), "urgency": Param("p3", type="string"), }, )
  • पाइपलाइन तब तक रुकी रहती है जब तक फ़ील्ड भरे नहीं जाते
Airflow के साथ Data Pipelines बनाना

रिस्पॉन्स एक्सेस करना

@task
def send_response(hitl_output):
    feedback = hitl_output["params_input"]["feedback"]
    urgency = hitl_output["params_input"]["urgency"]
    print(f"Sending response (urgency: {urgency}): {feedback}")

draft = generate_draft() chain(draft, review) send_response(review.output)

 

  • फॉर्म इनपुट params_input में होता है
  • ऑप्शंस (यदि हों) chosen_options में होते हैं
Airflow के साथ Data Pipelines बनाना

HITLBranchOperator

from airflow.providers.standard.operators.hitl import HITLBranchOperator

route = HITLBranchOperator( task_id="route_complaint", subject="Route customer complaint", body="A customer reported an issue. Choose the responsible team.", options=["Billing Issue", "General Inquiry", "Technical Issue"], options_mapping={ "Billing Issue": "handle_billing", "Technical Issue": "handle_technical", "General Inquiry": "handle_general", }, defaults=["General Inquiry"], )

 

  • इंसान एक ऑप्शन चुनता है, जिसे एक downstream task ID से मैप किया जाता है
  • न चुने गए पाथ स्किप हो जाते हैं
Airflow के साथ Data Pipelines बनाना

UI में Required Actions

 

  • HITL टास्क एक Required Action बनाते हैं
  • रिव्यूअर task instance view में जवाब देते हैं

 

  • सभी HITL ऑपरेटर deferrable हैं
  • इंतज़ार करते समय वे worker slot रिलीज़ करते हैं
  • Triggerer async wait संभालता है

Airflow UI में HITL

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

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

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

Preparing Video For Download...