Human-in-the-loop-workflows

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

HITL-use cases

$$

$$

  • Een door AI gegenereerd antwoord goedkeuren
  • Een supportticket naar het juiste team routeren
  • Een datakwaliteitsrapport checken vóór publicatie

Human in the loop visual

Data-pijplijnen bouwen met Airflow

De HITL-operatorfamilie

$$

Operator Use case
ApprovalOperator Binaire approve/reject
HITLBranchOperator Kies welke taken daarna draaien
HITLOperator Aangepaste opties + formulier
HITLEntryOperator Alleen formulierinvoer, geen opties

$$

  • Alle HITL-operators zijn geïmplementeerd als deferrable operators
Data-pijplijnen bouwen met 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"), }, )
  • De pipeline pauzeert tot de velden zijn ingevuld
Data-pijplijnen bouwen met Airflow

De response benaderen

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

 

  • Formulierinvoer staat in params_input
  • Opties (indien aanwezig) staan in chosen_options
Data-pijplijnen bouwen met 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"], )

 

  • Mens selecteert een optie, gemapt naar een downstream task ID
  • Niet-gekozen paden worden overgeslagen
Data-pijplijnen bouwen met Airflow

Required Actions in de UI

 

  • HITL-taken maken een Required Action
  • Reviewers reageren in de task instance view

 

  • Alle HITL-operators zijn deferrable
  • Ze maken de workerslot vrij tijdens het wachten
  • De Triggerer handelt het asynchrone wachten af

HITL in the Airflow UI

Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...