Workflow con human-in-the-loop

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Casi d'uso HITL

$$

$$

  • Approvare una risposta generata dall'AI
  • Inoltrare un ticket di supporto al team giusto
  • Rivedere un report di qualità dei dati prima della pubblicazione

Visual di human-in-the-loop

Creare data pipeline con Airflow

La famiglia di operatori HITL

$$

Operator Use case
ApprovalOperator Binario approva/rifiuta
HITLBranchOperator Scegli quali task eseguire dopo
HITLOperator Opzioni personalizzate + form
HITLEntryOperator Solo input da form, nessuna opzione

$$

  • Tutti gli operatori HITL sono implementati come operatori differibili
Creare data pipeline con 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"), }, )
  • La pipeline si mette in pausa finché i campi non sono compilati
Creare data pipeline con Airflow

Accedere alla risposta

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

 

  • L'input del form è in params_input
  • Le opzioni (se presenti) sono in chosen_options
Creare data pipeline con 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"], )

 

  • L'umano seleziona un'opzione, mappata a un task ID a valle
  • I percorsi non selezionati vengono saltati
Creare data pipeline con Airflow

Required Actions nell'interfaccia

 

  • I task HITL creano una Required Action
  • I revisori rispondono nella vista dell'istanza del task

 

  • Tutti gli operatori HITL sono differibili
  • Liberano lo slot del worker durante l'attesa
  • Il Triggerer gestisce l'attesa asincrona

HITL nell'interfaccia Airflow

Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...