AlpinaShop ja té una plataforma de dades completa. I té, exactament per això, un problema nou.

Cada nit cal fer això: exportar les comandes del dia des d'alpinashop-pedidos a Cloud Storage, carregar-les a BigQuery, llançar el pipeline de Dataflow que neteja i agrega, executar les consultes d'agregació, refrescar la vista de negoci que consumeix el tauler de control, i —el primer dia de cada mes— llançar a més la feina de Spark que recalcula la matriu de productes comprats junts i el pipeline de Data Fusion amb el fitxer del transportista.

Són vuit processos amb dependències reals entre ells. No té sentit carregar a BigQuery un fitxer que encara s'està escrivint, ni agregar sobre dades a mitges, ni refrescar la vista amb la meitat de les comandes del dia.

Ara mateix això ho resol 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 totes les nits, fins a la primera en què l'exportació triga quaranta minuts en comptes de vint perquè hi va haver campanya. Llavors, a les 02:30, cargar_bigquery.sh llegeix un fitxer incomplet. A les 03:00, Dataflow processa dades parcials. A les 04:30, la vista de negoci es refresca amb la meitat de les comandes. A les 09:00, direcció obre el tauler de control i veu que les vendes d'ahir van caure un 45 %. Es convoca una reunió d'urgència. A les 11:30, la Marta descobreix que la dada és falsa.

Ningú no se'n va assabentar perquè cron no sap si una cosa ha funcionat. Només sap quina hora és.

En aquesta lliçó veuràs què fa un orquestrador de veritat, construiràs el DAG complet del procés nocturn d'AlpinaShop a Cloud Composer, coneixeràs el seu cost real —que és l'argument decisiu per a una pime—, i muntaràs l'alternativa lleugera amb Workflows i Cloud Scheduler, que és per on AlpinaShop començarà.

Contingut

  1. Per què cron no és un orquestrador
  2. Què fa un orquestrador de veritat
  3. Cloud Composer: Apache Airflow gestionat
  4. Conceptes d'Airflow: DAG, tasca, operador, sensor
  5. L'entorn alpinashop-composer i el seu cost
  6. El DAG del procés nocturn, línia a línia
  7. Operadors de Google Cloud
  8. XComs, variables i connexions
  9. TaskGroup, dependències i patrons de flux
  10. Idempotència, catchup i reprocessos
  11. Workflows: orquestració serverless en YAML
  12. Cloud Scheduler: el disparador per cron
  13. La taula de decisió i l'elecció d'AlpinaShop

  1. Per què cron no és un orquestrador

cron respon a una sola pregunta: quina hora és?. Un orquestrador en respon una de molt diferent: què es pot executar ara, atès el que ha passat?.

Situació Amb cron Amb un orquestrador
La tasca A triga més del previst B arrenca igualment, amb dades incompletes B espera que A acabi bé
La tasca A falla B arrenca igualment, sobre res B no arrenca; s'avisa
Fallada transitòria de xarxa El procés mor Reintent automàtic amb espera
Es va executar ahir a la nit? Mirar registres per SSH Tauler amb l'historial complet
Reprocessar el 12 de març Script manual amb paràmetres a mà backfill d'aquella data
Tres tasques independents En sèrie, sumant temps En paral·lel
Ningú no s'assabenta d'una fallada Correcte: ningú no se n'assabenta Alerta configurada
La VM del cron cau No s'executa res i ningú no ho sap Servei gestionat amb reintents
Quant triga cada pas? No se sap Mètriques per tasca i històric

Hi ha una fila especialment insidiosa: la VM del cron. És una màquina que algú va crear fa tres anys, que ningú no apedaça, que té els scripts a /opt sense control de versions, amb credencials en fitxers .env, i de l'existència de la qual només se'n recorden dues persones. Quan aquesta VM mor, mor en silenci i la plataforma de dades deixa d'actualitzar-se durant dies.

La segona fila insidiosa és la de l'espera a ull. El cron de dalt assumeix que l'exportació triga menys de 30 minuts. Aquest marge és una aposta, i les apostes es perden justament el dia de més volum, que és el dia en què les dades més importen.

  1. Què fa un orquestrador de veritat

Un orquestrador de fluxos de treball aporta cinc coses:

Dependències explícites. Es declara que B depèn d'A, i el sistema garanteix l'ordre. No hi ha hores calculades a ull: si A triga deu minuts o dues hores, B arrenca quan A acaba.

Gestió de fallades. Reintents amb espera creixent, nombre màxim d'intents, i què fer si tot i així falla: parar, continuar amb el que no en depengui, o executar una tasca de neteja.

Observabilitat. Un tauler on es veu cada execució, cada tasca, la seva durada, els seus registres i el seu historial. Respondre a "va funcionar ahir a la nit?" costa una ullada.

Reproductibilitat. Poder reexecutar el flux d'una data passada amb els paràmetres d'aquella data, sense editar res.

Notificació. Algú s'assabenta quan alguna cosa falla, i se n'assabenta a temps.

Google Cloud ofereix tres eines per a això, i són complementàries:

Eina Què és Model
Cloud Scheduler Un cron gestionat Dispara una acció a una hora
Workflows Orquestració serverless declarativa Encadena passos en YAML, sense servidor
Cloud Composer Apache Airflow gestionat Orquestració completa amb Python

  1. Cloud Composer: Apache Airflow gestionat

Apache Airflow és l'estàndard de facto de l'orquestració de dades. Va néixer a Airbnb el 2014, és open source, i la seva idea central és que els fluxos de treball es defineixen com a codi Python.

Això últim és la seva gran virtut. Un flux definit en Python es versiona a Git, es revisa en un pull request, es prova, es genera dinàmicament amb bucles, i admet tota la lògica del llenguatge. Davant d'una interfície visual d'arrossegar caixes, un DAG en Python és infinitament més mantenible quan l'equip sap programar.

Cloud Composer és Airflow gestionat per Google: el planificador, la base de dades de metadades, els treballadors i la interfície web funcionant sobre GKE, amb integració d'IAM, Cloud Logging i Cloud Monitoring, i amb els operadors de Google Cloud ja instal·lats.

Versions vigents el 2026: Composer 3, amb Airflow 2.x i 3.x. Composer 3 va simplificar força l'arquitectura respecte a Composer 2 —menys infraestructura visible, escalat més fi— però el model de cost continua sent el mateix en allò essencial, i aquest és el punt que cal mirar abans que cap altre.

  1. Conceptes d'Airflow: DAG, tasca, operador, sensor

DAG (graf acíclic dirigit) és el flux de treball complet. Dirigit perquè les dependències tenen sentit; acíclic perquè no hi pot haver bucles: si A depèn de B i B d'A, res no podria començar mai.

Tasca (task) és un node del DAG: una unitat de treball.

Operador (operator) és la plantilla que defineix què fa una tasca. Airflow en porta centenars: executar Bash, cridar una funció Python, llançar una consulta de BigQuery, crear un clúster de Dataproc, enviar un correu.

Sensor és un operador especial que espera que passi alguna cosa: que aparegui un fitxer en un bucket, que una taula tingui dades, que una API respongui. És la peça que resol el problema del cron.

Execució (DAG run) és una instància concreta del DAG per a una data 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

Fixa't que agregar_ventas i agregar_visitas s'executen en paral·lel: totes dues depenen de lanzar_dataflow i cap no depèn de l'altra. Airflow ho dedueix del graf sense que calgui dir-ho. Amb cron caldria decidir un ordre i sumar els temps.

  1. L'entorn alpinashop-composer i el seu cost

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

gcloud iam service-accounts create sa-composer \
  --display-name="Cloud Composer d 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ó triga de 20 a 30 minuts.

I ara la conversa incòmoda, que cal tenir abans d'escriure una línia de DAG.

Component Cost aproximat mensual (verificar a la documentació oficial)
Entorn small (planificador, servidor web, base de dades) ~250-350 €
Treballadors addicionals sota càrrega Variable
Emmagatzematge del bucket de DAG i registres Cèntims
Total realista d'un entorn petit ~300-400 €/mes

Tres-cents euros al mes, s'executi un DAG o cent. És un servei permanentment encès: el planificador ha d'estar viu per saber quina hora és.

Posem això en context amb la resta de la plataforma d'AlpinaShop:

Servei Cost mensual estimat
BigQuery (emmagatzematge + consultes) ~10 €
Cloud Storage (60 GB de catàleg) ~2 €
Pub/Sub Cèntims
Dataflow per lots (nocturn) ~2 €
Dataproc Serverless (mensual) ~1 €
Cloud Composer ~350 €

L'orquestrador costaria vint vegades més que tot el que orquestra. Això no és un argument contra Composer: és un argument contra utilitzar-lo en el moment equivocat. Per a una empresa amb 200 DAG, 50 enginyers i dependències entre equips, 350 € és irrisori davant del valor. Per a AlpinaShop, amb vuit tasques nocturnes, és desproporcionat.

Tot i així construirem el DAG complet, per tres raons: perquè Airflow és l'estàndard del sector i cal saber-lo; perquè l'exercici de modelar les dependències és vàlid per a qualsevol orquestrador; i perquè el dia que AlpinaShop creixi, aquesta serà l'eina. Al final de la lliçó tornarem a la decisió.

  1. El DAG del procés nocturn, línia a línia

Els DAG es despleguen copiant-los 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}/"

I el DAG:

"""
proceso_nocturno.py -- Proces nocturn de dades d AlpinaShop.

Flux:
  exportar Cloud SQL -> esperar fitxer -> carregar BigQuery -> Dataflow
  -> agregacions (paral·leles) -> refrescar vista -> comprovar qualitat

Desplegament: 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

PROJECTE_DADES = "alpinashop-datos"
PROJECTE_PROD = "alpinashop-prod"
DATASET = "alpinashop_analitica"
BUCKET = "alpinashop-datalake"
REGIO = "europe-west1"
ZONA_HORARIA = pendulum.timezone("Europe/Madrid")
# ---------------------------------------------------------------------------
# ARGUMENTS PER DEFECTE: s apliquen a TOTES les tasques del DAG.
# Definir-los aqui evita repetir-los a cada operador.
# ---------------------------------------------------------------------------
arguments_per_defecte = {
    "owner": "equipo-datos",
    "depends_on_past": False,      # una execucio no espera l anterior
    "email": ["[email protected]"],
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 3,                  # 3 reintents davant d una fallada
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,   # 5, 10, 20 min: no matxucar l origen
    "max_retry_delay": timedelta(minutes=30),
    "execution_timeout": timedelta(hours=2),   # cap tasca eterna
    "sla": timedelta(hours=3),     # avis si el DAG no acaba en 3 h
}

Cadascun d'aquests paràmetres evita un incident concret:

  • retries amb retrocés exponencial: una fallada transitòria de xarxa no trenca la nit, i els reintents no enfonsen un sistema ja saturat.
  • execution_timeout: una tasca penjada no bloqueja el DAG indefinidament ni consumeix un treballador per sempre.
  • sla: si a les 05:00 el procés no ha acabat, algú se n'assabenta abans que direcció obri l'informe.
  • depends_on_past: False: cada nit és independent. Si es posa a True, una fallada del dilluns bloquejaria el dimarts, el dimecres i tota la resta, cosa que gairebé mai no és el que es vol.
with DAG(
    dag_id="alpinashop_proceso_nocturno",
    description="Proces nocturn: Cloud SQL -> BigQuery -> Dataflow -> agregats",
    default_args=arguments_per_defecte,
    # Tots els dies a les 02:00 hora de Madrid.
    # Airflow treballa internament en UTC; la zona horaria evita el desfasament
    # d una hora als canvis d horari d estiu.
    schedule="0 2 * * *",
    start_date=datetime(2026, 3, 1, tzinfo=ZONA_HORARIA),
    catchup=False,             # veure apartat 10
    max_active_runs=1,         # mai dues nits encavalcades
    tags=["alpinashop", "produccion", "datos"],
    doc_md=__doc__,            # la docstring es veu a la interficie
) as dag:

    inici = EmptyOperator(task_id="inicio")
    # -----------------------------------------------------------------------
    # 1) EXPORTAR Cloud SQL a Cloud Storage
    #
    # {{ ds }} es una plantilla Jinja: Airflow la substitueix per la data
    # logica de l execucio (yyyy-MM-dd). Es LA peca que fa el DAG
    # reproduible: en reprocessar el 12 de marc, {{ ds }} val 2026-03-12
    # i el fitxer de sortida i la consulta apunten a aquell dia, no a avui.
    # -----------------------------------------------------------------------
    exportar_comandes = CloudSQLExportInstanceOperator(
        task_id="exportar_pedidos_cloudsql",
        project_id=PROJECTE_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 }}'"
                    )
                },
            }
        },
    )

S'exporta des de la rèplica de lectura, no des de la instància principal. És la mateixa disciplina de 02-03 i 04-01: els processos analítics no toquen la base de dades que atén les compres.

    # -----------------------------------------------------------------------
    # 2) SENSOR: esperar que el fitxer existeixi de veritat.
    #
    # Aquesta tasca es la que resol el problema del cron. L exportacio
    # es asincrona: l operador anterior la llanca, pero el fitxer pot
    # trigar. En comptes de "esperem 30 minuts i creuem els dits",
    # es COMPROVA.
    # -----------------------------------------------------------------------
    esperar_fitxer = GCSObjectExistenceSensor(
        task_id="esperar_fichero_exportado",
        bucket=BUCKET,
        object="exportaciones/{{ ds_nodash }}/pedidos.csv",
        # 'reschedule': allibera el treballador entre comprovacions en comptes
        # d ocupar-lo esperant. Imprescindible en esperes llargues.
        mode="reschedule",
        poke_interval=60,          # comprovar cada minut
        timeout=60 * 60,           # rendir-se a l hora
    )

El mode reschedule davant de poke és un detall amb conseqüències reals: en mode poke, el sensor ocupa un treballador durant tota l'espera. Amb quatre sensors esperant una hora en un entorn petit, no queda cap treballador lliure i el DAG es bloqueja a si mateix. En mode reschedule, el sensor s'adorm i allibera el forat.

    # -----------------------------------------------------------------------
    # 3) CARREGAR a BigQuery
    # -----------------------------------------------------------------------
    carregar_comandes = GCSToBigQueryOperator(
        task_id="cargar_pedidos_bigquery",
        bucket=BUCKET,
        source_objects=["exportaciones/{{ ds_nodash }}/pedidos.csv"],
        destination_project_dataset_table=f"{PROJECTE_DADES}.{DATASET}.pedidos_staging",
        source_format="CSV",
        skip_leading_rows=0,
        field_delimiter=",",
        null_marker="\\N",
        # WRITE_TRUNCATE a la taula de staging: cada nit es reemplaca.
        # Aixo fa la tasca IDEMPOTENT: reexecutar-la no duplica res.
        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=REGIO,
    )

    # -----------------------------------------------------------------------
    # 4) MERGE del staging a la taula final: idempotent per disseny
    # -----------------------------------------------------------------------
    consolidar_comandes = BigQueryInsertJobOperator(
        task_id="consolidar_pedidos",
        location=REGIO,
        configuration={
            "query": {
                "query": f"""
                    MERGE `{PROJECTE_DADES}.{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 `{PROJECTE_DADES}.{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 és la clau de la idempotència: si la tasca es reexecuta, les comandes existents s'actualitzen i les noves s'insereixen. Mai no es duplica res. És exactament el mateix principi que l'Upsert de Data Fusion a 04-05 i la idempotència dels consumidors de Pub/Sub a 04-04. El mateix concepte apareix a les quatre lliçons perquè és la propietat que fa que un sistema de dades es pugui operar sense por.

    # -----------------------------------------------------------------------
    # 5) DATAFLOW: plantilla flexible creada a 04-02
    # -----------------------------------------------------------------------
    llancar_dataflow = DataflowStartFlexTemplateOperator(
        task_id="lanzar_pipeline_dataflow",
        location=REGIO,
        project_id=PROJECTE_DADES,
        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/{REGIO}/subnetworks/sn-datos-euw1",
                    "ipConfiguration": "WORKER_IP_PRIVATE",
                    "maxWorkers": 10,
                    "additionalUserLabels": {
                        "entorno": "produccion", "equipo": "datos",
                    },
                },
            }
        },
        # L operador ESPERA que el job acabi abans de donar la tasca
        # per completada. Sense aixo, el DAG continuaria amb dades a mitges.
        wait_until_finished=True,
    )

    # -----------------------------------------------------------------------
    # 6) AGREGACIONS EN PARAL·LEL, agrupades perque el graf sigui llegible
    # -----------------------------------------------------------------------
    with TaskGroup(group_id="agregaciones") as grup_agregacions:

        agregar_vendes = BigQueryInsertJobOperator(
            task_id="ventas_por_categoria",
            location=REGIO,
            configuration={
                "query": {
                    "query": f"""
                        CREATE OR REPLACE TABLE
                          `{PROJECTE_DADES}.{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 `{PROJECTE_DADES}.{DATASET}.lineas_pedido` AS l
                        JOIN `{PROJECTE_DADES}.{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_visites = BigQueryInsertJobOperator(
            task_id="embudo_conversion",
            location=REGIO,
            configuration={
                "query": {
                    "query": f"""
                        CREATE OR REPLACE TABLE
                          `{PROJECTE_DADES}.{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 `{PROJECTE_DADES}.{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 materialitzada que consumeix el tauler de control
    # -----------------------------------------------------------------------
    refrescar_vista = BigQueryInsertJobOperator(
        task_id="refrescar_vista_negocio",
        location=REGIO,
        configuration={
            "query": {
                "query": f"""
                    CALL BQ.REFRESH_MATERIALIZED_VIEW(
                      '{PROJECTE_DADES}.{DATASET}.mv_ventas_diarias_sku')
                """,
                "useLegacySql": False,
            }
        },
    )

    # -----------------------------------------------------------------------
    # 8) COMPROVACIO DE QUALITAT: l ultima defensa abans de l informe.
    #
    # Si aquesta tasca falla, el DAG queda marcat com a fallit i salta
    # l alerta. Es preferible una alerta a les 05:00 que un comite de
    # direccio mirant xifres falses a les 09:00.
    # -----------------------------------------------------------------------
    comprovar_qualitat = BigQueryCheckOperator(
        task_id="comprobar_calidad_datos",
        location=REGIO,
        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 `{PROJECTE_DADES}.{DATASET}.agg_ventas_categoria_dia`
              WHERE dia = DATE '{{{{ ds }}}}'
              GROUP BY dia
            )
        """,
    )

    def _avisar_exit(**context):
        """Callback informatiu: en produccio enviaria a Slack o Chat."""
        print(f"Proces nocturn completat per a {context['ds']}")

    fi = PythonOperator(
        task_id="fin",
        python_callable=_avisar_exit,
        trigger_rule="all_success",
    )

    # -----------------------------------------------------------------------
    # DEPENDENCIES: l operador >> significa "i despres".
    # Aquesta es la declaracio completa del graf, en quatre linies.
    # -----------------------------------------------------------------------
    (
        inici
        >> exportar_comandes
        >> esperar_fitxer
        >> carregar_comandes
        >> consolidar_comandes
        >> llancar_dataflow
        >> grup_agregacions
        >> refrescar_vista
        >> comprovar_qualitat
        >> fi
    )

La tasca 8, la comprovació de qualitat, és la que converteix aquest DAG en una cosa seriosa. Verifica tres regles de negoci: que hi va haver comandes, que cap venda no és negativa, i que cap categoria no supera 500 comandes en un dia (impossible amb el volum d'AlpinaShop, per tant indicaria duplicació). Si alguna cosa no quadra, el DAG falla sorollosament abans que ningú miri l'informe. És la diferència entre detectar un error a les 05:00 i descobrir-lo en una reunió.

  1. Operadors de Google Cloud

Airflow porta un catàleg enorme d'operadors per a Google Cloud. Els més útils per a AlpinaShop:

Operador Què fa
BigQueryInsertJobOperator Executa qualsevol consulta o feina de BigQuery. El més versàtil
BigQueryCheckOperator Executa un SQL que ha de retornar cert; si no, falla
BigQueryValueCheckOperator Compara un resultat amb un valor esperat i una tolerància
BigQueryTableExistenceSensor Espera que existeixi una taula
GCSToBigQueryOperator Carrega fitxers del bucket a una taula
BigQueryToGCSOperator Exporta una taula a fitxers
GCSObjectExistenceSensor Espera que aparegui un objecte
GCSToGCSOperator Copia o mou entre buckets
DataflowStartFlexTemplateOperator Llança una plantilla flexible de Dataflow
DataprocCreateBatchOperator Llança una feina a Dataproc Serverless
DataprocCreateClusterOperator / DeleteCluster Clúster efímer des del DAG
CloudSQLExportInstanceOperator Exporta una instància de Cloud SQL
PubSubPublishMessageOperator Publica en un topic
CloudRunExecuteJobOperator Executa un job de Cloud Run

Per al procés mensual d'AlpinaShop, l'anàlisi de cistella de 04-03 es llança així:

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

analisi_cistella = DataprocCreateBatchOperator(
    task_id="analisis_cesta_mensual",
    region=REGIO,
    project_id=PROJECTE_DADES,
    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 data relativa a la d'execució. És el que manté el DAG reproduïble: en reprocessar el gener, el rang serà el del gener, no el d'avui.

  1. XComs, variables i connexions

XCom (cross-communication) permet que una tasca passi un valor petit a una altra:

def _comptar_files(**context):
    from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook
    hook = BigQueryHook(location=REGIO, use_legacy_sql=False)
    files = hook.get_first(
        f"""SELECT COUNT(*) FROM `{PROJECTE_DADES}.{DATASET}.pedidos_staging`"""
    )[0]
    # El que retorna la funcio es desa automaticament com a XCom
    return int(files)


def _validar_volum(**context):
    files = context["ti"].xcom_pull(task_ids="contar_filas_cargadas")
    if files == 0:
        raise ValueError(f"Cap comanda carregada el {context['ds']}")
    if files > 5000:
        raise ValueError(f"{files} comandes: volum anomal, revisar duplicats")
    print(f"Volum correcte: {files} comandes")


comptar_files = PythonOperator(task_id="contar_filas_cargadas",
                               python_callable=_comptar_files)
validar_volum = PythonOperator(task_id="validar_volumen",
                               python_callable=_validar_volum)

Els XCom són per a valors petits: identificadors, comptadors, rutes. Es desen a la base de dades de metadades d'Airflow. Passar un DataFrame per XCom és un error clàssic que infla la base de dades i degrada tot l'entorn. Per a dades grans, s'escriu a Cloud Storage i es passa la ruta.

Variables són configuració global, editable des de la interfície sense tocar codi:

from airflow.models import Variable

llindar = 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

Connexions guarden credencials i paràmetres d'accés a sistemes externs, xifrades. Per al que sigui sensible de veritat, el correcte és que la connexió apunti a Secret Manager (03-06) mitjançant el backend de secrets d'Airflow, en comptes de guardar la contrasenya a la base de dades de metadades.

  1. TaskGroup, dependències i patrons de flux

Els operadors de dependència:

a >> b                    # b despres d a
a << b                    # a despres de b
a >> [b, c] >> d          # b i c en paral·lel, tots dos despres d a; d despres d ambdos
[a, b] >> c               # c espera que acabin a i b

Les regles de dispar (trigger_rule) controlen sota quina condició s'executa una tasca:

Regla S'executa quan
all_success (per defecte) Totes les tasques anteriors han tingut èxit
all_failed Totes han fallat
all_done Totes han acabat, amb èxit o sense
one_success Almenys una ha tingut èxit
one_failed Almenys una ha fallat
none_failed_min_one_success Cap no ha fallat i almenys una s'ha executat

La regla all_done és imprescindible per a tasques de neteja:

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

esborrar_cluster = DataprocDeleteClusterOperator(
    task_id="borrar_cluster",
    cluster_name="cluster-efimero-{{ ds_nodash }}",
    region=REGIO,
    project_id=PROJECTE_DADES,
    # CRITIC: s esborra el cluster PASSI EL QUE PASSI.
    # Amb all_success, una fallada del job deixaria el cluster ences
    # i facturant indefinidament.
    trigger_rule="all_done",
)

I one_failed per a notificacions d'error:

avisar_fallada = PythonOperator(
    task_id="avisar_fallo",
    python_callable=_notificar_a_slack,
    trigger_rule="one_failed",
)
[carregar_comandes, llancar_dataflow, refrescar_vista] >> avisar_fallada

Els TaskGroup agrupen visualment tasques relacionades, col·lapsant-les en un sol node a la interfície. Amb un DAG de trenta tasques, la diferència entre un graf llegible i un embolic.

  1. Idempotència, catchup i reprocessos

La data lògica. Airflow executa cada DAG per a una data lògica, disponible com a {{ ds }}. És el que fa que el DAG sigui reproduïble: en reprocessar el 12 de març, totes les plantilles valen 2026-03-12.

La regla d'or: una tasca mai no ha d'utilitzar CURRENT_DATE() ni datetime.now(). Ha d'utilitzar {{ ds }}. Amb CURRENT_DATE(), reprocessar el 12 de març recalcularia amb la data d'avui i produiria un resultat incorrecte en silenci.

# MALAMENT: no reproduible
"query": "SELECT ... WHERE fecha = CURRENT_DATE()"

# BE: reproduible
"query": "SELECT ... WHERE fecha = DATE '{{ ds }}'"

catchup. Si es posa catchup=True i el start_date és de fa tres mesos, Airflow executarà totes les dates pendents en activar el DAG. Pot ser útil per emplenar un històric, o pot llançar noranta execucions simultànies i esgotar la quota de BigQuery. Per a AlpinaShop: catchup=False, i els emplenaments es fan a mà i amb control.

backfill. Reprocessar un rang explícitament:

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

Aquesta comanda només és segura si totes les tasques són idempotents. Per això el DAG utilitza WRITE_TRUNCATE al staging, MERGE a la consolidació i CREATE OR REPLACE TABLE a les agregacions. Un DAG amb INSERT en comptes de MERGE duplicaria dades a cada reprocés, i el reprocés —que hauria de ser l'eina de reparació— es convertiria en la causa d'un problema pitjor.

Reexecutar una sola tasca, sense tot 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: orquestració serverless en YAML

Si Composer costa 350 € al mes i AlpinaShop té vuit tasques, hi ha una alternativa: Cloud Workflows, orquestració declarativa sense servidor i sense cost fix. Es paga per pas executat, i són cèntims.

Un flux es defineix en YAML i s'executa quan se'l invoca:

# proceso-nocturno.yaml -- Proces nocturn d AlpinaShop amb 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 es passa data, s utilitza ahir. Permet reprocessar
          # invocant el flux amb {"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. L API retorna una OPERACIO de llarga
    #    durada: cal esperar-la, no donar-la per feta.
    # -----------------------------------------------------------------
    - 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) Comprovar que el fitxer existeix. L equivalent 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) Carregar a 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) Llancar Dataflow amb 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) Agregacions EN PARAL·LEL amb la branca '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) Comprovacio de qualitat: si falla, es llanca un error explicit
    # -----------------------------------------------------------------
    - 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: ${"QUALITAT: no hi ha vendes agregades per a " + fecha}

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

# ---------------------------------------------------------------------
# SUBFLUXOS reutilitzables
# ---------------------------------------------------------------------
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}

Desplegament i execució:

gcloud iam service-accounts create sa-workflows-nocturno \
  --display-name="Orquestracio nocturna amb 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

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

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

Els límits honestos de Workflows, que cal conèixer abans de comprometre-s'hi:

Límit Valor aproximat Implicació
Durada màxima d'una execució 1 any Sense problema
Durada màxima d'una crida HTTP 30 min Una feina llarga cal sondejar-la, no esperar-la
Mida màxima d'una variable 512 KB No passar dades, només referències
Passos per execució ~100.000 Suficient
Interfície visual Graf simple Molt inferior al tauler d'Airflow
backfill d'un rang No existeix Cal invocar en bucle des d'un script
Sensors d'esdeveniments Sondeig manual amb retry Menys elegant que un sensor d'Airflow
Ecosistema d'operadors Crides a API Sense els centenars d'operadors d'Airflow

Les tres files que més pesen són el backfill inexistent, l'observabilitat més pobra i l'absència d'operadors llestos. Amb vuit tasques és assumible; amb vuitanta, dolorós.

  1. Cloud Scheduler: el disparador per cron

Workflows no té planificador propi: cal disparar-lo. Això ho fa Cloud Scheduler, un cron gestionat que costa pràcticament res (les primeres feines són gratuïtes i després són cèntims).

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

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" és important: la feina s'executa a les 02:00 hora local, i Google gestiona els canvis d'horari d'estiu. Amb UTC, el procés es desplaçaria una hora dues vegades l'any, cosa que sona inofensiva fins que coincideix amb la finestra de manteniment d'un altre sistema.

Cloud Scheduler també pot publicar a Pub/Sub o invocar Cloud Run directament, cosa que el converteix en el disparador universal de la plataforma.

I amb la notificació del bucket de 04-04, es pot muntar un flux reactiu en comptes de programat:

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

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

El flux es dispara a les 02:00 o quan arriba un fitxer, el que passi. És més robust que una hora fixa, perquè no depèn que el proveïdor sigui puntual.

  1. La taula de decisió i l'elecció d'AlpinaShop

Criteri Cloud Scheduler Workflows Cloud Composer Plantilles de Dataflow
Què resol "Executa això a aquesta hora" "Executa aquests passos en aquest ordre" Orquestració completa Un pipeline concret
Definició Cron + destinació YAML declaratiu Python Paràmetres
Cost fix ~0 € 0 € ~350 €/mes 0 €
Cost variable Cèntims Cèntims per pas Treballadors Recursos del job
Dependències complexes No Sí, amb límits Sí, sense límits No
Paral·lelisme No Sí (parallel) Intern
Reintents Sí, configurable Sí, molt fi
Sensors d'esdeveniments No Sondeig manual Sí, natius No
backfill No Manual Sí, natiu No
Observabilitat Registres Graf simple Tauler complet Interfície de Dataflow
Ecosistema N/A API de Google Centenars d'operadors N/A
Corba d'aprenentatge Minuts Hores Dies Minuts
Triar si… Una acció periòdica Poques desenes de passos Desenes de DAG i equips Un pipeline solt

La decisió d'AlpinaShop: començar amb Workflows + Cloud Scheduler.

El raonament, que és el que cal saber defensar:

  1. El cost mana. 350 € al mes per orquestrar vuit tasques és desproporcionat quan tota la plataforma de dades costa 15 €. Amb Workflows i Scheduler, l'orquestració costa cèntims.
  2. La complexitat actual no ho justifica. Vuit tasques amb dependències lineals i una bifurcació en paral·lel caben perfectament en un YAML de 150 línies. Airflow brilla amb cinquanta DAG que comparteixen dependències entre equips, i AlpinaShop en té un.
  3. No hi ha pèrdua funcional rellevant. Es cobreixen dependències, paral·lelisme, reintents amb retrocés, comprovació de qualitat i alertes. El que falta —backfill natiu, sensors, tauler ric— se supleix amb un script d'invocació en bucle i les alertes de Cloud Monitoring.
  4. La migració posterior és viable. Si demà calen vint DAG, es crea l'entorn de Composer i es migra. La lògica de les tasques —consultes SQL, plantilles de Dataflow, jobs de Dataproc— no canvia: només canvia qui les invoca. Res del que s'ha fet no es llença.

Els criteris objectius per fer el salt a Composer, escrits per endavant per no discutir-ho amb l'emoció del moment:

  • Més de 10 fluxos diferents amb dependències entre ells.
  • Necessitat recurrent de backfill de rangs llargs.
  • Més d'una persona mantenint els fluxos, amb revisió per pull request.
  • Necessitat d'operadors especialitzats (Salesforce, SAP, Kubernetes).
  • Que el cost de l'orquestrador baixi del 10 % del cost de la plataforma que orquestra.

Aquest últim criteri és el més útil i el més fàcil de comprovar: quan la plataforma de dades d'AlpinaShop costi 3.500 € al mes, Composer costarà el 10 % i estarà justificat. Avui costaria el 2.300 %.

Errors Habituals i Consells

Utilitzar CURRENT_DATE() en comptes de {{ ds }}. Trenca la reproductibilitat en silenci: el reprocés d'una data passada calcula amb la d'avui i ningú no ho nota.

Deixar catchup=True sense voler. En activar un DAG amb start_date antic, es llancen centenars d'execucions alhora, s'esgoten les quotes i es dispara el cost.

Sensors en mode poke amb esperes llargues. Ocupen treballadors. Amb diversos sensors simultanis, l'entorn es bloqueja a si mateix. Utilitza mode="reschedule".

Tasques no idempotents. Un INSERT en comptes d'un MERGE converteix cada reintent i cada reprocés en una duplicació de dades. La idempotència no és opcional en un DAG.

Passar dades grans per XCom. Es desen a la base de dades de metadades i la degraden. Passa rutes, no continguts.

Oblidar trigger_rule="all_done" a les tasques de neteja. Un clúster efímer que no s'esborra perquè el job va fallar queda encès facturant indefinidament.

No posar execution_timeout. Una tasca penjada reté un treballador per sempre i acaba bloquejant el DAG sencer.

Posar lògica pesada al cos del DAG. El fitxer es reavalua cada pocs segons pel planificador. Una consulta a una base de dades fora d'un operador s'executa constantment i enfonsa l'entorn. Tota la feina va dins d'operadors.

Deixar Composer encès "per si de cas". És l'equivalent al clúster de Dataproc permanent de 04-03 i al Data Fusion permanent de 04-05: el mateix error, tres vegades, i sempre la partida més cara de la factura.

Consell: guarda els DAG a Git i desplega amb CI. El bucket de DAG no és el lloc on viu el codi, és on es copia. A 06-01 ho automatitzarem amb Cloud Build.

Consell: posa una comprovació de qualitat al final de tot flux. És la tasca que converteix un procés automàtic en un de fiable. Detectar l'error a les 05:00 val molt més que descobrir-lo en un comitè.

Consell: anomena les tasques amb verbs. exportar_pedidos_cloudsql, no tarea_1. Quan alguna cosa falli a les 03:00, el nom és el primer que es llegeix.

Exercicis

Exercici 1: DAG mensual d'anàlisi de cistella

Escriu un DAG d'Airflow anomenat alpinashop_analisis_mensual que s'executi el dia 1 de cada mes a les 04:00 hora de Madrid i faci: (1) comprovar que la taula lineas_pedido té dades del mes anterior, fallant si no; (2) llançar la feina de Dataproc Serverless cesta_media.py amb les dates del mes anterior com a arguments; (3) en paral·lel, executar dues consultes de BigQuery que calculin el marge per categoria i els productes sense vendes del mes; (4) refrescar la vista materialitzada; (5) publicar un missatge en un topic informes-listos indicant que l'informe mensual està disponible. Inclou reintents, SLA i notificació de fallada.

Exercici 2: el mateix procés en Workflows

Implementa en YAML un flux alpinashop-informe-mensual equivalent als passos 1, 2 i 3 de l'exercici anterior, que accepti un paràmetre mes amb format yyyy-MM i utilitzi el mes anterior si no se'n passa cap. Ha d'executar les dues consultes en paral·lel, esperar correctament que acabi la feina de Dataproc (que és una operació de llarga durada) i llançar un error explícit si la comprovació inicial no troba dades. Afegeix la feina de Cloud Scheduler que el dispari.

Exercici 3: diagnòstic d'un DAG que menteix

Durant tres setmanes, el DAG alpinashop_proceso_nocturno apareix en verd totes les nits. Però el dilluns, la Lucía detecta que la taula agg_ventas_categoria_dia té les dades del 24 de febrer repetides a totes les particions des d'aquella data. Investigant trobes: la tasca consolidar_pedidos utilitza INSERT INTO en comptes de MERGE; la consulta d'agregació filtra per WHERE l.fecha_pedido = CURRENT_DATE() - 1; la tasca comprobar_calidad_datos verifica únicament que la taula no estigui buida; i el 24 de febrer algú va executar un backfill de les dues setmanes anteriors.

Explica exactament què ha passat i en quin ordre, per què el DAG apareixia en verd, i proposa les correccions concretes —amb el codi— per a cadascun dels quatre problemes.

Solucions

Solució 1

"""alpinashop_analisis_mensual.py -- Analisi de cistella i marge 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

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

arguments = {
    "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=arguments,
    schedule="0 4 1 * *",                     # dia 1 de cada mes a les 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 a l ultim dia
    PRIMER_DIA_MES_ANT = "{{ macros.ds_format(macros.ds_add(ds, -1), '%Y-%m-%d', '%Y-%m-01') }}"
    ULTIM_DIA_MES_ANT = "{{ macros.ds_add(ds, -1) }}"

    # 1) Comprovar que hi ha dades del mes anterior
    comprovar_dades = BigQueryCheckOperator(
        task_id="comprobar_datos_mes_anterior",
        location=REGIO,
        use_legacy_sql=False,
        sql=f"""
            SELECT COUNT(*) > 0
            FROM `{PROJECTE}.{DATASET}.lineas_pedido`
            WHERE fecha_pedido BETWEEN DATE '{PRIMER_DIA_MES_ANT}'
                                   AND DATE '{ULTIM_DIA_MES_ANT}'
        """,
    )

    # 2) Analisi de cistella a Dataproc Serverless
    analisi_cistella = DataprocCreateBatchOperator(
        task_id="analisis_cesta",
        region=REGIO,
        project_id=PROJECTE,
        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={ULTIM_DIA_MES_ANT}",
                ],
            },
            "runtime_config": {"version": "2.2"},
            "environment_config": {
                "execution_config": {
                    "service_account":
                        "[email protected]",
                    "subnetwork_uri": "sn-datos-euw1",
                }
            },
        },
    )

    # 3) Dues consultes en paral·lel
    with TaskGroup(group_id="informes_mensuales") as informes:

        marge_categoria = BigQueryInsertJobOperator(
            task_id="margen_por_categoria",
            location=REGIO,
            configuration={"query": {"useLegacySql": False, "query": f"""
                CREATE OR REPLACE TABLE `{PROJECTE}.{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 `{PROJECTE}.{DATASET}.lineas_pedido`               AS l
                JOIN `{PROJECTE}.{DATASET}.productos`                   AS pr USING (sku)
                JOIN `{PROJECTE}.{DATASET}.costes_producto_erp`         AS c  USING (sku)
                WHERE l.fecha_pedido BETWEEN DATE '{PRIMER_DIA_MES_ANT}'
                                         AND DATE '{ULTIM_DIA_MES_ANT}'
                GROUP BY mes, pr.categoria
            """}},
        )

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

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

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

    comprovar_dades >> analisi_cistella >> informes >> refrescar >> avisar

El punt que s'avalua és el càlcul del mes anterior amb macros.ds_add i macros.ds_format en comptes d'amb datetime.now(). Executant el DAG l'1 d'abril, ds val 2026-04-01, ds_add(ds, -1) dóna 2026-03-31 i el format a %Y-%m-01 dóna 2026-03-01. En reprocessar l'1 de febrer, els valors seran els del gener. Amb dates calculades en temps real, el reprocés seria inútil.

Solució 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) Comprovar que hi ha dades
    - 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: ${"Sense linies de comanda per al mes " + mes}

    # 2) Dataproc Serverless: operacio de llarga durada
    - 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) Dues consultes en paral·lel
    - 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}

# Sondeig de l operacio de llarga durada: Workflows NO pot
# esperar 20 minuts en una sola crida HTTP (limit de 30 min,
# i a mes el batch podria trigar mes).
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 fallit: " + default(map.get(estado.body, "stateMessage"), "sense detall")}
          - condition: ${estado.body.state == "CANCELLED"}
            raise: "Batch de Dataproc cancel·lat"
    - 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

El que s'avalua: el subflux espera_batch amb sondeig i sys.sleep. És la diferència clau amb Airflow, on DataprocCreateBatchOperator espera sol. A Workflows cal implementar el bucle a mà, i cal fer-ho bé: comprovar els tres estats finals (SUCCEEDED, FAILED, CANCELLED), no només el d'èxit, perquè si no un batch fallit deixaria el flux sondejant eternament.

Solució 3

Què ha passat, en ordre cronològic:

Pas 1 — El backfill del 24 de febrer va ser la causa desencadenant. Es van reprocessar dues setmanes. Per a cadascuna de les 14 dates lògiques, el DAG va executar consolidar_pedidos, que utilitza INSERT INTO. Com que les taules ja contenien aquelles comandes de les execucions originals, cada comanda d'aquelles dues setmanes va quedar duplicada a pedidos.

Pas 2 — L'agregació va amplificar el problema i el va congelar. La consulta filtra per WHERE l.fecha_pedido = CURRENT_DATE() - 1. Durant el backfill, CURRENT_DATE() era el 24 de febrer a les catorze execucions, perquè CURRENT_DATE() no sap res de la data lògica. És a dir: les catorze execucions van calcular el mateix dia, el 23 de febrer, i van escriure catorze vegades el mateix resultat sobre la mateixa partició.

Pas 3 — I va continuar malament cada nit. Aquí hi ha el detall que explica les tres setmanes: com que la consulta utilitza CREATE OR REPLACE TABLE sense filtre per data lògica a l'escriptura, cada execució nocturna posterior va sobreescriure la taula amb el que s'havia calculat per a "ahir" segons CURRENT_DATE(). Això hauria d'haver-se anat actualitzant... llevat que la taula resultant conserva l'històric previ tal com va quedar després del backfill, i les particions antigues van quedar congelades amb la dada del 24 de febrer. D'aquí que la Lucía vegi el 24 de febrer repetit a totes les particions des de llavors.

Per què el DAG apareixia en verd. Perquè cap tasca no va fallar. INSERT INTO no dóna error en inserir duplicats: no hi ha clau primària a BigQuery. La consulta d'agregació es va executar correctament. I la comprovació de qualitat només verificava que la taula no estigués buida —i no ho estava: estava plena de dades incorrectes. Verd no significa correcte; significa que no hi va haver excepcions. Aquesta distinció és la lliçó de l'exercici.

Correccions, una per problema:

A) INSERT INTOMERGE. Fa la consolidació idempotent:

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 }}, i escriptura només de la partició corresponent. Amb CREATE OR REPLACE TABLE es reescriu la taula sencera; el correcte és escriure únicament la partició de la data lògica:

agregar_vendes = BigQueryInsertJobOperator(
    task_id="ventas_por_categoria",
    location=REGIO,
    configuration={
        "query": {
            "useLegacySql": False,
            # Escriu NOMES la particio de la data logica
            "destinationTable": {
                "projectId": PROJECTE,
                "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 `{PROJECTE}.{DATASET}.lineas_pedido` AS l
                JOIN `{PROJECTE}.{DATASET}.productos`     AS pr USING (sku)
                WHERE l.fecha_pedido = DATE '{{{{ ds }}}}'
                GROUP BY dia, pr.categoria
            """,
        }
    },
)

El decorador $ a agg_ventas_categoria_dia${{ ds_nodash }} és el decorador de partició de BigQuery: WRITE_TRUNCATE afecta només aquesta partició, no la taula. Així, reprocessar el 12 de març corregeix el 12 de març i no toca cap altre dia.

C) Comprovació de qualitat de veritat. No "hi ha files", sinó regles de negoci:

comprovar_qualitat = BigQueryCheckOperator(
    task_id="comprobar_calidad_datos",
    location=REGIO,
    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 `{PROJECTE}.{DATASET}.agg_ventas_categoria_dia`
          WHERE dia = DATE '{{{{ ds }}}}'
        ),
        control_duplicados AS (
          SELECT COUNT(*) AS duplicados FROM (
            SELECT pedido_id
            FROM `{PROJECTE}.{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ó duplicados = 0 és la que hauria detectat el problema la mateixa nit del 24 de febrer, en comptes de tres setmanes després.

D) Procediment de reparació. Corregir el codi no arregla les dades ja corrompudes:

-- 1) Deduplicar la taula de comandes conservant l ultima versio
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) Reprocessar les agregacions del periode afectat, ja amb el DAG corregit
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

Aquest backfill ara sí que és segur, perquè després de les correccions A i B totes les tasques són idempotents. Abans, era la causa del problema. La moralitat completa de l'exercici: el backfill no és perillós; el que és perillós és un DAG no idempotent, i el backfill només ho revela.

Conclusió

Has vist per què cron no és un orquestrador: respon a "quina hora és?" quan la pregunta correcta és "què es pot executar ara, atès el que ha passat?". I has vist el preu concret d'aquesta confusió: una exportació que triga vint minuts de més i un comitè de direcció mirant vendes falses.

Coneixes les cinc aportacions d'un orquestrador real —dependències explícites, gestió de fallades, observabilitat, reproductibilitat i notificació— i les tres eines de Google Cloud que les cobreixen en diferent grau.

Has après Cloud Composer, Airflow gestionat, amb els seus conceptes: DAG com a graf acíclic, tasques, operadors que defineixen què fa cadascuna, i sensors que esperen que passi alguna cosa en comptes d'apostar per una hora. Has escrit el DAG complet del procés nocturn d'AlpinaShop: exportació des de la rèplica de lectura, sensor en mode reschedule que espera el fitxer de veritat, càrrega amb WRITE_TRUNCATE a staging, MERGE idempotent cap a la taula final, llançament de la plantilla flexible de Dataflow esperant que acabi, dues agregacions en paral·lel dins d'un TaskGroup, refresc de la vista materialitzada i una comprovació de qualitat que fa fallar el DAG abans que ningú miri un informe erroni. Amb reintents exponencials, execution_timeout, SLA i notificació.

Coneixes els operadors de Google Cloud, els XCom per a valors petits —mai per a dades—, les variables i connexions amb Secret Manager al darrere, les regles de dispar amb all_done per a les neteges que s'han d'executar passi el que passi, i la regla d'or de la reproductibilitat: {{ ds }} sempre, CURRENT_DATE() mai, perquè d'ella depèn que un backfill repari en comptes d'espatllar.

I has posat el cost sobre la taula sense adorns: uns 350 € al mes per un entorn petit, davant dels 15 € que costa tota la plataforma de dades d'AlpinaShop. L'orquestrador costaria vint vegades més que allò orquestrat. Per això has muntat l'alternativa: Workflows, orquestració declarativa en YAML sense cost fix, amb branques paral·leles, reintents amb retrocés, subfluxos reutilitzables i sondeig explícit de les operacions llargues; disparada per Cloud Scheduler amb la zona horària de Madrid perquè els canvis d'horari no desplacin el procés. Coneixes els seus límits reals —sense backfill natiu, sense sensors, sense tauler ric, sense catàleg d'operadors— i saps que amb vuit tasques són assumibles i amb vuitanta no.

La decisió d'AlpinaShop queda fixada i argumentada: començar amb Workflows i Cloud Scheduler, amb criteris objectius escrits per endavant per saltar a Composer —més de deu fluxos interdependents, necessitat recurrent de backfill, diverses persones mantenint-los, o el criteri més net de tots: quan l'orquestrador baixi del 10 % del cost del que orquestra—. I sabent que la migració no llençarà res, perquè la lògica viu a les consultes, les plantilles i els jobs; només canvia qui els invoca.

Amb això, la maquinària està completa. Les dades entren soles, es transformen soles, s'agreguen soles i es comproven soles, totes les nits, amb alertes si alguna cosa falla.

I apareix l'últim problema del mòdul, que no és tècnic. Ara hi ha setze taules a alpinashop_analitica, mitja dotzena de vistes, un llac amb Parquet, taules de staging, agregats i una taula anomenada pedidos_evento que ningú no recorda per què existeix. Quan algú de màrqueting pregunta "on és la dada de conversió?", ningú no sap contestar sense obrir BigQuery i buscar. Ningú no ha escrit què significa exactament ventas_eur —si inclou IVA, si descompta devolucions—. Ningú no sap si email_cliente continua apareixent en algun lloc on no hauria de ser. I el tauler de control que la Lucía va prometre a direcció continua sense existir: les dades són perfectes i ningú no les veu.

A 04-07, Dataplex i Looker Studio, tancarem el mòdul amb això. Veuràs què és el govern de les dades i per què una pime també el necessita: catàleg, llinatge, qualitat, classificació i cicle de vida. Muntaràs llacs, zones i actius a Dataplex, documentaràs alpinashop_analitica amb etiquetes, definiràs regles de qualitat que es comproven soles —total_pedido mai negatiu, sku sempre present—, i utilitzaràs Sensitive Data Protection per descobrir on hi ha correus i telèfons i desidentificar-los abans d'exposar-los, amb l'advertiment de RGPD que correspon. I per fi construiràs a Looker Studio el tauler de control de direcció d'AlpinaShop, amb les seves bones pràctiques de disseny, de cost i —sobretot— de permisos, incloent-hi l'error clàssic de "credencials del propietari" que converteix un informe compartit en una fuga de dades.

Curs de Google Cloud Platform (GCP)

Mòdul 1: Introducció a Google Cloud Platform

Mòdul 2: Serveis principals de GCP

Mòdul 3: Xarxes i seguretat

Mòdul 4: Dades i anàlisi

Mòdul 5: Aprenentatge automàtic i IA

Mòdul 6: DevOps i monitoratge

Mòdul 7: Temes avançats de GCP

Mòdul 8: Projecte final

© Copyright 2026. Tots els drets reservats