Quy trình human-in-the-loop

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Tình huống dùng HITL

$$

$$

  • Phê duyệt một phản hồi do AI tạo
  • Chuyển ticket hỗ trợ đến đúng nhóm
  • Rà soát báo cáo chất lượng dữ liệu trước khi xuất bản

Minh họa human-in-the-loop

Xây dựng Data Pipeline với Airflow

Nhóm operator HITL

$$

Operator Tình huống dùng
ApprovalOperator Nhị phân approve/reject
HITLBranchOperator Chọn nhiệm vụ chạy tiếp theo
HITLOperator Tùy chọn + biểu mẫu tùy chỉnh
HITLEntryOperator Chỉ nhập biểu mẫu, không có tùy chọn

$$

  • Mọi operator HITL đều là operator hoãn thực thi (deferrable)
Xây dựng Data Pipeline với 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 sẽ tạm dừng cho đến khi các trường được điền
Xây dựng Data Pipeline với Airflow

Truy cập phản hồi

@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)

 

  • Dữ liệu nhập từ biểu mẫu nằm trong params_input
  • Các tùy chọn (nếu có) nằm trong chosen_options
Xây dựng Data Pipeline với 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"], )

 

  • Con người chọn một tùy chọn, được ánh xạ tới ID nhiệm vụ hạ lưu
  • Các nhánh không được chọn sẽ bị bỏ qua
Xây dựng Data Pipeline với Airflow

Required Actions trong giao diện

 

  • Nhiệm vụ HITL tạo một Required Action
  • Người rà soát phản hồi trong chế độ xem task instance

 

  • Tất cả operator HITL đều có thể hoãn thực thi
  • Chúng giải phóng slot worker trong lúc chờ
  • Triggerer xử lý việc chờ bất đồng bộ

HITL trong giao diện Airflow

Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...