AlpinaShop tiene ya una plataforma de datos completa. Y tiene, exactamente por eso, un problema nuevo.

Cada noche hay que hacer esto: exportar los pedidos del día desde alpinashop-pedidos a Cloud Storage, cargarlos en BigQuery, lanzar el pipeline de Dataflow que limpia y agrega, ejecutar las consultas de agregación, refrescar la vista de negocio que consume el cuadro de mando, y —el primer día de cada mes— lanzar además el trabajo de Spark que recalcula la matriz de productos comprados juntos y el pipeline de Data Fusion con el fichero del transportista.

Son ocho procesos con dependencias reales entre ellos. No tiene sentido cargar en BigQuery un fichero que aún se está escribiendo, ni agregar sobre datos a medias, ni refrescar la vista con la mitad de los pedidos del día.

Ahora mismo eso lo resuelve un cron en una máquina virtual:

# /etc/crontab de la VM vm-procesos-nocturnos
0  2 * * * /opt/scripts/exportar_cloudsql.sh
30 2 * * * /opt/scripts/cargar_bigquery.sh
0  3 * * * /opt/scripts/lanzar_dataflow.sh
0  4 * * * /opt/scripts/agregaciones.sh
30 4 * * * /opt/scripts/refrescar_vista.sh

Funciona. Funciona todas las noches, hasta la primera en que la exportación tarda cuarenta minutos en vez de veinte porque hubo campaña. Entonces, a las 02:30, cargar_bigquery.sh lee un fichero incompleto. A las 03:00, Dataflow procesa datos parciales. A las 04:30, la vista de negocio se refresca con la mitad de los pedidos. A las 09:00, dirección abre el cuadro de mando y ve que las ventas de ayer cayeron un 45 %. Se convoca una reunión de urgencia. A las 11:30, Marta descubre que el dato es falso.

Nadie se enteró de nada porque cron no sabe si algo funcionó. Solo sabe qué hora es.

En esta lección verás qué hace un orquestador de verdad, construirás el DAG completo del proceso nocturno de AlpinaShop en Cloud Composer, conocerás su coste real —que es el argumento decisivo para una pyme—, y montarás la alternativa ligera con Workflows y Cloud Scheduler, que es por donde AlpinaShop va a empezar.

Contenido

  1. Por qué cron no es un orquestador
  2. Qué hace un orquestador de verdad
  3. Cloud Composer: Apache Airflow gestionado
  4. Conceptos de Airflow: DAG, tarea, operador, sensor
  5. El entorno alpinashop-composer y su coste
  6. El DAG del proceso nocturno, línea a línea
  7. Operadores de Google Cloud
  8. XComs, variables y conexiones
  9. TaskGroup, dependencias y patrones de flujo
  10. Idempotencia, catchup y reprocesos
  11. Workflows: orquestación serverless en YAML
  12. Cloud Scheduler: el disparador por cron
  13. La tabla de decisión y la elección de AlpinaShop

  1. Por qué cron no es un orquestador

cron responde a una sola pregunta: ¿qué hora es?. Un orquestador responde a otra muy distinta: ¿qué se puede ejecutar ahora, dado lo que ha pasado?.

Situación Con cron Con un orquestador
La tarea A tarda más de lo previsto B arranca igual, con datos incompletos B espera a que A termine bien
La tarea A falla B arranca igual, sobre nada B no arranca; se avisa
Fallo transitorio de red El proceso muere Reintento automático con espera
¿Se ejecutó anoche? Mirar logs por SSH Panel con el historial completo
Reprocesar el 12 de marzo Script manual con parámetros a mano backfill de esa fecha
Tres tareas independientes En serie, sumando tiempos En paralelo
Nadie se entera de un fallo Correcto: nadie se entera Alerta configurada
La VM del cron se cae Nada se ejecuta y nadie lo sabe Servicio gestionado con reintentos
¿Cuánto tarda cada paso? No se sabe Métricas por tarea e histórico

Hay una fila especialmente insidiosa: la VM del cron. Es una máquina que alguien creó hace tres años, que nadie parchea, que tiene los scripts en /opt sin control de versiones, con credenciales en ficheros .env, y de cuya existencia solo se acuerdan dos personas. Cuando esa VM muere, muere en silencio y la plataforma de datos deja de actualizarse durante días.

La segunda fila insidiosa es la de la espera a ojo. El cron de arriba asume que la exportación tarda menos de 30 minutos. Ese margen es una apuesta, y las apuestas se pierden justo el día de más volumen, que es el día en que los datos más importan.

  1. Qué hace un orquestador de verdad

Un orquestador de flujos de trabajo aporta cinco cosas:

Dependencias explícitas. Se declara que B depende de A, y el sistema garantiza el orden. No hay horas calculadas a ojo: si A tarda diez minutos o dos horas, B arranca cuando A termina.

Gestión de fallos. Reintentos con espera creciente, número máximo de intentos, y qué hacer si aun así falla: parar, seguir con lo que no dependa, o ejecutar una tarea de limpieza.

Observabilidad. Un panel donde se ve cada ejecución, cada tarea, su duración, sus logs y su historial. Responder a "¿funcionó anoche?" cuesta un vistazo.

Reproducibilidad. Poder reejecutar el flujo de una fecha pasada con los parámetros de esa fecha, sin editar nada.

Notificación. Alguien se entera cuando algo falla, y se entera a tiempo.

Google Cloud ofrece tres herramientas para esto, y son complementarias:

Herramienta Qué es Modelo
Cloud Scheduler Un cron gestionado Dispara una acción a una hora
Workflows Orquestación serverless declarativa Encadena pasos en YAML, sin servidor
Cloud Composer Apache Airflow gestionado Orquestación completa con Python

  1. Cloud Composer: Apache Airflow gestionado

Apache Airflow es el estándar de facto de la orquestación de datos. Nació en Airbnb en 2014, es open source, y su idea central es que los flujos de trabajo se definen como código Python.

Eso último es su gran virtud. Un flujo definido en Python se versiona en Git, se revisa en un pull request, se prueba, se genera dinámicamente con bucles, y admite toda la lógica del lenguaje. Frente a una interfaz visual de arrastrar cajas, un DAG en Python es infinitamente más mantenible cuando el equipo sabe programar.

Cloud Composer es Airflow gestionado por Google: el planificador, la base de datos de metadatos, los trabajadores y la interfaz web funcionando sobre GKE, con integración de IAM, Cloud Logging y Cloud Monitoring, y con los operadores de Google Cloud ya instalados.

Versiones vigentes en 2026: Composer 3, con Airflow 2.x y 3.x. Composer 3 simplificó bastante la arquitectura respecto a Composer 2 —menos infraestructura visible, escalado más fino— pero el modelo de coste sigue siendo el mismo en lo esencial, y ese es el punto que hay que mirar antes que ninguno.

  1. Conceptos de Airflow: DAG, tarea, operador, sensor

DAG (grafo acíclico dirigido) es el flujo de trabajo completo. Dirigido porque las dependencias tienen sentido; acíclico porque no puede haber bucles: si A depende de B y B de A, nada podría empezar nunca.

Tarea (task) es un nodo del DAG: una unidad de trabajo.

Operador (operator) es la plantilla que define qué hace una tarea. Airflow trae cientos: ejecutar Bash, llamar a una función Python, lanzar una consulta de BigQuery, crear un clúster de Dataproc, enviar un correo.

Sensor es un operador especial que espera a que ocurra algo: que aparezca un fichero en un bucket, que una tabla tenga datos, que una API responda. Es la pieza que resuelve el problema del cron.

Ejecución (DAG run) es una instancia concreta del DAG para una fecha lógica determinada.

flowchart LR
    A["exportar_cloudsql<br/>BashOperator"]
    B["esperar_fichero<br/>GCSObjectExistenceSensor"]
    C["cargar_bigquery<br/>GCSToBigQueryOperator"]
    D["lanzar_dataflow<br/>DataflowFlexTemplateOperator"]
    E1["agregar_ventas<br/>BigQueryInsertJobOperator"]
    E2["agregar_visitas<br/>BigQueryInsertJobOperator"]
    F["refrescar_vista<br/>BigQueryInsertJobOperator"]
    G["comprobar_calidad<br/>BigQueryCheckOperator"]

    A --> B --> C --> D
    D --> E1 --> F
    D --> E2 --> F
    F --> G

Fíjate en que agregar_ventas y agregar_visitas se ejecutan en paralelo: ambas dependen de lanzar_dataflow y ninguna depende de la otra. Airflow lo deduce del grafo sin que haya que decirlo. Con cron habría que decidir un orden y sumar los tiempos.

  1. El entorno alpinashop-composer y su coste

gcloud config set project alpinashop-datos
gcloud services enable composer.googleapis.com

gcloud iam service-accounts create sa-composer \
  --display-name="Cloud Composer de AlpinaShop"

SA_COMP="[email protected]"

for ROL in roles/composer.worker roles/bigquery.dataEditor roles/bigquery.jobUser \
           roles/dataflow.developer roles/storage.objectAdmin \
           roles/cloudsql.viewer roles/dataproc.editor; do
  gcloud projects add-iam-policy-binding alpinashop-datos \
    --member="serviceAccount:${SA_COMP}" --role="$ROL"
done

gcloud composer environments create alpinashop-composer \
  --location=europe-west1 \
  --image-version=composer-3-airflow-2.10.5 \
  --service-account="$SA_COMP" \
  --network=alpinashop-vpc \
  --subnetwork=sn-datos-euw1 \
  --enable-private-environment \
  --environment-size=small \
  --labels=entorno=produccion,equipo=datos,centro-coste=analitica

La creación tarda de 20 a 30 minutos.

Y ahora la conversación incómoda, que hay que tener antes de escribir una línea de DAG.

Componente Coste aproximado mensual (verificar en la documentación oficial)
Entorno small (planificador, servidor web, base de datos) ~250-350 €
Trabajadores adicionales bajo carga Variable
Almacenamiento del bucket de DAG y logs Céntimos
Total realista de un entorno pequeño ~300-400 €/mes

Trescientos euros al mes, se ejecute un DAG o cien. Es un servicio permanentemente encendido: el planificador tiene que estar vivo para saber qué hora es.

Pongamos eso en contexto con el resto de la plataforma de AlpinaShop:

Servicio Coste mensual estimado
BigQuery (almacenamiento + consultas) ~10 €
Cloud Storage (60 GB de catálogo) ~2 €
Pub/Sub Céntimos
Dataflow por lotes (nocturno) ~2 €
Dataproc Serverless (mensual) ~1 €
Cloud Composer ~350 €

El orquestador costaría veinte veces más que todo lo que orquesta. Eso no es un argumento contra Composer: es un argumento contra usarlo en el momento equivocado. Para una empresa con 200 DAG, 50 ingenieros y dependencias entre equipos, 350 € es irrisorio frente al valor. Para AlpinaShop, con ocho tareas nocturnas, es desproporcionado.

Aun así vamos a construir el DAG completo, por tres razones: porque Airflow es el estándar del sector y hay que saberlo; porque el ejercicio de modelar las dependencias es válido para cualquier orquestador; y porque el día en que AlpinaShop crezca, esta será la herramienta. Al final de la lección volveremos a la decisión.

  1. El DAG del proceso nocturno, línea a línea

Los DAG se despliegan copiándolos al bucket que Composer crea:

BUCKET_DAGS=$(gcloud composer environments describe alpinashop-composer \
  --location=europe-west1 --format="value(config.dagGcsPrefix)")

gcloud storage cp dags/proceso_nocturno.py "${BUCKET_DAGS}/"

Y el DAG:

"""
proceso_nocturno.py -- Proceso nocturno de datos de AlpinaShop.

Flujo:
  exportar Cloud SQL -> esperar fichero -> cargar BigQuery -> Dataflow
  -> agregaciones (paralelas) -> refrescar vista -> comprobar calidad

Despliegue: copiar a gs://<bucket-composer>/dags/
"""
from datetime import datetime, timedelta

import pendulum
from airflow import DAG
from airflow.operators.empty import EmptyOperator
from airflow.operators.python import PythonOperator
from airflow.providers.google.cloud.operators.bigquery import (
    BigQueryCheckOperator,
    BigQueryInsertJobOperator,
)
from airflow.providers.google.cloud.operators.cloud_sql import (
    CloudSQLExportInstanceOperator,
)
from airflow.providers.google.cloud.operators.dataflow import (
    DataflowStartFlexTemplateOperator,
)
from airflow.providers.google.cloud.sensors.gcs import GCSObjectExistenceSensor
from airflow.providers.google.cloud.transfers.gcs_to_bigquery import (
    GCSToBigQueryOperator,
)
from airflow.utils.task_group import TaskGroup

PROYECTO_DATOS = "alpinashop-datos"
PROYECTO_PROD = "alpinashop-prod"
DATASET = "alpinashop_analitica"
BUCKET = "alpinashop-datalake"
REGION = "europe-west1"
ZONA_HORARIA = pendulum.timezone("Europe/Madrid")
# ---------------------------------------------------------------------------
# ARGUMENTOS POR DEFECTO: se aplican a TODAS las tareas del DAG.
# Definirlos aqui evita repetirlos en cada operador.
# ---------------------------------------------------------------------------
argumentos_por_defecto = {
    "owner": "equipo-datos",
    "depends_on_past": False,      # una ejecucion no espera a la anterior
    "email": ["[email protected]"],
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 3,                  # 3 reintentos ante fallo
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,   # 5, 10, 20 min: no machacar el origen
    "max_retry_delay": timedelta(minutes=30),
    "execution_timeout": timedelta(hours=2),   # ninguna tarea eterna
    "sla": timedelta(hours=3),     # aviso si el DAG no acaba en 3 h
}

Cada uno de estos parámetros evita un incidente concreto:

  • retries con retroceso exponencial: un fallo transitorio de red no rompe la noche, y los reintentos no hunden un sistema ya saturado.
  • execution_timeout: una tarea colgada no bloquea el DAG indefinidamente ni consume un trabajador para siempre.
  • sla: si a las 05:00 el proceso no ha terminado, alguien se entera antes de que dirección abra el informe.
  • depends_on_past: False: cada noche es independiente. Si se pone en True, un fallo del lunes bloquearía el martes, el miércoles y todo lo demás, lo que casi nunca es lo que se quiere.
with DAG(
    dag_id="alpinashop_proceso_nocturno",
    description="Proceso nocturno: Cloud SQL -> BigQuery -> Dataflow -> agregados",
    default_args=argumentos_por_defecto,
    # Todos los dias a las 02:00 hora de Madrid.
    # Airflow trabaja internamente en UTC; la zona horaria evita el desfase
    # de una hora en los cambios de horario de verano.
    schedule="0 2 * * *",
    start_date=datetime(2026, 3, 1, tzinfo=ZONA_HORARIA),
    catchup=False,             # ver apartado 10
    max_active_runs=1,         # nunca dos noches solapadas
    tags=["alpinashop", "produccion", "datos"],
    doc_md=__doc__,            # la docstring se ve en la interfaz
) as dag:

    inicio = EmptyOperator(task_id="inicio")
    # -----------------------------------------------------------------------
    # 1) EXPORTAR Cloud SQL a Cloud Storage
    #
    # {{ ds }} es una plantilla Jinja: Airflow la sustituye por la fecha
    # logica de la ejecucion (yyyy-MM-dd). Es LA pieza que hace el DAG
    # reproducible: al reprocesar el 12 de marzo, {{ ds }} vale 2026-03-12
    # y el fichero de salida y la consulta apuntan a ese dia, no a hoy.
    # -----------------------------------------------------------------------
    exportar_pedidos = CloudSQLExportInstanceOperator(
        task_id="exportar_pedidos_cloudsql",
        project_id=PROYECTO_PROD,
        instance="alpinashop-pedidos-replica-informes",
        body={
            "exportContext": {
                "fileType": "CSV",
                "uri": f"gs://{BUCKET}/exportaciones/{{{{ ds_nodash }}}}/pedidos.csv",
                "databases": ["tienda"],
                "csvExportOptions": {
                    "selectQuery": (
                        "SELECT pedido_id, creado_en, cliente_id, canal, estado, "
                        "pais, ciudad, codigo_postal, metodo_pago, "
                        "subtotal, descuento, iva, total "
                        "FROM pedidos "
                        "WHERE DATE(creado_en) = '{{ ds }}'"
                    )
                },
            }
        },
    )

Se exporta desde la réplica de lectura, no desde la instancia principal. Es la misma disciplina de 02-03 y 04-01: los procesos analíticos no tocan la base de datos que atiende las compras.

    # -----------------------------------------------------------------------
    # 2) SENSOR: esperar a que el fichero exista de verdad.
    #
    # Esta tarea es la que resuelve el problema del cron. La exportacion
    # es asincrona: el operador anterior la lanza, pero el fichero puede
    # tardar. En vez de "esperamos 30 minutos y cruzamos los dedos",
    # se COMPRUEBA.
    # -----------------------------------------------------------------------
    esperar_fichero = GCSObjectExistenceSensor(
        task_id="esperar_fichero_exportado",
        bucket=BUCKET,
        object="exportaciones/{{ ds_nodash }}/pedidos.csv",
        # 'reschedule': libera el trabajador entre comprobaciones en vez de
        # ocuparlo esperando. Imprescindible en esperas largas.
        mode="reschedule",
        poke_interval=60,          # comprobar cada minuto
        timeout=60 * 60,           # rendirse a la hora
    )

El modo reschedule frente a poke es un detalle con consecuencias reales: en modo poke, el sensor ocupa un trabajador durante toda la espera. Con cuatro sensores esperando una hora en un entorno pequeño, no queda ningún trabajador libre y el DAG se bloquea a sí mismo. En modo reschedule, el sensor se duerme y libera el hueco.

    # -----------------------------------------------------------------------
    # 3) CARGAR en BigQuery
    # -----------------------------------------------------------------------
    cargar_pedidos = GCSToBigQueryOperator(
        task_id="cargar_pedidos_bigquery",
        bucket=BUCKET,
        source_objects=["exportaciones/{{ ds_nodash }}/pedidos.csv"],
        destination_project_dataset_table=f"{PROYECTO_DATOS}.{DATASET}.pedidos_staging",
        source_format="CSV",
        skip_leading_rows=0,
        field_delimiter=",",
        null_marker="\\N",
        # WRITE_TRUNCATE en la tabla de staging: cada noche se reemplaza.
        # Esto hace la tarea IDEMPOTENTE: reejecutarla no duplica nada.
        write_disposition="WRITE_TRUNCATE",
        create_disposition="CREATE_IF_NEEDED",
        autodetect=False,
        schema_fields=[
            {"name": "pedido_id", "type": "STRING", "mode": "REQUIRED"},
            {"name": "creado_en", "type": "TIMESTAMP", "mode": "REQUIRED"},
            {"name": "cliente_id", "type": "STRING"},
            {"name": "canal", "type": "STRING"},
            {"name": "estado", "type": "STRING"},
            {"name": "pais", "type": "STRING"},
            {"name": "ciudad", "type": "STRING"},
            {"name": "codigo_postal", "type": "STRING"},
            {"name": "metodo_pago", "type": "STRING"},
            {"name": "subtotal", "type": "NUMERIC"},
            {"name": "descuento", "type": "NUMERIC"},
            {"name": "iva", "type": "NUMERIC"},
            {"name": "total", "type": "NUMERIC"},
        ],
        location=REGION,
    )

    # -----------------------------------------------------------------------
    # 4) MERGE del staging a la tabla final: idempotente por diseno
    # -----------------------------------------------------------------------
    consolidar_pedidos = BigQueryInsertJobOperator(
        task_id="consolidar_pedidos",
        location=REGION,
        configuration={
            "query": {
                "query": f"""
                    MERGE `{PROYECTO_DATOS}.{DATASET}.pedidos` AS destino
                    USING (
                      SELECT
                        pedido_id,
                        DATE(creado_en) AS fecha_pedido,
                        creado_en       AS momento_pedido,
                        cliente_id, canal, estado,
                        STRUCT(pais, NULL AS provincia, ciudad, codigo_postal,
                               NULL AS metodo, CAST(NULL AS NUMERIC) AS coste) AS envio,
                        metodo_pago, NULL AS cupon,
                        subtotal, descuento, iva, total AS total_pedido
                      FROM `{PROYECTO_DATOS}.{DATASET}.pedidos_staging`
                    ) AS origen
                    ON destino.pedido_id = origen.pedido_id
                       AND destino.fecha_pedido = origen.fecha_pedido
                    WHEN MATCHED THEN UPDATE SET
                      estado = origen.estado,
                      total_pedido = origen.total_pedido
                    WHEN NOT MATCHED THEN INSERT ROW
                """,
                "useLegacySql": False,
            }
        },
    )

El MERGE es la clave de la idempotencia: si la tarea se reejecuta, los pedidos existentes se actualizan y los nuevos se insertan. Nunca se duplica nada. Es exactamente el mismo principio que el Upsert de Data Fusion en 04-05 y la idempotencia de los consumidores de Pub/Sub en 04-04. El mismo concepto aparece en las cuatro lecciones porque es la propiedad que hace que un sistema de datos se pueda operar sin miedo.

    # -----------------------------------------------------------------------
    # 5) DATAFLOW: plantilla flexible creada en 04-02
    # -----------------------------------------------------------------------
    lanzar_dataflow = DataflowStartFlexTemplateOperator(
        task_id="lanzar_pipeline_dataflow",
        location=REGION,
        project_id=PROYECTO_DATOS,
        body={
            "launchParameter": {
                "jobName": "pedidos-lote-{{ ds_nodash }}",
                "containerSpecGcsPath":
                    f"gs://alpinashop-dataflow/plantillas/pedidos-lote.json",
                "parameters": {"fecha": "{{ ds }}"},
                "environment": {
                    "serviceAccountEmail":
                        "[email protected]",
                    "subnetwork":
                        f"regions/{REGION}/subnetworks/sn-datos-euw1",
                    "ipConfiguration": "WORKER_IP_PRIVATE",
                    "maxWorkers": 10,
                    "additionalUserLabels": {
                        "entorno": "produccion", "equipo": "datos",
                    },
                },
            }
        },
        # El operador ESPERA a que el job termine antes de dar la tarea
        # por completada. Sin esto, el DAG seguiria con datos a medias.
        wait_until_finished=True,
    )

    # -----------------------------------------------------------------------
    # 6) AGREGACIONES EN PARALELO, agrupadas para que el grafo sea legible
    # -----------------------------------------------------------------------
    with TaskGroup(group_id="agregaciones") as grupo_agregaciones:

        agregar_ventas = BigQueryInsertJobOperator(
            task_id="ventas_por_categoria",
            location=REGION,
            configuration={
                "query": {
                    "query": f"""
                        CREATE OR REPLACE TABLE
                          `{PROYECTO_DATOS}.{DATASET}.agg_ventas_categoria_dia`
                        PARTITION BY dia AS
                        SELECT
                          l.fecha_pedido                 AS dia,
                          pr.categoria,
                          COUNT(DISTINCT l.pedido_id)    AS pedidos,
                          SUM(l.cantidad)                AS unidades,
                          ROUND(SUM(l.importe_linea), 2) AS ventas_eur
                        FROM `{PROYECTO_DATOS}.{DATASET}.lineas_pedido` AS l
                        JOIN `{PROYECTO_DATOS}.{DATASET}.productos`     AS pr
                          USING (sku)
                        WHERE l.fecha_pedido >= DATE_SUB(DATE '{{{{ ds }}}}',
                                                         INTERVAL 400 DAY)
                        GROUP BY dia, pr.categoria
                    """,
                    "useLegacySql": False,
                }
            },
        )

        agregar_visitas = BigQueryInsertJobOperator(
            task_id="embudo_conversion",
            location=REGION,
            configuration={
                "query": {
                    "query": f"""
                        CREATE OR REPLACE TABLE
                          `{PROYECTO_DATOS}.{DATASET}.agg_embudo_dia`
                        PARTITION BY dia AS
                        WITH marcas AS (
                          SELECT
                            fecha AS dia, sesion_id,
                            LOGICAL_OR(e.tipo = 'ver_producto')   AS vio,
                            LOGICAL_OR(e.tipo = 'anadir_carrito') AS anadio,
                            LOGICAL_OR(e.tipo = 'iniciar_pago')   AS pago,
                            LOGICAL_OR(e.tipo = 'compra')         AS compro
                          FROM `{PROYECTO_DATOS}.{DATASET}.visitas`,
                               UNNEST(eventos) AS e
                          WHERE fecha BETWEEN DATE_SUB(DATE '{{{{ ds }}}}',
                                                       INTERVAL 90 DAY)
                                          AND DATE '{{{{ ds }}}}'
                          GROUP BY dia, sesion_id
                        )
                        SELECT
                          dia,
                          COUNT(*)             AS sesiones,
                          COUNTIF(vio)         AS vieron_producto,
                          COUNTIF(anadio)      AS anadieron_carrito,
                          COUNTIF(pago)        AS iniciaron_pago,
                          COUNTIF(compro)      AS compraron
                        FROM marcas
                        GROUP BY dia
                    """,
                    "useLegacySql": False,
                }
            },
        )

    # -----------------------------------------------------------------------
    # 7) REFRESCAR la vista materializada que consume el cuadro de mando
    # -----------------------------------------------------------------------
    refrescar_vista = BigQueryInsertJobOperator(
        task_id="refrescar_vista_negocio",
        location=REGION,
        configuration={
            "query": {
                "query": f"""
                    CALL BQ.REFRESH_MATERIALIZED_VIEW(
                      '{PROYECTO_DATOS}.{DATASET}.mv_ventas_diarias_sku')
                """,
                "useLegacySql": False,
            }
        },
    )

    # -----------------------------------------------------------------------
    # 8) COMPROBACION DE CALIDAD: la ultima defensa antes del informe.
    #
    # Si esta tarea falla, el DAG queda marcado como fallido y salta la
    # alerta. Es preferible una alerta a las 05:00 que un comite de
    # direccion mirando cifras falsas a las 09:00.
    # -----------------------------------------------------------------------
    comprobar_calidad = BigQueryCheckOperator(
        task_id="comprobar_calidad_datos",
        location=REGION,
        use_legacy_sql=False,
        sql=f"""
            SELECT
              COUNTIF(pedidos_dia = 0)  = 0 AND
              COUNTIF(ventas_dia < 0)   = 0 AND
              COUNTIF(pedidos_dia > 500) = 0
            FROM (
              SELECT
                dia,
                SUM(pedidos)    AS pedidos_dia,
                SUM(ventas_eur) AS ventas_dia
              FROM `{PROYECTO_DATOS}.{DATASET}.agg_ventas_categoria_dia`
              WHERE dia = DATE '{{{{ ds }}}}'
              GROUP BY dia
            )
        """,
    )

    def _avisar_exito(**contexto):
        """Callback informativo: en produccion enviaria a Slack o Chat."""
        print(f"Proceso nocturno completado para {contexto['ds']}")

    fin = PythonOperator(
        task_id="fin",
        python_callable=_avisar_exito,
        trigger_rule="all_success",
    )

    # -----------------------------------------------------------------------
    # DEPENDENCIAS: el operador >> significa "y despues".
    # Esta es la declaracion completa del grafo, en cuatro lineas.
    # -----------------------------------------------------------------------
    (
        inicio
        >> exportar_pedidos
        >> esperar_fichero
        >> cargar_pedidos
        >> consolidar_pedidos
        >> lanzar_dataflow
        >> grupo_agregaciones
        >> refrescar_vista
        >> comprobar_calidad
        >> fin
    )

La tarea 8, la comprobación de calidad, es la que convierte este DAG en algo serio. Verifica tres reglas de negocio: que hubo pedidos, que ninguna venta es negativa, y que ninguna categoría supera 500 pedidos en un día (imposible con el volumen de AlpinaShop, luego indicaría duplicación). Si algo no cuadra, el DAG falla ruidosamente antes de que nadie mire el informe. Es la diferencia entre detectar un error a las 05:00 y descubrirlo en una reunión.

  1. Operadores de Google Cloud

Airflow trae un catálogo enorme de operadores para Google Cloud. Los más útiles para AlpinaShop:

Operador Qué hace
BigQueryInsertJobOperator Ejecuta cualquier consulta o trabajo de BigQuery. El más versátil
BigQueryCheckOperator Ejecuta un SQL que debe devolver verdadero; si no, falla
BigQueryValueCheckOperator Compara un resultado con un valor esperado y una tolerancia
BigQueryTableExistenceSensor Espera a que exista una tabla
GCSToBigQueryOperator Carga ficheros del bucket a una tabla
BigQueryToGCSOperator Exporta una tabla a ficheros
GCSObjectExistenceSensor Espera a que aparezca un objeto
GCSToGCSOperator Copia o mueve entre buckets
DataflowStartFlexTemplateOperator Lanza una plantilla flexible de Dataflow
DataprocCreateBatchOperator Lanza un trabajo en Dataproc Serverless
DataprocCreateClusterOperator / DeleteCluster Clúster efímero desde el DAG
CloudSQLExportInstanceOperator Exporta una instancia de Cloud SQL
PubSubPublishMessageOperator Publica en un topic
CloudRunExecuteJobOperator Ejecuta un job de Cloud Run

Para el proceso mensual de AlpinaShop, el análisis de cesta de 04-03 se lanza así:

from airflow.providers.google.cloud.operators.dataproc import (
    DataprocCreateBatchOperator,
)

analisis_cesta = DataprocCreateBatchOperator(
    task_id="analisis_cesta_mensual",
    region=REGION,
    project_id=PROYECTO_DATOS,
    batch_id="cesta-{{ ds_nodash }}",
    batch={
        "pyspark_batch": {
            "main_python_file_uri": f"gs://{BUCKET}/jobs/cesta_media.py",
            "args": [
                "--fecha-desde={{ macros.ds_add(ds, -30) }}",
                "--fecha-hasta={{ ds }}",
            ],
        },
        "runtime_config": {"version": "2.2"},
        "environment_config": {
            "execution_config": {
                "service_account":
                    "[email protected]",
                "subnetwork_uri": "sn-datos-euw1",
            }
        },
    },
)

{{ macros.ds_add(ds, -30) }} calcula una fecha relativa a la de ejecución. Es lo que mantiene el DAG reproducible: al reprocesar enero, el rango será el de enero, no el de hoy.

  1. XComs, variables y conexiones

XCom (cross-communication) permite que una tarea pase un valor pequeño a otra:

def _contar_filas(**contexto):
    from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook
    hook = BigQueryHook(location=REGION, use_legacy_sql=False)
    filas = hook.get_first(
        f"""SELECT COUNT(*) FROM `{PROYECTO_DATOS}.{DATASET}.pedidos_staging`"""
    )[0]
    # Lo devuelto por la funcion se guarda automaticamente como XCom
    return int(filas)


def _validar_volumen(**contexto):
    filas = contexto["ti"].xcom_pull(task_ids="contar_filas_cargadas")
    if filas == 0:
        raise ValueError(f"Ningun pedido cargado el {contexto['ds']}")
    if filas > 5000:
        raise ValueError(f"{filas} pedidos: volumen anomalo, revisar duplicados")
    print(f"Volumen correcto: {filas} pedidos")


contar_filas = PythonOperator(task_id="contar_filas_cargadas",
                              python_callable=_contar_filas)
validar_volumen = PythonOperator(task_id="validar_volumen",
                                 python_callable=_validar_volumen)

Los XCom son para valores pequeños: identificadores, contadores, rutas. Se guardan en la base de datos de metadatos de Airflow. Pasar un DataFrame por XCom es un error clásico que hincha la base de datos y degrada todo el entorno. Para datos grandes, se escribe en Cloud Storage y se pasa la ruta.

Variables son configuración global, editable desde la interfaz sin tocar código:

from airflow.models import Variable

umbral = int(Variable.get("alpinashop_umbral_pedidos_dia", default_var="500"))
gcloud composer environments run alpinashop-composer \
  --location=europe-west1 variables -- set alpinashop_umbral_pedidos_dia 500

Conexiones guardan credenciales y parámetros de acceso a sistemas externos, cifradas. Para lo que sea sensible de verdad, lo correcto es que la conexión apunte a Secret Manager (03-06) mediante el backend de secretos de Airflow, en lugar de guardar la contraseña en la base de datos de metadatos.

  1. TaskGroup, dependencias y patrones de flujo

Los operadores de dependencia:

a >> b                    # b despues de a
a << b                    # a despues de b
a >> [b, c] >> d          # b y c en paralelo, ambos tras a; d tras ambos
[a, b] >> c               # c espera a que terminen a y b

Las reglas de disparo (trigger_rule) controlan bajo qué condición se ejecuta una tarea:

Regla Se ejecuta cuando
all_success (por defecto) Todas las tareas anteriores tuvieron éxito
all_failed Todas fallaron
all_done Todas terminaron, con éxito o sin él
one_success Al menos una tuvo éxito
one_failed Al menos una falló
none_failed_min_one_success Ninguna falló y al menos una se ejecutó

La regla all_done es imprescindible para tareas de limpieza:

from airflow.providers.google.cloud.operators.dataproc import (
    DataprocDeleteClusterOperator,
)

borrar_cluster = DataprocDeleteClusterOperator(
    task_id="borrar_cluster",
    cluster_name="cluster-efimero-{{ ds_nodash }}",
    region=REGION,
    project_id=PROYECTO_DATOS,
    # CRITICO: se borra el cluster PASE LO QUE PASE.
    # Con all_success, un fallo del job dejaria el cluster encendido
    # y facturando indefinidamente.
    trigger_rule="all_done",
)

Y one_failed para notificaciones de error:

avisar_fallo = PythonOperator(
    task_id="avisar_fallo",
    python_callable=_notificar_a_slack,
    trigger_rule="one_failed",
)
[cargar_pedidos, lanzar_dataflow, refrescar_vista] >> avisar_fallo

Los TaskGroup agrupan visualmente tareas relacionadas, colapsándolas en un solo nodo en la interfaz. Con un DAG de treinta tareas, la diferencia entre un grafo legible y una maraña.

  1. Idempotencia, catchup y reprocesos

La fecha lógica. Airflow ejecuta cada DAG para una fecha lógica, disponible como {{ ds }}. Es lo que hace que el DAG sea reproducible: al reprocesar el 12 de marzo, todas las plantillas valen 2026-03-12.

La regla de oro: una tarea nunca debe usar CURRENT_DATE() ni datetime.now(). Debe usar {{ ds }}. Con CURRENT_DATE(), reprocesar el 12 de marzo recalcularía con la fecha de hoy y produciría un resultado incorrecto en silencio.

# MAL: no reproducible
"query": "SELECT ... WHERE fecha = CURRENT_DATE()"

# BIEN: reproducible
"query": "SELECT ... WHERE fecha = DATE '{{ ds }}'"

catchup. Si se pone catchup=True y el start_date es de hace tres meses, Airflow ejecutará todas las fechas pendientes al activar el DAG. Puede ser útil para rellenar un histórico, o puede lanzar noventa ejecuciones simultáneas y agotar la cuota de BigQuery. Para AlpinaShop: catchup=False, y los rellenos se hacen a mano y con control.

backfill. Reprocesar un rango explícitamente:

gcloud composer environments run alpinashop-composer \
  --location=europe-west1 dags backfill -- \
  --start-date 2026-03-10 --end-date 2026-03-15 \
  --reset-dagruns \
  alpinashop_proceso_nocturno

Ese comando solo es seguro si todas las tareas son idempotentes. Por eso el DAG usa WRITE_TRUNCATE en el staging, MERGE en la consolidación y CREATE OR REPLACE TABLE en las agregaciones. Un DAG con INSERT en lugar de MERGE duplicaría datos en cada reproceso, y el reproceso —que debería ser la herramienta de reparación— se convertiría en la causa de un problema peor.

Reejecutar una sola tarea, sin todo el DAG:

gcloud composer environments run alpinashop-composer \
  --location=europe-west1 tasks clear -- \
  --task-regex "refrescar_vista_negocio" \
  --start-date 2026-03-14 --end-date 2026-03-14 --yes \
  alpinashop_proceso_nocturno

  1. Workflows: orquestación serverless en YAML

Si Composer cuesta 350 € al mes y AlpinaShop tiene ocho tareas, hay una alternativa: Cloud Workflows, orquestación declarativa sin servidor y sin coste fijo. Se paga por paso ejecutado, y son céntimos.

Un flujo se define en YAML y se ejecuta cuando se le invoca:

# proceso-nocturno.yaml -- Proceso nocturno de AlpinaShop con Workflows
main:
  params: [entrada]
  steps:
    - inicializar:
        assign:
          - proyecto: "alpinashop-datos"
          - proyecto_prod: "alpinashop-prod"
          - region: "europe-west1"
          - dataset: "alpinashop_analitica"
          - bucket: "alpinashop-datalake"
          # Si no se pasa fecha, se usa ayer. Permite reprocesar
          # invocando el flujo con {"fecha": "2026-03-12"}.
          - fecha: ${default(map.get(entrada, "fecha"), text.substring(time.format(sys.now() - 86400), 0, 10))}
          - fecha_compacta: ${text.replace_all(fecha, "-", "")}

    # -----------------------------------------------------------------
    # 1) Exportar Cloud SQL. La API devuelve una OPERACION de larga
    #    duracion: hay que esperarla, no darla por hecha.
    # -----------------------------------------------------------------
    - exportar_cloudsql:
        call: googleapis.sqladmin.v1.instances.export
        args:
          project: ${proyecto_prod}
          instance: "alpinashop-pedidos-replica-informes"
          body:
            exportContext:
              fileType: "CSV"
              uri: ${"gs://" + bucket + "/exportaciones/" + fecha_compacta + "/pedidos.csv"}
              databases: ["tienda"]
              csvExportOptions:
                selectQuery: ${"SELECT pedido_id, creado_en, cliente_id, canal, estado, pais, ciudad, codigo_postal, metodo_pago, subtotal, descuento, iva, total FROM pedidos WHERE DATE(creado_en) = '" + fecha + "'"}
        result: operacion_export

    - esperar_export:
        call: sys.sleep
        args:
          seconds: 30

    # -----------------------------------------------------------------
    # 2) Comprobar que el fichero existe. El equivalente al sensor.
    # -----------------------------------------------------------------
    - comprobar_fichero:
        try:
          call: googleapis.storage.v1.objects.get
          args:
            bucket: ${bucket}
            object: ${"exportaciones%2F" + fecha_compacta + "%2Fpedidos.csv"}
          result: info_fichero
        retry:
          predicate: ${http.default_retry_predicate}
          max_retries: 20
          backoff:
            initial_delay: 30
            max_delay: 120
            multiplier: 1.5

    # -----------------------------------------------------------------
    # 3) Cargar en BigQuery
    # -----------------------------------------------------------------
    - cargar_bigquery:
        call: googleapis.bigquery.v2.jobs.insert
        args:
          projectId: ${proyecto}
          body:
            configuration:
              load:
                sourceUris:
                  - ${"gs://" + bucket + "/exportaciones/" + fecha_compacta + "/pedidos.csv"}
                destinationTable:
                  projectId: ${proyecto}
                  datasetId: ${dataset}
                  tableId: "pedidos_staging"
                sourceFormat: "CSV"
                writeDisposition: "WRITE_TRUNCATE"
                autodetect: true
        result: job_carga

    - esperar_carga:
        call: espera_job_bigquery
        args:
          proyecto: ${proyecto}
          job_id: ${job_carga.jobReference.jobId}
        result: estado_carga

    # -----------------------------------------------------------------
    # 4) Lanzar Dataflow con la plantilla flexible
    # -----------------------------------------------------------------
    - lanzar_dataflow:
        call: http.post
        args:
          url: ${"https://dataflow.googleapis.com/v1b3/projects/" + proyecto + "/locations/" + region + "/flexTemplates:launch"}
          auth:
            type: OAuth2
          body:
            launchParameter:
              jobName: ${"pedidos-lote-" + fecha_compacta}
              containerSpecGcsPath: "gs://alpinashop-dataflow/plantillas/pedidos-lote.json"
              parameters:
                fecha: ${fecha}
              environment:
                serviceAccountEmail: "[email protected]"
                subnetwork: ${"regions/" + region + "/subnetworks/sn-datos-euw1"}
                ipConfiguration: "WORKER_IP_PRIVATE"
        result: job_dataflow

    # -----------------------------------------------------------------
    # 5) Agregaciones EN PARALELO con la rama 'parallel'
    # -----------------------------------------------------------------
    - agregaciones:
        parallel:
          branches:
            - ventas:
                steps:
                  - consulta_ventas:
                      call: ejecuta_sql
                      args:
                        proyecto: ${proyecto}
                        sql: ${"CREATE OR REPLACE TABLE `" + proyecto + "." + dataset + ".agg_ventas_categoria_dia` PARTITION BY dia AS SELECT l.fecha_pedido AS dia, pr.categoria, COUNT(DISTINCT l.pedido_id) AS pedidos, SUM(l.cantidad) AS unidades, ROUND(SUM(l.importe_linea),2) AS ventas_eur FROM `" + proyecto + "." + dataset + ".lineas_pedido` l JOIN `" + proyecto + "." + dataset + ".productos` pr USING (sku) WHERE l.fecha_pedido >= DATE_SUB(DATE '" + fecha + "', INTERVAL 400 DAY) GROUP BY dia, pr.categoria"}
            - embudo:
                steps:
                  - consulta_embudo:
                      call: ejecuta_sql
                      args:
                        proyecto: ${proyecto}
                        sql: ${"CREATE OR REPLACE TABLE `" + proyecto + "." + dataset + ".agg_embudo_dia` PARTITION BY dia AS SELECT fecha AS dia, COUNT(DISTINCT sesion_id) AS sesiones FROM `" + proyecto + "." + dataset + ".visitas` WHERE fecha = DATE '" + fecha + "' GROUP BY dia"}

    # -----------------------------------------------------------------
    # 6) Comprobacion de calidad: si falla, se lanza un error explicito
    # -----------------------------------------------------------------
    - comprobar_calidad:
        call: ejecuta_sql
        args:
          proyecto: ${proyecto}
          sql: ${"SELECT COUNT(*) AS n FROM `" + proyecto + "." + dataset + ".agg_ventas_categoria_dia` WHERE dia = DATE '" + fecha + "'"}
        result: resultado_calidad

    - evaluar_calidad:
        switch:
          - condition: ${int(resultado_calidad.rows[0].f[0].v) == 0}
            raise: ${"CALIDAD: no hay ventas agregadas para " + fecha}

    - devolver:
        return:
          fecha: ${fecha}
          estado: "OK"
          job_dataflow: ${job_dataflow.body.job.id}

# ---------------------------------------------------------------------
# SUBFLUJOS reutilizables
# ---------------------------------------------------------------------
ejecuta_sql:
  params: [proyecto, sql]
  steps:
    - lanzar:
        call: googleapis.bigquery.v2.jobs.query
        args:
          projectId: ${proyecto}
          body:
            query: ${sql}
            useLegacySql: false
            location: "europe-west1"
            timeoutMs: 300000
        result: r
    - devolver:
        return: ${r}

espera_job_bigquery:
  params: [proyecto, job_id]
  steps:
    - consultar:
        call: googleapis.bigquery.v2.jobs.get
        args:
          projectId: ${proyecto}
          jobId: ${job_id}
        result: estado
    - evaluar:
        switch:
          - condition: ${estado.status.state == "DONE"}
            next: comprobar_error
        next: esperar
    - esperar:
        call: sys.sleep
        args:
          seconds: 15
        next: consultar
    - comprobar_error:
        switch:
          - condition: ${"errorResult" in estado.status}
            raise: ${estado.status.errorResult.message}
    - devolver:
        return: ${estado}

Despliegue y ejecución:

gcloud iam service-accounts create sa-workflows-nocturno \
  --display-name="Orquestacion nocturna con Workflows"

SA_WF="[email protected]"

for ROL in roles/bigquery.dataEditor roles/bigquery.jobUser \
           roles/dataflow.developer roles/storage.objectAdmin \
           roles/cloudsql.editor roles/logging.logWriter; do
  gcloud projects add-iam-policy-binding alpinashop-datos \
    --member="serviceAccount:${SA_WF}" --role="$ROL"
done

gcloud workflows deploy alpinashop-proceso-nocturno \
  --source=proceso-nocturno.yaml \
  --location=europe-west1 \
  --service-account="$SA_WF" \
  --labels=entorno=produccion,equipo=datos,centro-coste=analitica

# Ejecucion manual con fecha explicita (reproceso)
gcloud workflows run alpinashop-proceso-nocturno \
  --location=europe-west1 \
  --data='{"fecha":"2026-03-12"}'

# Ver el resultado
gcloud workflows executions list alpinashop-proceso-nocturno \
  --location=europe-west1 --limit=5 \
  --format="table(name.basename(), state, startTime, endTime)"

Los límites honestos de Workflows, que hay que conocer antes de comprometerse:

Límite Valor aproximado Implicación
Duración máxima de una ejecución 1 año Sin problema
Duración máxima de una llamada HTTP 30 min Un job largo hay que sondearlo, no esperarlo
Tamaño máximo de una variable 512 KB No pasar datos, solo referencias
Pasos por ejecución ~100.000 Suficiente
Interfaz visual Grafo simple Muy inferior al panel de Airflow
backfill de un rango No existe Hay que invocar en bucle desde un script
Sensores de eventos Sondeo manual con retry Menos elegante que un sensor de Airflow
Ecosistema de operadores Llamadas a API Sin los cientos de operadores de Airflow

Las tres filas que más pesan son el backfill inexistente, la observabilidad más pobre y la ausencia de operadores listos. Con ocho tareas es asumible; con ochenta, doloroso.

  1. Cloud Scheduler: el disparador por cron

Workflows no tiene planificador propio: hay que dispararlo. Eso lo hace Cloud Scheduler, un cron gestionado que cuesta prácticamente nada (los primeros trabajos son gratuitos y después son céntimos).

gcloud iam service-accounts create sa-scheduler-nocturno \
  --display-name="Disparador del proceso nocturno"

SA_SCH="[email protected]"

gcloud projects add-iam-policy-binding alpinashop-datos \
  --member="serviceAccount:${SA_SCH}" --role="roles/workflows.invoker"

gcloud scheduler jobs create http disparar-proceso-nocturno \
  --location=europe-west1 \
  --schedule="0 2 * * *" \
  --time-zone="Europe/Madrid" \
  --uri="https://workflowexecutions.googleapis.com/v1/projects/alpinashop-datos/locations/europe-west1/workflows/alpinashop-proceso-nocturno/executions" \
  --http-method=POST \
  --oauth-service-account-email="$SA_SCH" \
  --message-body='{"argument":"{}"}' \
  --max-retry-attempts=3 \
  --min-backoff=60s \
  --max-backoff=600s \
  --attempt-deadline=60s

--time-zone="Europe/Madrid" es importante: el trabajo se ejecuta a las 02:00 hora local, y Google gestiona los cambios de horario de verano. Con UTC, el proceso se desplazaría una hora dos veces al año, lo que suena inofensivo hasta que coincide con la ventana de mantenimiento de otro sistema.

Cloud Scheduler también puede publicar en Pub/Sub o invocar Cloud Run directamente, lo que lo convierte en el disparador universal de la plataforma.

Y con la notificación del bucket de 04-04, se puede montar un flujo reactivo en lugar de programado:

flowchart LR
    S["Cloud Scheduler<br/>02:00 Europe/Madrid"]
    W["Workflows<br/>proceso nocturno"]
    G["Cloud Storage<br/>fichero del transportista"]
    P["Pub/Sub<br/>imagenes-subidas"]
    F["Cloud Function<br/>06-03"]

    S -->|cron| W
    G -->|notificacion| P --> F -->|invoca| W

El flujo se dispara a las 02:00 o cuando llega un fichero, lo que ocurra. Es más robusto que una hora fija, porque no depende de que el proveedor sea puntual.

  1. La tabla de decisión y la elección de AlpinaShop

Criterio Cloud Scheduler Workflows Cloud Composer Plantillas de Dataflow
Qué resuelve "Ejecuta esto a esta hora" "Ejecuta estos pasos en este orden" Orquestación completa Un pipeline concreto
Definición Cron + destino YAML declarativo Python Parámetros
Coste fijo ~0 € 0 € ~350 €/mes 0 €
Coste variable Céntimos Céntimos por paso Trabajadores Recursos del job
Dependencias complejas No Sí, con límites Sí, sin límites No
Paralelismo No Sí (parallel) Interno
Reintentos Sí, configurable Sí, muy fino
Sensores de eventos No Sondeo manual Sí, nativos No
backfill No Manual Sí, nativo No
Observabilidad Logs Grafo simple Panel completo Interfaz de Dataflow
Ecosistema N/A APIs de Google Cientos de operadores N/A
Curva de aprendizaje Minutos Horas Días Minutos
Elegir si… Una acción periódica Pocas decenas de pasos Decenas de DAG y equipos Un pipeline suelto

La decisión de AlpinaShop: empezar con Workflows + Cloud Scheduler.

El razonamiento, que es el que hay que saber defender:

  1. El coste manda. 350 € al mes por orquestar ocho tareas es desproporcionado cuando toda la plataforma de datos cuesta 15 €. Con Workflows y Scheduler, la orquestación cuesta céntimos.
  2. La complejidad actual no lo justifica. Ocho tareas con dependencias lineales y una bifurcación en paralelo caben perfectamente en un YAML de 150 líneas. Airflow brilla con cincuenta DAG que comparten dependencias entre equipos, y AlpinaShop tiene uno.
  3. No hay pérdida funcional relevante. Se cubren dependencias, paralelismo, reintentos con retroceso, comprobación de calidad y alertas. Lo que falta —backfill nativo, sensores, panel rico— se suple con un script de invocación en bucle y las alertas de Cloud Monitoring.
  4. La migración posterior es viable. Si mañana hacen falta veinte DAG, se crea el entorno de Composer y se migra. La lógica de las tareas —consultas SQL, plantillas de Dataflow, jobs de Dataproc— no cambia: solo cambia quién las invoca. Nada de lo hecho se tira.

Los criterios objetivos para dar el salto a Composer, escritos de antemano para no discutirlo con la emoción del momento:

  • Más de 10 flujos distintos con dependencias entre ellos.
  • Necesidad recurrente de backfill de rangos largos.
  • Más de una persona manteniendo los flujos, con revisión por pull request.
  • Necesidad de operadores especializados (Salesforce, SAP, Kubernetes).
  • Que el coste del orquestador baje del 10 % del coste de la plataforma que orquesta.

Ese último criterio es el más útil y el más fácil de comprobar: cuando la plataforma de datos de AlpinaShop cueste 3.500 € al mes, Composer costará el 10 % y estará justificado. Hoy costaría el 2.300 %.

Errores Comunes y Consejos

Usar CURRENT_DATE() en lugar de {{ ds }}. Rompe la reproducibilidad en silencio: el reproceso de una fecha pasada calcula con la de hoy y nadie lo nota.

Dejar catchup=True sin querer. Al activar un DAG con start_date antiguo, se lanzan cientos de ejecuciones a la vez, se agotan las cuotas y se dispara el coste.

Sensores en modo poke con esperas largas. Ocupan trabajadores. Con varios sensores simultáneos, el entorno se bloquea a sí mismo. Usa mode="reschedule".

Tareas no idempotentes. Un INSERT en lugar de un MERGE convierte cada reintento y cada reproceso en una duplicación de datos. La idempotencia no es opcional en un DAG.

Pasar datos grandes por XCom. Se guardan en la base de datos de metadatos y la degradan. Pasa rutas, no contenidos.

Olvidar trigger_rule="all_done" en las tareas de limpieza. Un clúster efímero que no se borra porque el job falló queda encendido facturando indefinidamente.

No poner execution_timeout. Una tarea colgada retiene un trabajador para siempre y acaba bloqueando el DAG entero.

Poner lógica pesada en el cuerpo del DAG. El fichero se reevalúa cada pocos segundos por el planificador. Una consulta a una base de datos fuera de un operador se ejecuta constantemente y hunde el entorno. Todo el trabajo va dentro de operadores.

Dejar Composer encendido "por si acaso". Es el equivalente al clúster de Dataproc permanente de 04-03 y al Data Fusion permanente de 04-05: el mismo error, tres veces, y siempre la partida más cara de la factura.

Consejo: guarda los DAG en Git y despliega con CI. El bucket de DAG no es el sitio donde vive el código, es donde se copia. En 06-01 lo automatizaremos con Cloud Build.

Consejo: pon una comprobación de calidad al final de todo flujo. Es la tarea que convierte un proceso automático en uno fiable. Detectar el error a las 05:00 vale mucho más que descubrirlo en un comité.

Consejo: nombra las tareas con verbos. exportar_pedidos_cloudsql, no tarea_1. Cuando algo falle a las 03:00, el nombre es lo primero que se lee.

Ejercicios

Ejercicio 1: DAG mensual de análisis de cesta

Escribe un DAG de Airflow llamado alpinashop_analisis_mensual que se ejecute el día 1 de cada mes a las 04:00 hora de Madrid y realice: (1) comprobar que la tabla lineas_pedido tiene datos del mes anterior, fallando si no; (2) lanzar el trabajo de Dataproc Serverless cesta_media.py con las fechas del mes anterior como argumentos; (3) en paralelo, ejecutar dos consultas de BigQuery que calculen el margen por categoría y los productos sin ventas del mes; (4) refrescar la vista materializada; (5) publicar un mensaje en un topic informes-listos indicando que el informe mensual está disponible. Incluye reintentos, SLA y notificación de fallo.

Ejercicio 2: el mismo proceso en Workflows

Implementa en YAML un flujo alpinashop-informe-mensual equivalente a los pasos 1, 2 y 3 del ejercicio anterior, que acepte un parámetro mes con formato yyyy-MM y use el mes anterior si no se pasa. Debe ejecutar las dos consultas en paralelo, esperar correctamente a que termine el trabajo de Dataproc (que es una operación de larga duración) y lanzar un error explícito si la comprobación inicial no encuentra datos. Añade el trabajo de Cloud Scheduler que lo dispare.

Ejercicio 3: diagnóstico de un DAG que miente

Durante tres semanas, el DAG alpinashop_proceso_nocturno aparece en verde todas las noches. Pero el lunes, Lucía detecta que la tabla agg_ventas_categoria_dia tiene los datos del 24 de febrero repetidos en todas las particiones desde esa fecha. Investigando encuentras: la tarea consolidar_pedidos usa INSERT INTO en lugar de MERGE; la consulta de agregación filtra por WHERE l.fecha_pedido = CURRENT_DATE() - 1; la tarea comprobar_calidad_datos verifica únicamente que la tabla no esté vacía; y el 24 de febrero alguien ejecutó un backfill de las dos semanas anteriores.

Explica exactamente qué ha pasado y en qué orden, por qué el DAG aparecía en verde, y propón las correcciones concretas —con el código— para cada uno de los cuatro problemas.

Soluciones

Solución 1

"""alpinashop_analisis_mensual.py -- Analisis de cesta y margen mensual."""
from datetime import datetime, timedelta

import pendulum
from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import (
    BigQueryCheckOperator, BigQueryInsertJobOperator,
)
from airflow.providers.google.cloud.operators.dataproc import (
    DataprocCreateBatchOperator,
)
from airflow.providers.google.cloud.operators.pubsub import (
    PubSubPublishMessageOperator,
)
from airflow.utils.task_group import TaskGroup

PROYECTO = "alpinashop-datos"
DATASET = "alpinashop_analitica"
BUCKET = "alpinashop-datalake"
REGION = "europe-west1"
TZ = pendulum.timezone("Europe/Madrid")

argumentos = {
    "owner": "equipo-datos",
    "email": ["[email protected]"],
    "email_on_failure": True,
    "retries": 2,
    "retry_delay": timedelta(minutes=10),
    "retry_exponential_backoff": True,
    "execution_timeout": timedelta(hours=3),
    "sla": timedelta(hours=4),
}

with DAG(
    dag_id="alpinashop_analisis_mensual",
    default_args=argumentos,
    schedule="0 4 1 * *",                     # dia 1 de cada mes a las 04:00
    start_date=datetime(2026, 3, 1, tzinfo=TZ),
    catchup=False,
    max_active_runs=1,
    tags=["alpinashop", "mensual", "datos"],
) as dag:

    # ds del dia 1 -> el mes anterior va del dia 1 anterior al ultimo dia
    PRIMER_DIA_MES_ANT = "{{ macros.ds_format(macros.ds_add(ds, -1), '%Y-%m-%d', '%Y-%m-01') }}"
    ULTIMO_DIA_MES_ANT = "{{ macros.ds_add(ds, -1) }}"

    # 1) Comprobar que hay datos del mes anterior
    comprobar_datos = BigQueryCheckOperator(
        task_id="comprobar_datos_mes_anterior",
        location=REGION,
        use_legacy_sql=False,
        sql=f"""
            SELECT COUNT(*) > 0
            FROM `{PROYECTO}.{DATASET}.lineas_pedido`
            WHERE fecha_pedido BETWEEN DATE '{PRIMER_DIA_MES_ANT}'
                                   AND DATE '{ULTIMO_DIA_MES_ANT}'
        """,
    )

    # 2) Analisis de cesta en Dataproc Serverless
    analisis_cesta = DataprocCreateBatchOperator(
        task_id="analisis_cesta",
        region=REGION,
        project_id=PROYECTO,
        batch_id="cesta-{{ ds_nodash }}",
        batch={
            "pyspark_batch": {
                "main_python_file_uri": f"gs://{BUCKET}/jobs/cesta_media.py",
                "args": [
                    f"--fecha-desde={PRIMER_DIA_MES_ANT}",
                    f"--fecha-hasta={ULTIMO_DIA_MES_ANT}",
                ],
            },
            "runtime_config": {"version": "2.2"},
            "environment_config": {
                "execution_config": {
                    "service_account":
                        "[email protected]",
                    "subnetwork_uri": "sn-datos-euw1",
                }
            },
        },
    )

    # 3) Dos consultas en paralelo
    with TaskGroup(group_id="informes_mensuales") as informes:

        margen_categoria = BigQueryInsertJobOperator(
            task_id="margen_por_categoria",
            location=REGION,
            configuration={"query": {"useLegacySql": False, "query": f"""
                CREATE OR REPLACE TABLE `{PROYECTO}.{DATASET}.agg_margen_mes` AS
                SELECT
                  DATE '{PRIMER_DIA_MES_ANT}'                          AS mes,
                  pr.categoria,
                  SUM(l.cantidad)                                      AS unidades,
                  ROUND(SUM(l.importe_linea), 2)                       AS ventas_eur,
                  ROUND(SUM(l.cantidad * c.coste_medio_eur), 2)        AS coste_eur,
                  ROUND(SUM(l.importe_linea)
                        - SUM(l.cantidad * c.coste_medio_eur), 2)      AS margen_eur
                FROM `{PROYECTO}.{DATASET}.lineas_pedido`               AS l
                JOIN `{PROYECTO}.{DATASET}.productos`                   AS pr USING (sku)
                JOIN `{PROYECTO}.{DATASET}.costes_producto_erp`         AS c  USING (sku)
                WHERE l.fecha_pedido BETWEEN DATE '{PRIMER_DIA_MES_ANT}'
                                         AND DATE '{ULTIMO_DIA_MES_ANT}'
                GROUP BY mes, pr.categoria
            """}},
        )

        sin_ventas = BigQueryInsertJobOperator(
            task_id="productos_sin_ventas",
            location=REGION,
            configuration={"query": {"useLegacySql": False, "query": f"""
                CREATE OR REPLACE TABLE `{PROYECTO}.{DATASET}.agg_sin_ventas_mes` AS
                SELECT
                  DATE '{PRIMER_DIA_MES_ANT}' AS mes,
                  pr.sku, pr.nombre, pr.categoria, pr.precio_catalogo
                FROM `{PROYECTO}.{DATASET}.productos` AS pr
                WHERE pr.activo = TRUE
                  AND NOT EXISTS (
                    SELECT 1 FROM `{PROYECTO}.{DATASET}.lineas_pedido` AS l
                    WHERE l.sku = pr.sku
                      AND l.fecha_pedido BETWEEN DATE '{PRIMER_DIA_MES_ANT}'
                                             AND DATE '{ULTIMO_DIA_MES_ANT}'
                  )
            """}},
        )

    # 4) Refrescar la vista materializada
    refrescar = BigQueryInsertJobOperator(
        task_id="refrescar_vista",
        location=REGION,
        configuration={"query": {"useLegacySql": False, "query": f"""
            CALL BQ.REFRESH_MATERIALIZED_VIEW(
              '{PROYECTO}.{DATASET}.mv_ventas_diarias_sku')
        """}},
    )

    # 5) Avisar de que el informe esta listo
    avisar = PubSubPublishMessageOperator(
        task_id="avisar_informe_listo",
        project_id=PROYECTO,
        topic="informes-listos",
        messages=[{
            "data": b'{"informe":"mensual","estado":"completado"}',
            "attributes": {
                "tipo_evento": "informe_listo",
                "periodo": PRIMER_DIA_MES_ANT,
            },
        }],
    )

    comprobar_datos >> analisis_cesta >> informes >> refrescar >> avisar

El punto que se evalúa es el cálculo del mes anterior con macros.ds_add y macros.ds_format en lugar de con datetime.now(). Ejecutando el DAG el 1 de abril, ds vale 2026-04-01, ds_add(ds, -1) da 2026-03-31 y el formateo a %Y-%m-01 da 2026-03-01. Al reprocesar el 1 de febrero, los valores serán los de enero. Con fechas calculadas en tiempo real, el reproceso sería inútil.

Solución 2

# informe-mensual.yaml
main:
  params: [entrada]
  steps:
    - inicializar:
        assign:
          - proyecto: "alpinashop-datos"
          - dataset: "alpinashop_analitica"
          - region: "europe-west1"
          - bucket: "alpinashop-datalake"
          - mes: ${default(map.get(entrada, "mes"), text.substring(time.format(sys.now() - 2592000), 0, 7))}
          - primer_dia: ${mes + "-01"}
          - mes_compacto: ${text.replace_all(mes, "-", "")}

    - calcular_ultimo_dia:
        call: ejecuta_sql
        args:
          proyecto: ${proyecto}
          sql: ${"SELECT CAST(LAST_DAY(DATE '" + primer_dia + "') AS STRING) AS ultimo"}
        result: r_ultimo

    - asignar_ultimo:
        assign:
          - ultimo_dia: ${r_ultimo.rows[0].f[0].v}

    # 1) Comprobar que hay datos
    - comprobar_datos:
        call: ejecuta_sql
        args:
          proyecto: ${proyecto}
          sql: ${"SELECT COUNT(*) AS n FROM `" + proyecto + "." + dataset + ".lineas_pedido` WHERE fecha_pedido BETWEEN DATE '" + primer_dia + "' AND DATE '" + ultimo_dia + "'"}
        result: r_datos

    - evaluar_datos:
        switch:
          - condition: ${int(r_datos.rows[0].f[0].v) == 0}
            raise: ${"Sin lineas de pedido para el mes " + mes}

    # 2) Dataproc Serverless: operacion de larga duracion
    - lanzar_cesta:
        call: http.post
        args:
          url: ${"https://dataproc.googleapis.com/v1/projects/" + proyecto + "/locations/" + region + "/batches?batchId=cesta-" + mes_compacto}
          auth:
            type: OAuth2
          body:
            pysparkBatch:
              mainPythonFileUri: ${"gs://" + bucket + "/jobs/cesta_media.py"}
              args:
                - ${"--fecha-desde=" + primer_dia}
                - ${"--fecha-hasta=" + ultimo_dia}
            runtimeConfig:
              version: "2.2"
            environmentConfig:
              executionConfig:
                serviceAccount: "[email protected]"
                subnetworkUri: "sn-datos-euw1"
        result: r_batch

    - esperar_cesta:
        call: espera_batch
        args:
          proyecto: ${proyecto}
          region: ${region}
          batch_id: ${"cesta-" + mes_compacto}

    # 3) Dos consultas en paralelo
    - informes:
        parallel:
          branches:
            - margen:
                steps:
                  - q_margen:
                      call: ejecuta_sql
                      args:
                        proyecto: ${proyecto}
                        sql: ${"CREATE OR REPLACE TABLE `" + proyecto + "." + dataset + ".agg_margen_mes` AS SELECT DATE '" + primer_dia + "' AS mes, pr.categoria, SUM(l.cantidad) AS unidades, ROUND(SUM(l.importe_linea),2) AS ventas_eur, ROUND(SUM(l.cantidad*c.coste_medio_eur),2) AS coste_eur FROM `" + proyecto + "." + dataset + ".lineas_pedido` l JOIN `" + proyecto + "." + dataset + ".productos` pr USING (sku) JOIN `" + proyecto + "." + dataset + ".costes_producto_erp` c USING (sku) WHERE l.fecha_pedido BETWEEN DATE '" + primer_dia + "' AND DATE '" + ultimo_dia + "' GROUP BY mes, pr.categoria"}
            - sin_ventas:
                steps:
                  - q_sin_ventas:
                      call: ejecuta_sql
                      args:
                        proyecto: ${proyecto}
                        sql: ${"CREATE OR REPLACE TABLE `" + proyecto + "." + dataset + ".agg_sin_ventas_mes` AS SELECT DATE '" + primer_dia + "' AS mes, pr.sku, pr.nombre, pr.categoria FROM `" + proyecto + "." + dataset + ".productos` pr WHERE pr.activo = TRUE AND NOT EXISTS (SELECT 1 FROM `" + proyecto + "." + dataset + ".lineas_pedido` l WHERE l.sku = pr.sku AND l.fecha_pedido BETWEEN DATE '" + primer_dia + "' AND DATE '" + ultimo_dia + "')"}

    - devolver:
        return:
          mes: ${mes}
          estado: "OK"

ejecuta_sql:
  params: [proyecto, sql]
  steps:
    - lanzar:
        call: googleapis.bigquery.v2.jobs.query
        args:
          projectId: ${proyecto}
          body:
            query: ${sql}
            useLegacySql: false
            location: "europe-west1"
            timeoutMs: 300000
        result: r
    - devolver:
        return: ${r}

# Sondeo de la operacion de larga duracion: Workflows NO puede
# esperar 20 minutos en una sola llamada HTTP (limite de 30 min,
# y ademas el batch podria tardar mas).
espera_batch:
  params: [proyecto, region, batch_id]
  steps:
    - consultar:
        call: http.get
        args:
          url: ${"https://dataproc.googleapis.com/v1/projects/" + proyecto + "/locations/" + region + "/batches/" + batch_id}
          auth:
            type: OAuth2
        result: estado
    - evaluar:
        switch:
          - condition: ${estado.body.state == "SUCCEEDED"}
            return: ${estado.body}
          - condition: ${estado.body.state == "FAILED"}
            raise: ${"Batch de Dataproc fallido: " + default(map.get(estado.body, "stateMessage"), "sin detalle")}
          - condition: ${estado.body.state == "CANCELLED"}
            raise: "Batch de Dataproc cancelado"
    - esperar:
        call: sys.sleep
        args:
          seconds: 30
        next: consultar
gcloud workflows deploy alpinashop-informe-mensual \
  --source=informe-mensual.yaml --location=europe-west1 \
  --service-account="[email protected]"

gcloud scheduler jobs create http disparar-informe-mensual \
  --location=europe-west1 \
  --schedule="0 4 1 * *" \
  --time-zone="Europe/Madrid" \
  --uri="https://workflowexecutions.googleapis.com/v1/projects/alpinashop-datos/locations/europe-west1/workflows/alpinashop-informe-mensual/executions" \
  --http-method=POST \
  --oauth-service-account-email="[email protected]" \
  --message-body='{"argument":"{}"}' \
  --max-retry-attempts=3

Lo que se evalúa: el subflujo espera_batch con sondeo y sys.sleep. Es la diferencia clave con Airflow, donde DataprocCreateBatchOperator espera solo. En Workflows hay que implementar el bucle a mano, y hay que hacerlo bien: comprobar los tres estados finales (SUCCEEDED, FAILED, CANCELLED), no solo el de éxito, porque si no un batch fallido dejaría el flujo sondeando eternamente.

Solución 3

Qué ha pasado, en orden cronológico:

Paso 1 — El backfill del 24 de febrero fue la causa desencadenante. Se reprocesaron dos semanas. Para cada una de las 14 fechas lógicas, el DAG ejecutó consolidar_pedidos, que usa INSERT INTO. Como las tablas ya contenían esos pedidos de las ejecuciones originales, cada pedido de esas dos semanas quedó duplicado en pedidos.

Paso 2 — La agregación amplificó el problema y lo congeló. La consulta filtra por WHERE l.fecha_pedido = CURRENT_DATE() - 1. Durante el backfill, CURRENT_DATE() era el 24 de febrero en las catorce ejecuciones, porque CURRENT_DATE() no sabe nada de la fecha lógica. Es decir: las catorce ejecuciones calcularon el mismo día, el 23 de febrero, y escribieron catorce veces el mismo resultado sobre la misma partición.

Paso 3 — Y siguió mal cada noche. Aquí está el detalle que explica las tres semanas: como la consulta usa CREATE OR REPLACE TABLE sin filtro por fecha lógica en la escritura, cada ejecución nocturna posterior sobrescribió la tabla con lo calculado para "ayer" según CURRENT_DATE(). Eso debería haber ido actualizándose... salvo que la tabla resultante conserva el histórico previo tal como quedó tras el backfill, y las particiones antiguas quedaron congeladas con el dato del 24 de febrero. De ahí que Lucía vea el 24 de febrero repetido en todas las particiones desde entonces.

Por qué el DAG aparecía en verde. Porque ninguna tarea falló. INSERT INTO no da error al insertar duplicados: no hay clave primaria en BigQuery. La consulta de agregación se ejecutó correctamente. Y la comprobación de calidad solo verificaba que la tabla no estuviera vacía —y no lo estaba: estaba llena de datos incorrectos. Verde no significa correcto; significa que no hubo excepciones. Esa distinción es la lección del ejercicio.

Correcciones, una por problema:

A) INSERT INTOMERGE. Hace la consolidación idempotente:

MERGE `alpinashop-datos.alpinashop_analitica.pedidos` AS d
USING (SELECT * FROM `alpinashop-datos.alpinashop_analitica.pedidos_staging`) AS o
ON d.pedido_id = o.pedido_id AND d.fecha_pedido = o.fecha_pedido
WHEN MATCHED THEN UPDATE SET estado = o.estado, total_pedido = o.total_pedido
WHEN NOT MATCHED THEN INSERT ROW

B) CURRENT_DATE(){{ ds }}, y escritura solo de la partición correspondiente. Con CREATE OR REPLACE TABLE se reescribe la tabla entera; lo correcto es escribir únicamente la partición de la fecha lógica:

agregar_ventas = BigQueryInsertJobOperator(
    task_id="ventas_por_categoria",
    location=REGION,
    configuration={
        "query": {
            "useLegacySql": False,
            # Escribe SOLO la particion de la fecha logica
            "destinationTable": {
                "projectId": PROYECTO,
                "datasetId": DATASET,
                "tableId": "agg_ventas_categoria_dia${{ ds_nodash }}",
            },
            "writeDisposition": "WRITE_TRUNCATE",
            "timePartitioning": {"type": "DAY", "field": "dia"},
            "query": f"""
                SELECT
                  l.fecha_pedido                 AS dia,
                  pr.categoria,
                  COUNT(DISTINCT l.pedido_id)    AS pedidos,
                  SUM(l.cantidad)                AS unidades,
                  ROUND(SUM(l.importe_linea), 2) AS ventas_eur
                FROM `{PROYECTO}.{DATASET}.lineas_pedido` AS l
                JOIN `{PROYECTO}.{DATASET}.productos`     AS pr USING (sku)
                WHERE l.fecha_pedido = DATE '{{{{ ds }}}}'
                GROUP BY dia, pr.categoria
            """,
        }
    },
)

El decorador $ en agg_ventas_categoria_dia${{ ds_nodash }} es el decorador de partición de BigQuery: WRITE_TRUNCATE afecta solo a esa partición, no a la tabla. Así, reprocesar el 12 de marzo corrige el 12 de marzo y no toca ningún otro día.

C) Comprobación de calidad de verdad. No "hay filas", sino reglas de negocio:

comprobar_calidad = BigQueryCheckOperator(
    task_id="comprobar_calidad_datos",
    location=REGION,
    use_legacy_sql=False,
    sql=f"""
        WITH dia AS (
          SELECT
            SUM(pedidos)                      AS pedidos_dia,
            SUM(ventas_eur)                   AS ventas_dia,
            COUNT(*)                          AS filas
          FROM `{PROYECTO}.{DATASET}.agg_ventas_categoria_dia`
          WHERE dia = DATE '{{{{ ds }}}}'
        ),
        control_duplicados AS (
          SELECT COUNT(*) AS duplicados FROM (
            SELECT pedido_id
            FROM `{PROYECTO}.{DATASET}.pedidos`
            WHERE fecha_pedido = DATE '{{{{ ds }}}}'
            GROUP BY pedido_id
            HAVING COUNT(*) > 1
          )
        )
        SELECT
          d.filas       > 0    AND
          d.pedidos_dia > 0    AND
          d.pedidos_dia < 500  AND
          d.ventas_dia  > 0    AND
          c.duplicados  = 0
        FROM dia AS d CROSS JOIN control_duplicados AS c
    """,
)

La condición duplicados = 0 es la que habría detectado el problema la misma noche del 24 de febrero, en lugar de tres semanas después.

D) Procedimiento de reparación. Corregir el código no arregla los datos ya corrompidos:

-- 1) Deduplicar la tabla de pedidos conservando la ultima version
CREATE OR REPLACE TABLE `alpinashop-datos.alpinashop_analitica.pedidos`
PARTITION BY fecha_pedido
CLUSTER BY estado, canal AS
SELECT * EXCEPT(rn) FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY pedido_id, fecha_pedido
                       ORDER BY momento_pedido DESC) AS rn
  FROM `alpinashop-datos.alpinashop_analitica.pedidos`
)
WHERE rn = 1;
# 2) Reprocesar las agregaciones del periodo afectado, ya con el DAG corregido
gcloud composer environments run alpinashop-composer \
  --location=europe-west1 dags backfill -- \
  --start-date 2026-02-10 --end-date 2026-03-16 \
  --reset-dagruns alpinashop_proceso_nocturno

Ese backfill ahora sí es seguro, porque tras las correcciones A y B todas las tareas son idempotentes. Antes, era la causa del problema. La moraleja completa del ejercicio: el backfill no es peligroso; lo peligroso es un DAG no idempotente, y el backfill solo lo revela.

Conclusión

Has visto por qué cron no es un orquestador: responde a "¿qué hora es?" cuando la pregunta correcta es "¿qué puede ejecutarse ahora, dado lo que ha pasado?". Y has visto el precio concreto de esa confusión: una exportación que tarda veinte minutos de más y un comité de dirección mirando ventas falsas.

Conoces los cinco aportes de un orquestador real —dependencias explícitas, gestión de fallos, observabilidad, reproducibilidad y notificación— y las tres herramientas de Google Cloud que los cubren en distinto grado.

Has aprendido Cloud Composer, Airflow gestionado, con sus conceptos: DAG como grafo acíclico, tareas, operadores que definen qué hace cada una, y sensores que esperan a que algo ocurra en lugar de apostar por una hora. Has escrito el DAG completo del proceso nocturno de AlpinaShop: exportación desde la réplica de lectura, sensor en modo reschedule que espera al fichero de verdad, carga con WRITE_TRUNCATE en staging, MERGE idempotente hacia la tabla final, lanzamiento de la plantilla flexible de Dataflow esperando a que termine, dos agregaciones en paralelo dentro de un TaskGroup, refresco de la vista materializada y una comprobación de calidad que hace fallar el DAG antes de que nadie mire un informe erróneo. Con reintentos exponenciales, execution_timeout, SLA y notificación.

Conoces los operadores de Google Cloud, los XCom para valores pequeños —nunca para datos—, las variables y conexiones con Secret Manager detrás, las reglas de disparo con all_done para las limpiezas que deben ejecutarse pase lo que pase, y la regla de oro de la reproducibilidad: {{ ds }} siempre, CURRENT_DATE() jamás, porque de ella depende que un backfill repare en lugar de estropear.

Y has puesto el coste sobre la mesa sin adornos: unos 350 € al mes por un entorno pequeño, frente a los 15 € que cuesta toda la plataforma de datos de AlpinaShop. El orquestador costaría veinte veces más que lo orquestado. Por eso has montado la alternativa: Workflows, orquestación declarativa en YAML sin coste fijo, con ramas paralelas, reintentos con retroceso, subflujos reutilizables y sondeo explícito de las operaciones largas; disparada por Cloud Scheduler con la zona horaria de Madrid para que los cambios de horario no desplacen el proceso. Conoces sus límites reales —sin backfill nativo, sin sensores, sin panel rico, sin catálogo de operadores— y sabes que con ocho tareas son asumibles y con ochenta no.

La decisión de AlpinaShop queda fijada y argumentada: empezar con Workflows y Cloud Scheduler, con criterios objetivos escritos de antemano para saltar a Composer —más de diez flujos interdependientes, necesidad recurrente de backfill, varias personas manteniéndolos, o el criterio más limpio de todos: cuando el orquestador baje del 10 % del coste de lo que orquesta—. Y sabiendo que la migración no tirará nada, porque la lógica vive en las consultas, las plantillas y los jobs; solo cambia quién los invoca.

Con esto, la maquinaria está completa. Los datos entran solos, se transforman solos, se agregan solos y se comprueban solos, todas las noches, con alertas si algo falla.

Y aparece el último problema del módulo, que no es técnico. Hay ahora dieciséis tablas en alpinashop_analitica, media docena de vistas, un lago con Parquet, tablas de staging, agregados y una tabla llamada pedidos_evento que nadie recuerda por qué existe. Cuando alguien de marketing pregunta "¿dónde está el dato de conversión?", nadie sabe contestar sin abrir BigQuery y buscar. Nadie ha escrito qué significa exactamente ventas_eur —si incluye IVA, si descuenta devoluciones—. Nadie sabe si email_cliente sigue apareciendo en algún sitio donde no debería. Y el cuadro de mando que Lucía prometió a dirección sigue sin existir: los datos están perfectos y nadie los ve.

En 04-07, Dataplex y Looker Studio, cerraremos el módulo con eso. Verás qué es el gobierno del dato y por qué una pyme también lo necesita: catálogo, linaje, calidad, clasificación y ciclo de vida. Montarás lagos, zonas y activos en Dataplex, documentarás alpinashop_analitica con etiquetas, definirás reglas de calidad que se comprueban solas —total_pedido nunca negativo, sku siempre presente—, y usarás Sensitive Data Protection para descubrir dónde hay emails y teléfonos y desidentificarlos antes de exponerlos, con la advertencia de RGPD que corresponde. Y por fin construirás en Looker Studio el cuadro de mando de dirección de AlpinaShop, con sus buenas prácticas de diseño, de coste y —sobre todo— de permisos, incluyendo el error clásico de "credenciales del propietario" que convierte un informe compartido en una fuga de datos.

Curso de Google Cloud Platform (GCP)

Módulo 1: Introducción a Google Cloud Platform

Módulo 2: Servicios principales de GCP

Módulo 3: Redes y seguridad

Módulo 4: Datos y análisis

Módulo 5: Aprendizaje automático e IA

Módulo 6: DevOps y monitoreo

Módulo 7: Temas avanzados de GCP

Módulo 8: Proyecto final

© Copyright 2026. Todos los derechos reservados