Fluxuri de lucru human-in-the-loop

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Cazuri de utilizare HITL

$$

$$

  • Aprobarea unui răspuns generat de AI
  • Rutarea unui tichet de suport către echipa potrivită
  • Verificarea unui raport de calitate a datelor înainte de publicare

Vizualizare human-in-the-loop

Construirea pipeline-urilor de date cu Airflow

Familia de operatori HITL

$$

Operator Caz de utilizare
ApprovalOperator Aprobare/respingere binară
HITLBranchOperator Alege ce task-uri urmează
HITLOperator Opțiuni + formular personalizat
HITLEntryOperator Introducere formular, fără opțiuni

$$

  • Toți operatorii HITL sunt implementați ca operatori deferibili
Construirea pipeline-urilor de date cu 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"), }, )
  • Pipeline-ul se oprește până când câmpurile sunt completate
Construirea pipeline-urilor de date cu Airflow

Accesarea răspunsului

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

 

  • Datele din formular se află în params_input
  • Opțiunile (dacă există) se află în chosen_options
Construirea pipeline-urilor de date cu 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"], )

 

  • Utilizatorul selectează o opțiune, mapată la un ID de task din aval
  • Căile neselectate sunt omise
Construirea pipeline-urilor de date cu Airflow

Acțiuni necesare în interfață

 

  • Task-urile HITL creează o acțiune necesară
  • Revizorii răspund în vizualizarea instanței de task

 

  • Toți operatorii HITL sunt deferibili
  • Eliberează slotul de worker în timpul așteptării
  • Triggerer-ul gestionează așteptarea asincronă

HITL în interfața Airflow

Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...