人机协同工作流

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

HITL 用例

$$

$$

  • 审批 AI 生成的回复
  • 将支持工单路由到合适团队
  • 发布前审阅数据质量报告

人机协同示意图

使用 Airflow 构建数据流水线

HITL 运算符家族

$$

Operator 用例
ApprovalOperator 二选一:批准/拒绝
HITLBranchOperator 选择接下来运行的任务
HITLOperator 自定义选项 + 表单
HITLEntryOperator 表单输入,无选项

$$

  • 所有 HITL 运算符均实现为可延迟运算符
使用 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 构建数据流水线

界面中的 Required Actions

 

  • HITL 任务会创建一个 Required Action
  • 审阅者在任务实例视图中作答

 

  • 所有 HITL 运算符都是可延迟
  • 等待期间会释放 worker 槽位
  • Triggerer 负责处理异步等待

Airflow 界面中的 HITL

使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...