Przepływy z człowiekiem w pętli

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Przypadki użycia HITL

$$

$$

  • Zatwierdzanie odpowiedzi wygenerowanej przez AI
  • Kierowanie zgłoszenia do właściwego zespołu
  • Przeglądanie raportu jakości danych przed publikacją

Schemat human-in-the-loop

Budowanie potoków danych z Airflow

Rodzina operatorów HITL

$$

Operator Zastosowanie
ApprovalOperator Binarne zatwierdzenie/odrzucenie
HITLBranchOperator Wybór kolejnych zadań
HITLOperator Własne opcje + formularz
HITLEntryOperator Wyłącznie dane z formularza

$$

  • Wszystkie operatory HITL są zaimplementowane jako operatory odroczone
Budowanie potoków danych z 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"), }, )
  • Potok wstrzymuje się, dopóki pola nie zostaną wypełnione
Budowanie potoków danych z Airflow

Odczytywanie odpowiedzi

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

 

  • Dane z formularza są w params_input
  • Wybrane opcje (jeśli są) — w chosen_options
Budowanie potoków danych z 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"], )

 

  • Człowiek wybiera opcję mapowaną na ID zadania następczego
  • Nieprzybrane ścieżki są pomijane
Budowanie potoków danych z Airflow

Wymagane akcje w interfejsie

 

  • Zadania HITL tworzą wymaganą akcję
  • Recenzenci odpowiadają w widoku instancji zadania

 

  • Wszystkie operatory HITL są odroczone
  • Zwalniają slot workera podczas oczekiwania
  • Triggerer obsługuje asynchroniczne oczekiwanie

HITL w interfejsie Airflow

Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...