人機協作工作流程

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

HITL 應用情境

$$

$$

  • 核准 AI 產生的回覆
  • 將支援單導向正確團隊
  • 發佈前審查資料品質報告

人機協作示意

使用 Airflow 建置資料管線

HITL 運算元家族

$$

Operator Use case
ApprovalOperator 二元 核准/拒絕
HITLBranchOperator 選擇下一步要執行的工作
HITLOperator 自訂選項表單
HITLEntryOperator 表單輸入,無選項

$$

  • 所有 HITL operators 都是以可延後(deferrable)運算元實作
使用 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"), }, )
  • 管線會暫停,直到欄位填完為止
使用 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
使用 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"], )

 

  • 由人工選擇選項,對應到下游工作 ID
  • 未選路徑會被略過
使用 Airflow 建置資料管線

UI 中的必要動作

 

  • HITL 工作會建立必要動作(Required Action)
  • 審查者在工作實例檢視中回覆

 

  • 所有 HITL 運算元皆為可延後
  • 等待期間會釋放 worker 名額
  • Triggerer 處理非同步等待

Airflow UI 中的 HITL

使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...