Human-in-the-loop-Workflows

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

HITL-Anwendungsfälle

$$

$$

  • Eine KI-Antwort freigeben
  • Ein Support-Ticket an das richtige Team leiten
  • Einen Datenqualitätsbericht vor der Veröffentlichung prüfen

Mensch-in-der-Schleife Visualisierung

Data-Pipelines mit Airflow aufbauen

Die HITL-Operatorfamilie

$$

Operator Anwendungsfall
ApprovalOperator Binär freigeben/ablehnen
HITLBranchOperator Welche Tasks als Nächstes laufen
HITLOperator Eigene Optionen + Formular
HITLEntryOperator Reine Formulareingabe, keine Optionen

$$

  • Alle HITL-Operatoren sind als aufschiebbar (deferrable) implementiert
Data-Pipelines mit Airflow aufbauen

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"), }, )
  • Die Pipeline pausiert, bis die Felder ausgefüllt sind
Data-Pipelines mit Airflow aufbauen

Auf die Antwort zugreifen

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

 

  • Formulareingaben stehen in params_input
  • Optionen (falls vorhanden) stehen in chosen_options
Data-Pipelines mit Airflow aufbauen

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"], )

 

  • Mensch wählt eine Option, gemappt auf eine downstream-Task-ID
  • Nicht gewählte Pfade werden übersprungen
Data-Pipelines mit Airflow aufbauen

Required Actions in der UI

 

  • HITL-Tasks erzeugen eine Required Action
  • Reviewer antworten in der Task-Instanzansicht

 

  • Alle HITL-Operatoren sind deferrable
  • Sie geben den Worker-Slot frei, während gewartet wird
  • Der Triggerer übernimmt das asynchrone Warten

HITL in der Airflow-UI

Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...