Human-in-the-loop workflows

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

กรณีการใช้งาน HITL

$$

$$

  • การอนุมัติการตอบกลับจาก AI
  • การกำหนดเส้นทาง support ticket ไปยังทีมที่เหมาะสม
  • การตรวจสอบรายงานคุณภาพข้อมูลก่อนเผยแพร่

ภาพแสดง Human in the loop

การสร้าง Data Pipeline ด้วย Airflow

กลุ่ม HITL operator

$$

Operator กรณีการใช้งาน
ApprovalOperator อนุมัติ/ปฏิเสธ
HITLBranchOperator เลือกงานที่จะรันต่อ
HITLOperator ตัวเลือกและฟอร์มแบบกำหนดเอง
HITLEntryOperator รับค่าจากฟอร์มอย่างเดียว ไม่มีตัวเลือก

$$

  • HITL operator ทุกตัวถูก implement เป็น deferrable operators
การสร้าง Data Pipeline ด้วย Airflow

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"), }, )
  • pipeline จะหยุดรอจนกว่าจะกรอกข้อมูลครบทุกช่อง
การสร้าง Data Pipeline ด้วย Airflow

การเข้าถึงข้อมูลที่ตอบกลับ

@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
การสร้าง Data Pipeline ด้วย Airflow

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"], )

 

  • ผู้ใช้เลือกตัวเลือก ซึ่ง map ไปยัง task ID ที่อยู่ถัดไป
  • เส้นทางที่ไม่ถูกเลือกจะถูกข้าม
การสร้าง Data Pipeline ด้วย Airflow

Required Actions ใน UI

 

  • งาน HITL จะสร้าง Required Action
  • ผู้ตรวจสอบตอบกลับใน task instance view

 

  • HITL operator ทุกตัวเป็นแบบ deferrable
  • จะคืน worker slot ระหว่างรอ
  • Triggerer จัดการการรอแบบ async

HITL ใน Airflow UI

การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...