Organizar Dags complejos con Task Groups

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

El problema de la complejidad

Dag complejo

  • Los Dags grandes son difíciles de navegar en la vista Graph
  • Los nombres de tareas se confunden
  • A quien se incorpora le cuesta ver qué le corresponde
Creación de canalizaciones de datos con Airflow

@task_group

from airflow.sdk import dag, task, task_group

@task_group(
    group_id="ingest_orders",
    default_args={"retries": 3},
)

def process_orders(): @task def extract_orders(): return [{"id": 1, "amount": 99.99}] @task def transform_orders(orders): return [{"id": o["id"], "total": o["amount"] * 1.08} for o in orders] return transform_orders(extract_orders())
  • group_id establece un identificador personalizado; por defecto usa el nombre de la función
  • default_args se aplica a todas las tareas del grupo para evitar repetir configuración
  • Los task groups aparecen como bloques desplegables en la UI
Creación de canalizaciones de datos con Airflow

Cómo se ve en la UI de Airflow

Sin task groups

Gráfico simple con tres tareas en línea: extract_orders, transform_orders, load, todas al mismo nivel

  • Todas las tareas al mismo nivel

Con task groups

Gráfico con un bloque contraído llamado process_orders que contiene extract y transform, seguido de una tarea load fuera del grupo

  • Bloques colapsables en la UI
Creación de canalizaciones de datos con Airflow

Anidación y nombres visibles personalizados

@task_group(group_display_name="Process All")
def process_all():

    @task_group(group_display_name="Ingest Orders")
    def orders():
        return transform(extract())

    @task_group(group_display_name="Process Returns")
    def returns():
        return transform(extract())

    return {
        "orders": orders(),
        "returns": returns(),
    }
  • Puedes anidar task groups para canalizaciones más complejas
  • group_display_name define una etiqueta legible en la UI (incluso admite emojis)
Creación de canalizaciones de datos con Airflow

El patrón de factoría

@task_group
def process_source(source_name, source_path):

    @task
    def extract():
        return read_data(source_path)

    @task
    def transform(data):
        return clean_data(data)

    return transform(extract())

# Reuse the same pattern for different sources orders = process_source("orders", "/data/orders.csv") returns = process_source("returns", "/data/returns.csv") events = process_source("events", "/data/events.csv")
  • Los task groups son funciones de Python decoradas
  • Puedes llamarlos varias veces con distintos parámetros
  • Este patrón de factoría permite reutilizar lógica de canalización
Creación de canalizaciones de datos con Airflow

Guías para agrupar

Agrupar por dominio o preocupación

  • Agrupa por dominio o preocupación, no por tipo de operador
  • Usa group_id para nombres programáticos claros y group_display_name para etiquetas de UI
  • Usa default_args para compartir configuración como reintentos en todas las tareas del grupo
  • Aplica la ley de Miller: más de 7 elementos de nivel superior puede indicar que necesitas un task group
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...