Human-in-the-loop-arbetsflöden

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

HITL-användningsfall

$$

$$

  • Godkänna ett AI-genererat svar
  • Dirigera ett supportärende till rätt team
  • Granska en datakvalitetsrapport innan publicering

Human-in-the-loop-illustration

Bygg datapipelines med Airflow

HITL-operatorfamiljen

$$

Operator Användningsfall
ApprovalOperator Binärt godkänn/avvisa
HITLBranchOperator Välj vilka uppgifter som körs härnäst
HITLOperator Anpassade alternativ + formulär
HITLEntryOperator Enbart formulärinmatning, inga alternativ

$$

  • Alla HITL-operatorer implementeras som deferbara operatorer
Bygg datapipelines med 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"), }, )
  • Pipelinen pausar tills fälten är ifyllda
Bygg datapipelines med Airflow

Läsa svaret

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

 

  • Formulärinmatning finns i params_input
  • Alternativ (om sådana finns) finns i chosen_options
Bygg datapipelines med 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"], )

 

  • En person väljer ett alternativ som mappas till ett nedströms uppgifts-ID
  • Ej valda grenar hoppas över
Bygg datapipelines med Airflow

Required Actions i gränssnittet

 

  • HITL-uppgifter skapar en Required Action
  • Granskare svarar i uppgiftsinstansvyn

 

  • Alla HITL-operatorer är deferbara
  • De frigör worker-platsen under väntetiden
  • Triggerer hanterar det asynkrona väntetillståndet

HITL i Airflow-gränssnittet

Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...