Flux de travail avec humain dans la boucle

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Cas d'usage HITL

$$

$$

  • Approuver une réponse générée par l'IA
  • Acheminer un billet de soutien à la bonne équipe
  • Vérifier un rapport de qualité des données avant publication

Visuel « Human in the loop »

Créer des pipelines de données avec Airflow

La famille d'opérateurs HITL

$$

Opérateur Cas d'usage
ApprovalOperator Décision binaire approuver/refuser
HITLBranchOperator Choisir quelles tâches s'exécutent ensuite
HITLOperator Options personnalisées + formulaire
HITLEntryOperator Saisie de formulaire uniquement, sans options

$$

  • Tous les opérateurs HITL sont des opérateurs ajournables
Créer des pipelines de données avec 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"), }, )
  • Le pipeline se met en pause jusqu'à ce que les champs soient remplis
Créer des pipelines de données avec Airflow

Accéder à la réponse

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

 

  • Les données du formulaire se trouvent dans params_input
  • Les options (s'il y en a) sont dans chosen_options
Créer des pipelines de données avec 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"], )

 

  • Une personne choisit une option, associée à un ID de tâche en aval
  • Les chemins non retenus sont ignorés
Créer des pipelines de données avec Airflow

Actions requises dans l'interface

 

  • Les tâches HITL créent une action requise
  • Les réviseurs répondent dans la vue d'instance de tâche

 

  • Tous les opérateurs HITL sont ajournables
  • Ils libèrent l'emplacement de travailleur pendant l'attente
  • Le Triggerer gère l'attente asynchrone

HITL dans l'interface Airflow

Créer des pipelines de données avec Airflow

Passons à la pratique !

Créer des pipelines de données avec Airflow

Preparing Video For Download...