Human-in-the-loop workflows

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Případy použití HITL

$$

$$

  • Schválení odpovědi vygenerované AI
  • Přesměrování tiketu podpory správnému týmu
  • Kontrola reportu kvality dat před publikací

Vizualizace human-in-the-loop

Tvorba datových pipeline s Airflow

Rodina HITL operátorů

$$

Operátor Případ použití
ApprovalOperator Binární schválení/zamítnutí
HITLBranchOperator Výběr dalších úloh
HITLOperator Vlastní možnosti + formulář
HITLEntryOperator Čistý vstup formuláře, bez možností

$$

  • Všechny HITL operátory jsou implementovány jako deferrable operators
Tvorba datových pipeline s 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 se pozastaví, dokud nejsou pole vyplněna
Tvorba datových pipeline s Airflow

Přístup k odpovědi

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

 

  • Vstup formuláře je v params_input
  • Zvolené možnosti (pokud existují) jsou v chosen_options
Tvorba datových pipeline s 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"], )

 

  • Člověk vybere možnost, která se mapuje na ID navazující úlohy
  • Nevybrané větve jsou přeskočeny
Tvorba datových pipeline s Airflow

Required Actions v UI

 

  • HITL úlohy vytvoří Required Action
  • Kontroloři reagují v zobrazení instance úlohy

 

  • Všechny HITL operátory jsou deferrable
  • Během čekání uvolňují worker slot
  • Asynchronní čekání zajišťuje Triggerer

HITL v Airflow UI

Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...