Flujos con personas en el bucle

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Casos de uso de HITL

$$

$$

  • Aprobar una respuesta generada por IA
  • Derivar un ticket de soporte al equipo adecuado
  • Revisar un informe de calidad de datos antes de publicar

Persona en el bucle (visual)

Creación de canalizaciones de datos con Airflow

La familia de operadores HITL

$$

Operador Caso de uso
ApprovalOperator Aprobar/rechazar binario
HITLBranchOperator Elegir qué tareas se ejecutan después
HITLOperator Opciones personalizadas + formulario
HITLEntryOperator Solo entrada de formulario, sin opciones

$$

  • Todos los operadores HITL se implementan como operadores aplazables
Creación de canalizaciones de datos con 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"), }, )
  • La canalización se pausa hasta que se completen los campos
Creación de canalizaciones de datos con Airflow

Acceder a la respuesta

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

 

  • La entrada del formulario está en params_input
  • Las opciones (si las hay) están en chosen_options
Creación de canalizaciones de datos con 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"], )

 

  • Una persona selecciona una opción, mapeada a un ID de tarea siguiente
  • Las rutas no seleccionadas se omiten
Creación de canalizaciones de datos con Airflow

Acciones requeridas en la interfaz

 

  • Las tareas HITL crean una Acción requerida
  • Quien revisa responde en la vista de instancia de tarea

 

  • Todos los operadores HITL son aplazables
  • Liberan el slot del worker mientras esperan
  • El Triggerer gestiona la espera asíncrona

HITL en la interfaz de Airflow

Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...