alpinashop_analitica ja existeix i té a dins l'històric de comandes. Però aquest històric hi va arribar d'una manera que no es pot repetir cada nit: la Lucía va executar a mà una consulta federada contra la rèplica de PostgreSQL. Va funcionar una vegada. La pregunta és què passa demà, i demà passat, i el dia que algú demani veure les vendes de la campanya avui i no demà.

La resposta ingènua és escriure un script. Un fitxer sincronizar.py que es connecti a Cloud SQL, llegeixi les comandes del dia, les transformi i les insereixi a BigQuery, llançat per un cron en una màquina virtual. Aquest script funcionarà durant setmanes, i després fallarà. Fallarà perquè les dades creixen i l'script triga cinc hores; perquè la VM es reinicia a mig procés i ningú no sap si va escriure la meitat de les files; perquè una comanda amb un caràcter estrany llança una excepció i es perd la resta del bolcat; perquè ningú no sap si es va executar ahir; i perquè el dia que algú demani dades en temps real, tot el disseny s'ha de llençar.

Cloud Dataflow és el servei gestionat de Google per executar pipelines de dades escrits amb Apache Beam. La seva promesa concreta és doble: escrius la lògica de transformació una sola vegada i Google s'encarrega del paral·lelisme, del reintent, de l'escalat i de la coherència; i aquest mateix codi serveix tant per processar dos anys d'històric com per processar les comandes a mesura que entren.

En aquesta lliçó entendràs el model de Beam, escriuràs un pipeline per lots que carrega l'històric d'AlpinaShop, el provaràs en local, el llançaràs a Dataflow, i després entraràs a la part que de debò separa el processament de dades seriós de l'amateur: el temps. Què vol dir "les vendes de les 10:00" quan un esdeveniment generat a les 10:00 arriba a les 10:07, i com es respon a això sense mentir.

Contingut

  1. Què és un pipeline de dades i per què un script no n'hi ha prou
  2. Apache Beam: el model unificat
  3. Els quatre conceptes: Pipeline, PCollection, PTransform, runner
  4. Les transformacions que faràs servir el 90 % del temps
  5. Primer pipeline per lots, explicat línia a línia
  6. Execució local amb DirectRunner
  7. Execució gestionada amb DataflowRunner
  8. El temps en streaming: esdeveniment davant de procés
  9. Marques d'aigua, finestres, activadors i dades tardanes
  10. El pipeline en temps real de pedidos-nuevos
  11. Plantilles de Dataflow: la via pràctica
  12. Escalat automàtic i Dataflow Prime
  13. Monitoratge, paral·lelisme i biaix de dades
  14. Cost: què es paga exactament
  15. Quan Dataflow no és la resposta

  1. Què és un pipeline de dades i per què un script no n'hi ha prou

Un pipeline de dades és una seqüència declarada d'operacions que porten dades d'un origen a un destí transformant-les pel camí. L'adjectiu important és declarada: descrius què vols que passi, no com es reparteix la feina entre màquines.

Comparem amb honestedat l'script en una VM i el pipeline gestionat:

Aspecte Script en una VM Pipeline a Dataflow
Paral·lelisme El que programis tu, a mà, amb fils Automàtic: reparteix per treballadors
Escalat Canviar la VM i reiniciar Horitzontal i automàtic durant l'execució
Fallada d'una màquina Es perd tot el procés Es reintenta el fragment afectat
Fallada d'un registre Excepció que tomba el procés Es desvia a una sortida d'errors i continua
Semàntica d'escriptura La que aconsegueixis Exactament una vegada a les connexions natives
Estat si es reinicia Desconegut Gestionat pel servei
Lot i streaming Dos programes diferents El mateix codi
Cost en repòs La VM encesa sempre Zero: no hi ha res encès
Observabilitat Els print que hi hagis posat Graf, mètriques i registres integrats

La fila que més fa mal a la pràctica és la de la fallada d'un registre. Un bolcat de 800.000 comandes en què la comanda 412.337 té un import amb coma en comptes de punt no ha de perdre les 387.663 comandes restants. Un script mal fet les perd; un de ben fet requereix un esforç considerable de tractament d'errors que a Beam ve de sèrie amb les sortides etiquetades.

I la fila del cost en repòs té el seu matís: Dataflow en mode lot costa zero quan no s'executa, però un pipeline en temps real està permanentment encès i factura sense parar. Hi tornarem a l'apartat 14, perquè és la sorpresa més habitual.

  1. Apache Beam: el model unificat

Apache Beam és un model de programació open source —donat per Google a l'Apache Software Foundation— per definir pipelines de dades que després s'executen en motors diferents.

La paraula clau és unificat. Abans de Beam, processar per lots i processar en temps real eren mons separats amb eines separades, i les empreses mantenien dues implementacions de la mateixa lògica de negoci: una per a l'històric i una altra per al temps real, que inevitablement divergien i donaven números diferents. Beam parteix de la idea que un lot és simplement un flux acotat i un stream és un flux no acotat, i que la lògica de transformació és la mateixa en tots dos casos.

flowchart LR
    subgraph SDK["SDK d Apache Beam"]
        P["Codi del pipeline<br/>Python / Java / Go"]
    end
    subgraph Runners["Runners"]
        D["DirectRunner<br/>local, proves"]
        DF["DataflowRunner<br/>Google Cloud"]
        FL["FlinkRunner"]
        SP["SparkRunner"]
    end
    P --> D
    P --> DF
    P --> FL
    P --> SP

El runner és el motor que executa el pipeline. El mateix fitxer Python pot córrer al teu portàtil amb DirectRunner, a Dataflow amb DataflowRunner, o en un clúster de Flink o Spark. Això redueix el bloqueig amb el proveïdor: si AlpinaShop hagués de sortir de Google Cloud, la lògica dels seus pipelines viatjaria.

Beam té SDK per a Java, Python i Go. Farem servir Python, coherent amb la resta de l'aplicació d'AlpinaShop, que ja és Flask.

# Instal·lacio del SDK amb les dependencies de Google Cloud
python -m venv venv-beam
source venv-beam/bin/activate
pip install 'apache-beam[gcp]==2.64.0'

Fixar la versió és deliberat: els treballadors de Dataflow faran servir exactament la versió del SDK amb què llances el pipeline, i les diferències entre versions són una font clàssica de fallades que només apareixen al núvol.

  1. Els quatre conceptes: Pipeline, PCollection, PTransform, runner

Pipeline és l'objecte que conté el graf complet. Es construeix, es declara i s'executa. No passa res mentre l'escrius: estàs dibuixant un plànol.

PCollection és un conjunt de dades distribuït i immutable. No és una llista de Python: pot tenir zero elements o bilions, pot estar repartida per cent màquines, i no es pot modificar. Cada transformació produeix una PCollection nova. La immutabilitat és el que permet reintentar un fragment fallit sense corrompre res.

Una PCollection pot ser:

  • Acotada (bounded): té un final conegut. Un fitxer, una taula. És el cas de lot.
  • No acotada (unbounded): no acaba mai. Un topic de Pub/Sub. És el cas de temps real.

PTransform és una operació que pren una o més PCollection i produeix una o més PCollection. S'aplica amb l'operador |, que a Beam està sobrecarregat per significar "aplica aquesta transformació".

Runner és el motor d'execució, ja vist.

La sintaxi, que resulta estranya la primera vegada:

resultat = entrada | "Nom descriptiu del pas" >> beam.Map(funcio)

L'operador >> associa un nom al pas. No és decoratiu: aquest nom és el que apareix al graf de la consola de Dataflow i a les mètriques. Un pipeline amb passos anomenats Map(<lambda at main.py:34>) és impossible de depurar en producció. Anomena tots els passos, sempre.

flowchart TD
    A["PCollection: linies de text crues<br/>gs://alpinashop-catalogo/exportaciones/..."]
    B["PTransform: AnalitzarCSV<br/>ParDo"]
    C["PCollection: diccionaris de comanda"]
    D["PTransform: ValidarINetejar<br/>ParDo amb sortides multiples"]
    E["PCollection: comandes valides"]
    F["PCollection: comandes rebutjades"]
    G["PTransform: WriteToBigQuery"]
    H["PTransform: WriteToText<br/>quarantena al bucket"]

    A --> B --> C --> D
    D --> E --> G
    D --> F --> H

Aquest graf és exactament el que escriuràs a l'apartat 5. Fixa't que la branca d'errors és part del disseny, no un afegit.

  1. Les transformacions que faràs servir el 90 % del temps

Transformació Què fa Exemple a AlpinaShop
beam.Map(f) Aplica f a cada element; retorna un Convertir una línia CSV en diccionari
beam.FlatMap(f) Aplica f; retorna 0, 1 o N elements Desglossar una comanda en les seves línies
beam.Filter(f) Es queda amb els que compleixen f Descartar comandes cancel·lades
beam.ParDo(DoFn) La forma general: classe amb estat, cicle de vida i sortides múltiples Validar i separar vàlids de rebutjats
beam.GroupByKey() Agrupa parells (clau, valor) per clau Agrupar línies per sku
beam.CombinePerKey(f) Agrupa i redueix amb una funció associativa Sumar vendes per sku
beam.CoGroupByKey() Uneix diverses PCollection per clau (el JOIN de Beam) Creuar comandes amb productes
beam.Keys() / beam.Values() Extreu claus o valors
beam.Distinct() Elimina duplicats Sessions úniques
beam.io.ReadFromText / WriteToText Lectura/escriptura de fitxers, inclòs Cloud Storage
beam.io.ReadFromPubSub Llegeix un topic o subscripció El pipeline de pedidos-nuevos
beam.io.WriteToBigQuery Escriu en una taula Destí de tot

Dos aclariments que eviten errors conceptuals importants:

GroupByKey davant de CombinePerKey. GroupByKey porta tots els valors d'una clau a una sola màquina i els materialitza en memòria. Si un sku té tres milions de línies, aquella màquina pot rebentar. CombinePerKey amb una funció associativa i commutativa (sum, max, min) fa agregació parcial a cada treballador abans de moure res per la xarxa: cada worker suma el seu i només viatgen els resultats parcials. Sempre que puguis fer servir CombinePerKey, fes-lo servir; és ordres de magnitud més eficient i no es trenca amb claus calentes.

Map davant de ParDo. Map és sucre sintàctic sobre ParDo. Fes servir Map per a transformacions simples i sense estat; fes servir ParDo amb una classe DoFn quan necessitis inicialització costosa (obrir un client d'API una vegada per treballador, no una vegada per element), mètriques pròpies, o diverses sortides.

  1. Primer pipeline per lots, explicat línia a línia

Objectiu concret: llegir els fitxers CSV d'exportació de comandes que hi ha a gs://alpinashop-catalogo/exportaciones/2026/03/14/, validar-los, netejar-los, calcular un agregat de vendes per SKU i escriure dues coses a alpinashop_analitica: les línies netes a lineas_pedido i l'agregat en una taula nova ventas_diarias_sku. Els registres defectuosos van a un fitxer de quarantena al bucket.

"""
Pipeline per lots d AlpinaShop: exportacions CSV -> BigQuery.
Execucio local:
  python pipeline_pedidos.py --fecha 2026-03-14
Execucio a Dataflow:
  python pipeline_pedidos.py --fecha 2026-03-14 --runner DataflowRunner ...
"""
import argparse
import csv
import io
import logging
from datetime import datetime
from decimal import Decimal, InvalidOperation

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions

PROJECTE = "alpinashop-datos"
DATASET = "alpinashop_analitica"
BUCKET = "alpinashop-catalogo"

Tot el de dalt és Python normal. Les constants en majúscules van al principi perquè canviar de projecte no obligui a buscar cadenes pel fitxer.

class AnalitzarLiniaCSV(beam.DoFn):
    """Converteix una linia de text CSV en un diccionari.

    S implementa com a DoFn i no com a Map perque necessitem
    comptadors propis i perque volem dues sortides: valids i rebutjats.
    """

    SORTIDA_REBUIGS = "rebuigs"

    def __init__(self):
        # Els comptadors de Beam s agreguen entre tots els treballadors
        # i es veuen a la consola de Dataflow. Son la forma correcta
        # d instrumentar un pipeline; els print no serveixen de res.
        self.comptador_ok = beam.metrics.Metrics.counter("analisi", "files_ok")
        self.comptador_ko = beam.metrics.Metrics.counter("analisi", "files_rebutjades")

    def process(self, linia):
        try:
            camps = next(csv.reader(io.StringIO(linia), delimiter=","))
        except Exception as exc:
            self.comptador_ko.inc()
            yield beam.pvalue.TaggedOutput(
                self.SORTIDA_REBUIGS,
                {"linia": linia, "motiu": f"csv_illegible: {exc}"},
            )
            return

        if len(camps) != 8:
            self.comptador_ko.inc()
            yield beam.pvalue.TaggedOutput(
                self.SORTIDA_REBUIGS,
                {"linia": linia, "motiu": f"esperades 8 columnes, n hi ha {len(camps)}"},
            )
            return

        yield {
            "pedido_id": camps[0].strip(),
            "linea_num": camps[1].strip(),
            "fecha_pedido": camps[2].strip(),
            "sku": camps[3].strip().upper(),
            "cantidad": camps[4].strip(),
            "precio_unitario": camps[5].strip().replace(",", "."),
            "descuento_linea": camps[6].strip().replace(",", "."),
            "estado": camps[7].strip().lower(),
        }

Punts importants d'aquesta classe:

  • yield en comptes de return. Un DoFn és un generador: pot emetre zero, un o molts elements per entrada. Un return amb valor no funciona com esperes.
  • TaggedOutput marca l'element perquè surti per una branca diferent del graf. És el mecanisme de Beam per a "això ha anat malament però el pipeline continua". Sense ell, una excepció reintenta el paquet quatre vegades i després mata el treball sencer.
  • replace(",", ".") normalitza els decimals de l'ERP espanyol, que exporta 89,90. És exactament el tipus de brutícia real que hi ha en qualsevol exportació.
  • Els comptadors (Metrics.counter) pugen a la consola de Dataflow i permeten respondre "quantes files es van rebutjar ahir a la nit?" sense obrir un registre.
class ValidarComanda(beam.DoFn):
    """Converteix tipus i aplica regles de negoci."""

    SORTIDA_REBUIGS = "rebuigs"

    def process(self, fila):
        motius = []

        # 1) Data
        try:
            data = datetime.strptime(fila["fecha_pedido"], "%Y-%m-%d").date()
        except ValueError:
            motius.append("data invalida")
            data = None

        # 2) Quantitat: enter estrictament positiu
        try:
            quantitat = int(fila["cantidad"])
            if quantitat <= 0:
                motius.append("quantitat no positiva")
        except ValueError:
            motius.append("quantitat no numerica")
            quantitat = None

        # 3) Imports en Decimal, MAI en float (veure 04-01)
        try:
            preu = Decimal(fila["precio_unitario"])
            descompte = Decimal(fila["descuento_linea"] or "0")
            if preu < 0 or descompte < 0:
                motius.append("import negatiu")
        except InvalidOperation:
            motius.append("import no numeric")
            preu = descompte = None

        # 4) SKU amb el format del cataleg d AlpinaShop
        if not fila["sku"] or len(fila["sku"]) < 4:
            motius.append("sku absent o massa curt")

        if motius:
            yield beam.pvalue.TaggedOutput(
                self.SORTIDA_REBUIGS,
                {"linia": str(fila), "motiu": "; ".join(motius)},
            )
            return

        import_linia = (preu * quantitat) - descompte
        yield {
            "pedido_id": fila["pedido_id"],
            "linea_num": int(fila["linea_num"]),
            "fecha_pedido": data.isoformat(),
            "sku": fila["sku"],
            "cantidad": quantitat,
            "precio_unitario": str(preu),        # BigQuery accepta NUMERIC com a cadena
            "descuento_linea": str(descompte),
            "importe_linea": str(import_linia),
        }

Detall que costa car descobrir en producció: els valors NUMERIC es passen a WriteToBigQuery com a cadena, no com a float. Si converteixes a float per serialitzar, has reintroduït l'error de coma flotant que vam evitar amb tanta cura a 04-01.

Ara el pipeline complet:

def construir_pipeline(pipeline, data):
    ruta_entrada = f"gs://{BUCKET}/exportaciones/{data.replace('-', '/')}/pedidos-*.csv"
    ruta_quarantena = f"gs://{BUCKET}/cuarentena/{data}/rechazos"

    # 1) LLEGIR: cada linia del CSV es un element de la PCollection
    cru = (
        pipeline
        | "LlegirCSV" >> beam.io.ReadFromText(ruta_entrada, skip_header_lines=1)
    )

    # 2) ANALITZAR amb dues sortides
    analitzat = (
        cru
        | "AnalitzarCSV" >> beam.ParDo(AnalitzarLiniaCSV()).with_outputs(
            AnalitzarLiniaCSV.SORTIDA_REBUIGS, main="valids"
        )
    )

    # 3) VALIDAR, tambe amb dues sortides
    validat = (
        analitzat.valids
        | "ValidarComanda" >> beam.ParDo(ValidarComanda()).with_outputs(
            ValidarComanda.SORTIDA_REBUIGS, main="nets"
        )
    )

    linies_netes = validat.nets

    # 4) ESCRIURE el detall a BigQuery
    (
        linies_netes
        | "EscriureLinies" >> beam.io.WriteToBigQuery(
            table=f"{PROJECTE}:{DATASET}.lineas_pedido",
            schema="pedido_id:STRING,linea_num:INTEGER,fecha_pedido:DATE,"
                   "sku:STRING,cantidad:INTEGER,precio_unitario:NUMERIC,"
                   "descuento_linea:NUMERIC,importe_linea:NUMERIC",
            create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
            additional_bq_parameters={
                "timePartitioning": {"type": "DAY", "field": "fecha_pedido"},
                "clustering": {"fields": ["sku"]},
            },
        )
    )

    # 5) AGREGAR vendes per SKU amb CombinePerKey (no GroupByKey)
    (
        linies_netes
        | "ClauSKU" >> beam.Map(
            lambda f: ((f["fecha_pedido"], f["sku"]), Decimal(f["importe_linea"]))
        )
        | "SumarPerSKU" >> beam.CombinePerKey(sum)
        | "FormatarAgregat" >> beam.Map(
            lambda kv: {
                "dia": kv[0][0],
                "sku": kv[0][1],
                "ventas_eur": str(kv[1]),
            }
        )
        | "EscriureAgregat" >> beam.io.WriteToBigQuery(
            table=f"{PROJECTE}:{DATASET}.ventas_diarias_sku",
            schema="dia:DATE,sku:STRING,ventas_eur:NUMERIC",
            create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
            write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE,
        )
    )

    # 6) QUARANTENA: els rebuigs dels dos passos, junts, a Cloud Storage
    (
        (analitzat[AnalitzarLiniaCSV.SORTIDA_REBUIGS],
         validat[ValidarComanda.SORTIDA_REBUIGS])
        | "UnirRebuigs" >> beam.Flatten()
        | "SerialitzarRebuigs" >> beam.Map(
            lambda r: f'{r["motiu"]}\t{r["linia"]}'
        )
        | "EscriureQuarantena" >> beam.io.WriteToText(
            ruta_quarantena, file_name_suffix=".tsv"
        )
    )


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--fecha", required=True, help="Data a processar, yyyy-MM-dd")
    coneguts, resta = parser.parse_known_args()

    opcions = PipelineOptions(resta)
    # Necessari perque els treballadors instal·lin les dependencies del fitxer
    opcions.view_as(SetupOptions).save_main_session = True

    with beam.Pipeline(options=opcions) as p:
        construir_pipeline(p, coneguts.fecha)


if __name__ == "__main__":
    logging.getLogger().setLevel(logging.INFO)
    main()

Cinc coses que mereixen comentari:

  1. with beam.Pipeline(...) as p: en sortir del bloc with, Beam crida run() i espera. Sense el with, cal cridar p.run().wait_until_finish() explícitament.
  2. WRITE_APPEND davant de WRITE_TRUNCATE: el detall s'afegeix (volem històric acumulat); l'agregat es reescriu sencer (volem la foto vigent). Triar malament aquí duplica dades silenciosament.
  3. additional_bq_parameters crea la taula ja particionada i agrupada en clústers si no existia. És la manera de no perdre el que hem après a 04-01 quan la taula la crea el pipeline.
  4. beam.Flatten() uneix diverses PCollection del mateix tipus en una. És la unió de branques del graf, l'equivalent a un UNION ALL.
  5. save_main_session=True serialitza l'àmbit global del mòdul per als treballadors. Sense això, un pipeline que funciona en local falla a Dataflow amb NameError sobre les constants o els imports. És l'error de novell número u.

  1. Execució local amb DirectRunner

Abans de gastar un cèntim, es prova en local. El DirectRunner executa el pipeline a la teva màquina, amb un subconjunt de dades.

# Dades de prova locals
mkdir -p ./pruebas && cat > ./pruebas/pedidos-test.csv <<'EOF'
pedido_id,linea_num,fecha_pedido,sku,cantidad,precio_unitario,descuento_linea,estado
PED-2026-0042,1,2026-03-14,MOCH-40L-AZ,1,89,90,0,confirmado
PED-2026-0042,2,2026-03-14,FRON-300L,2,34.50,5.00,confirmado
PED-2026-0043,1,2026-03-14,CRAM-12P,-1,120.00,0,confirmado
PED-2026-0044,1,data-dolenta,TIEN-2P,1,240.00,0,confirmado
EOF

python pipeline_pedidos.py \
  --fecha 2026-03-14 \
  --runner DirectRunner

Aquest fitxer de prova està brut expressament, i així és com ha de ser qualsevol joc de proves d'un pipeline:

  • La primera línia té nou camps perquè 89,90 hi posa una coma de més. L'analitzador la rebutjarà amb "esperades 8 columnes". És exactament la fallada real d'un ERP mal configurat.
  • La tercera té quantitat negativa: la rebutja el validador.
  • La quarta té una data invàlida: la rebutja el validador.
  • Només la segona passa.

El DirectRunner és deliberadament estricte: comprova la immutabilitat dels elements, serialitza i deserialitza entre passos, i desordena les dades expressament. Si el teu pipeline depèn de l'ordre d'arribada o modifica un objecte in situ, el DirectRunner ho detecta i falla, mentre que al núvol produiria resultats incorrectes de manera intermitent. Que sigui lent i primmirat és la funcionalitat, no un defecte.

Bona pràctica: a més de provar el pipeline sencer, prova les transformacions amb unittest i les utilitats de Beam:

import unittest
import apache_beam as beam
from apache_beam.testing.test_pipeline import TestPipeline
from apache_beam.testing.util import assert_that, equal_to


class TestValidarComanda(unittest.TestCase):
    def test_rebutja_quantitat_negativa(self):
        entrada = [{
            "pedido_id": "PED-1", "linea_num": "1", "fecha_pedido": "2026-03-14",
            "sku": "CRAM-12P", "cantidad": "-1",
            "precio_unitario": "120.00", "descuento_linea": "0", "estado": "confirmado",
        }]
        with TestPipeline() as p:
            sortides = (
                p | beam.Create(entrada)
                  | beam.ParDo(ValidarComanda()).with_outputs(
                        ValidarComanda.SORTIDA_REBUIGS, main="nets")
            )
            assert_that(sortides.nets, equal_to([]), label="sense valids")

beam.Create(...) fabrica una PCollection a partir d'una llista de Python: és la manera d'injectar dades de prova. assert_that amb equal_to compara el contingut sense importar l'ordre, que és el correcte en un sistema distribuït.

  1. Execució gestionada amb DataflowRunner

Amb les proves en verd, al servei real.

# Bucket propi per als artefactes de Dataflow (no barrejar amb el cataleg)
gcloud storage buckets create gs://alpinashop-dataflow \
  --project=alpinashop-datos --location=europe-west1 \
  --uniform-bucket-level-access

# Compte de servei dedicat al pipeline, amb minim privilegi (03-04)
gcloud iam service-accounts create sa-dataflow-pedidos \
  --project=alpinashop-datos \
  --display-name="Pipelines de Dataflow de comandes"

SA="[email protected]"

gcloud projects add-iam-policy-binding alpinashop-datos \
  --member="serviceAccount:$SA" --role="roles/dataflow.worker"
gcloud projects add-iam-policy-binding alpinashop-datos \
  --member="serviceAccount:$SA" --role="roles/bigquery.dataEditor"
gcloud projects add-iam-policy-binding alpinashop-datos \
  --member="serviceAccount:$SA" --role="roles/bigquery.jobUser"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-dataflow \
  --member="serviceAccount:$SA" --role="roles/storage.objectAdmin"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-catalogo \
  --member="serviceAccount:$SA" --role="roles/storage.objectViewer"

I el llançament:

python pipeline_pedidos.py \
  --fecha 2026-03-14 \
  --runner DataflowRunner \
  --project alpinashop-datos \
  --region europe-west1 \
  --temp_location gs://alpinashop-dataflow/temp \
  --staging_location gs://alpinashop-dataflow/staging \
  --service_account_email "$SA" \
  --subnetwork regions/europe-west1/subnetworks/sn-datos-euw1 \
  --no_use_public_ips \
  --machine_type n2-standard-2 \
  --max_num_workers 10 \
  --job_name alpinashop-pedidos-20260314 \
  --labels entorno=produccion,equipo=datos,centro-coste=analitica

Opció per opció, perquè cadascuna té conseqüències:

Opció Què fa i per què
--region europe-west1 On s'executen els treballadors. Ha de coincidir amb la ubicació del conjunt de dades i amb la del bucket, o pagaràs sortida entre regions i afegiràs latència
--temp_location Fitxers temporals i de shuffle intermedi. Obligatori
--staging_location On es puja el codi del pipeline i les seves dependències
--service_account_email La identitat dels treballadors. Sense això es fa servir el compte de servei de Compute per defecte, que sol ser Editor del projecte: un incompliment directe del mínim privilegi de 03-04
--subnetwork Els treballadors arrenquen dins d'alpinashop-vpc, a la subxarxa de dades. No en una xarxa per defecte
--no_use_public_ips Sense IP pública: surten a internet pel Cloud NAT que ja vas configurar a 03-01, i arriben a les API de Google per Private Google Access
--machine_type Tipus de VM del treballador
--max_num_workers Sostre de l'escalat automàtic: la xarxa de seguretat contra una factura desbocada
--job_name Nom visible. Incloure-hi la data ajuda a localitzar reprocessos
--labels Etiquetes de facturació, mateix esquema que la resta d'AlpinaShop

Comprovació de l'estat:

gcloud dataflow jobs list --region=europe-west1 \
  --format="table(id, name, type, state, createTime)" --limit=5

# Detall i metriques d un job concret
gcloud dataflow jobs describe JOB_ID --region=europe-west1
gcloud dataflow metrics list JOB_ID --region=europe-west1 \
  --format="table(name.name, scalar)" --filter="name.name~files_"

Aquesta última comanda retorna els comptadors files_ok i files_rebutjades que hem instrumentat. Aquest és el retorn d'haver fet servir mètriques de Beam en comptes de print.

  1. El temps en streaming: esdeveniment davant de procés

Aquí comença la part difícil, i també la que fa que valgui la pena aprendre Beam en comptes d'improvisar.

En lot, el temps és simple: tens el fitxer, té un final, processes i acabes. En temps real no hi ha final, i apareix una distinció que ho canvia tot:

  • Temps de l'esdeveniment (event time): quan va passar el fet al món real. El client va prémer "Comprar" a les 10:00:00.
  • Temps de procés (processing time): quan la dada arriba al teu pipeline. 10:00:03, o 10:07:12 si el mòbil era en un túnel, o 11:30 si l'aplicació va desar l'esdeveniment en local fins a recuperar cobertura.

Amb AlpinaShop és molt concret. Un client al metro de Barcelona navega pel catàleg, afegeix una motxilla a la cistella a les 18:42, perd cobertura, i l'aplicació envia els esdeveniments acumulats a les 18:51.

flowchart LR
    subgraph Real["Temps de l esdeveniment (mon real)"]
        E1["18:41 ver_producto"]
        E2["18:42 anadir_carrito"]
        E3["18:44 iniciar_pago"]
    end
    subgraph Proces["Temps de proces (arribada al pipeline)"]
        P1["18:51 tots tres alhora"]
    end
    E1 --> P1
    E2 --> P1
    E3 --> P1

Ara la pregunta de negoci: quants productes es van veure entre les 18:40 i les 18:45? Si agrupes per temps de procés, la resposta és zero, i és falsa. Si agrupes per temps de l'esdeveniment, la resposta inclou aquest client, i és la correcta.

Beam agrupa per temps de l'esdeveniment per defecte. Aquesta és la seva decisió de disseny més important i la raó per la qual els seus números quadren amb els de l'informe per lots de l'endemà. Un sistema que agrupa per temps de procés dona resultats que depenen de la xarxa i que mai no són reproduïbles: reprocessar el mateix dia dona un resultat diferent.

  1. Marques d'aigua, finestres, activadors i dades tardanes

Si esperes per temps de l'esdeveniment, sorgeix la pregunta inevitable: quan deixes d'esperar? Un esdeveniment de les 18:42 podria arribar demà. Tanques la finestra o esperes eternament?

La resposta de Beam són quatre mecanismes que es combinen.

Finestres (Windowing)

Trossegen la PCollection no acotada en trossos finits sobre els quals sí que es pot agregar.

Tipus Definició Ús a AlpinaShop
Fixa (fixed/tumbling) Intervals contigus que no se solapen Comandes per hora per al tauler de campanya
Lliscant (sliding) Intervals que se solapen Mitjana mòbil de visites: finestra de 30 min cada 5 min
De sessió (session) S'agrupen esdeveniments separats per menys d'un buit donat Sessions de navegació reals: tot el que fa un usuari fins a estar 30 min inactiu
Global Una sola finestra infinita Només amb activadors explícits
from apache_beam import window

# Finestra fixa d 1 hora: per al comptador de comandes
per_hora = esdeveniments | "FinestraHora" >> beam.WindowInto(
    window.FixedWindows(60 * 60)
)

# Finestra lliscant: mitjana mobil de 30 min, actualitzada cada 5
mitjana_mobil = esdeveniments | "FinestraLliscant" >> beam.WindowInto(
    window.SlidingWindows(size=30 * 60, period=5 * 60)
)

# Finestra de sessio: agrupa l activitat d un usuari
sessions = esdeveniments | "FinestraSessio" >> beam.WindowInto(
    window.Sessions(gap_size=30 * 60)
)

La finestra de sessió és especialment elegant i no té equivalent senzill en SQL: la seva durada no es fixa per endavant, la determinen les mateixes dades. És exactament la definició de "sessió de navegació" que necessita la taula visitas de 04-01.

Marca d'aigua (watermark)

La marca d'aigua és l'estimació que fa el sistema de "ja no espero esdeveniments anteriors a aquest instant". Dataflow la calcula automàticament observant les marques de temps de les dades que van arribant: si de Pub/Sub fa deu minuts que arriben esdeveniments de les 18:50 endavant, la marca d'aigua avança més enllà de les 18:45 i les finestres anteriors es poden tancar.

No és una garantia, és una heurística. Sempre pot arribar alguna cosa després. Per això existeixen els altres dos mecanismes.

Activadors (triggers)

L'activador decideix quan emetre el resultat d'una finestra. Per defecte, en passar la marca d'aigua: un resultat per finestra, quan es considera completa.

Però el tauler de la campanya de tardor no pot esperar una hora a veure el primer número. Amb un activador s'emeten resultats parcials:

from apache_beam.transforms.trigger import (
    AfterWatermark, AfterProcessingTime, AccumulationMode
)

comandes_per_hora = (
    esdeveniments
    | "Finestra" >> beam.WindowInto(
        window.FixedWindows(60 * 60),
        trigger=AfterWatermark(
            early=AfterProcessingTime(60),      # avanc parcial cada minut
            late=AfterProcessingTime(10 * 60),  # correccions cada 10 min si arriba tard
        ),
        allowed_lateness=2 * 60 * 60,           # acceptem fins a 2 h de retard
        accumulation_mode=AccumulationMode.ACCUMULATING,
    )
    | "Comptar" >> beam.CombinePerKey(sum)
)

Interpretació pràctica d'aquesta configuració per a AlpinaShop:

  • Cada minut s'emet un resultat provisional: el tauler es mou i la gent veu que el sistema és viu.
  • Quan passa la marca d'aigua, s'emet el resultat considerat definitiu.
  • Durant dues hores més (allowed_lateness), si arriben esdeveniments d'aquella hora —el client del metro—, s'emet una correcció cada deu minuts.
  • Passades les dues hores, el que arribi es descarta. Aquest descart és una decisió de negoci explícita, no un accident, i cal instrumentar-lo amb un comptador per saber quant es perd.

AccumulationMode.ACCUMULATING significa que cada emissió conté el total acumulat de la finestra, així que el destí ha de sobreescriure. L'alternativa, DISCARDING, emet només el que és nou des de l'última emissió, i el destí ha de sumar. Confondre-les produeix xifres doblades o dividides, i és un error molt difícil de veure en un tauler.

Dades tardanes

Tot el que arriba després de la marca d'aigua és una dada tardana. Amb allowed_lateness decideixes quant de temps l'acceptes. La regla és simple: com més esperes, més correctes són els números i més recursos (estat en memòria) consumeix el pipeline. Dues hores per a un comerç electrònic europeu és un valor raonable; dos dies seria caríssim i no canviaria les decisions.

  1. El pipeline en temps real de pedidos-nuevos

A la propera lliçó, 04-04, crearem el topic de Pub/Sub pedidos-nuevos, on l'aplicació Flask publicarà un missatge JSON cada vegada que es confirmi una comanda. Anticipem el consumidor, perquè és el cas d'ús canònic de Dataflow.

"""
Pipeline en temps real: pedidos-nuevos (Pub/Sub) -> BigQuery.
S executa de forma continua, sense fi.
"""
import json
import logging
from datetime import datetime

import apache_beam as beam
from apache_beam import window
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions
from apache_beam.transforms.trigger import AfterWatermark, AfterProcessingTime, AccumulationMode

PROJECTE = "alpinashop-datos"
SUBSCRIPCIO = f"projects/{PROJECTE}/subscriptions/sub-analitica"


class DescodificarComanda(beam.DoFn):
    SORTIDA_REBUIGS = "rebuigs"

    def process(self, missatge, moment=beam.DoFn.TimestampParam):
        try:
            dades = json.loads(missatge.data.decode("utf-8"))
        except Exception as exc:
            yield beam.pvalue.TaggedOutput(
                self.SORTIDA_REBUIGS,
                {"payload": str(missatge.data[:500]), "motiu": f"json invalid: {exc}"},
            )
            return

        # Els atributs del missatge de Pub/Sub viatgen a part del cos
        origen = missatge.attributes.get("origen", "desconegut")

        yield {
            "pedido_id": dades["pedido_id"],
            "fecha_pedido": dades["fecha"][:10],
            "momento_pedido": dades["fecha"],
            "cliente_id": dades.get("cliente_id"),
            "canal": origen,
            "estado": "confirmado",
            "total_pedido": str(dades["total"]),
            "momento_ingesta": datetime.utcnow().isoformat(),
        }


def main():
    opcions = PipelineOptions(streaming=True, save_main_session=True)
    opcions.view_as(StandardOptions).streaming = True

    with beam.Pipeline(options=opcions) as p:
        missatges = (
            p
            | "LlegirPubSub" >> beam.io.ReadFromPubSub(
                subscription=SUBSCRIPCIO,
                with_attributes=True,
                timestamp_attribute="momento_evento",   # <-- clau
            )
        )

        descodificats = (
            missatges
            | "Descodificar" >> beam.ParDo(DescodificarComanda()).with_outputs(
                DescodificarComanda.SORTIDA_REBUIGS, main="valids")
        )

        # A) Detall en temps real, fila a fila
        (
            descodificats.valids
            | "EscriureComandes" >> beam.io.WriteToBigQuery(
                table=f"{PROJECTE}:alpinashop_analitica.pedidos_streaming",
                schema="pedido_id:STRING,fecha_pedido:DATE,momento_pedido:TIMESTAMP,"
                       "cliente_id:STRING,canal:STRING,estado:STRING,"
                       "total_pedido:NUMERIC,momento_ingesta:TIMESTAMP",
                create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
                write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
                method="STORAGE_WRITE_API",
            )
        )

        # B) Agregat per hora i canal, per al tauler de campanya
        (
            descodificats.valids
            | "FinestraHora" >> beam.WindowInto(
                window.FixedWindows(3600),
                trigger=AfterWatermark(early=AfterProcessingTime(60)),
                allowed_lateness=7200,
                accumulation_mode=AccumulationMode.ACCUMULATING,
            )
            | "ClauCanal" >> beam.Map(lambda d: (d["canal"], float(d["total_pedido"])))
            | "SumarPerCanal" >> beam.CombinePerKey(sum)
            | "Formatar" >> beam.Map(
                lambda kv, w=beam.DoFn.WindowParam: {
                    "hora_inicio": w.start.to_utc_datetime().isoformat(),
                    "canal": kv[0],
                    "ventas_eur": round(kv[1], 2),
                }
            )
            | "EscriureAgregat" >> beam.io.WriteToBigQuery(
                table=f"{PROJECTE}:alpinashop_analitica.ventas_por_hora",
                schema="hora_inicio:TIMESTAMP,canal:STRING,ventas_eur:FLOAT",
                create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
                write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
                method="STORAGE_WRITE_API",
            )
        )


if __name__ == "__main__":
    logging.getLogger().setLevel(logging.INFO)
    main()

Tres decisions que cal entendre:

  • timestamp_attribute="momento_evento": li diu a Beam que faci servir l'atribut del missatge com a temps de l'esdeveniment, en comptes del moment de publicació a Pub/Sub. Sense això, un missatge retingut nou minuts al mòbil es comptabilitzaria a l'hora equivocada. És la línia que fa que tot l'apartat 9 funcioni de debò, i s'oblida constantment.
  • method="STORAGE_WRITE_API": fa servir l'API moderna d'escriptura, més barata i amb semàntica d'exactament una vegada, en comptes de les insercions en streaming clàssiques.
  • Es llegeix d'una subscripció, no del topic. Amb una subscripció, si el pipeline s'atura, els missatges s'acumulen i es recuperen en arrencar. Llegir del topic directament fa que Dataflow creï una subscripció efímera i es perdi tot el que s'ha publicat mentre el pipeline està aturat.

  1. Plantilles de Dataflow: la via pràctica

Tot l'anterior és potent i també és feina. Per a tasques comunes existeix una cosa molt més simple: les plantilles.

Una plantilla és un pipeline ja compilat i parametritzat que es llança sense escriure ni compilar codi. Google en publica desenes.

Plantilla Què fa
Pub/Sub Subscription to BigQuery Llegeix JSON d'una subscripció i l'insereix en una taula
Cloud Storage Text to BigQuery Carrega fitxers amb una funció JavaScript de transformació
JDBC to BigQuery Bolca una base de dades relacional
BigQuery to Cloud Storage (Parquet) Exporta
Datastream to BigQuery Aplica CDC (ho veurem a 04-05)
Bulk Compress/Decompress Utilitats sobre el bucket

Llançar la de Pub/Sub a BigQuery per a AlpinaShop és una línia:

gcloud dataflow jobs run alpinashop-pedidos-a-bq \
  --gcs-location gs://dataflow-templates-europe-west1/latest/PubSub_Subscription_to_BigQuery \
  --region europe-west1 \
  --service-account-email "[email protected]" \
  --subnetwork regions/europe-west1/subnetworks/sn-datos-euw1 \
  --disable-public-ips \
  --max-workers 5 \
  --parameters \
inputSubscription=projects/alpinashop-datos/subscriptions/sub-analitica,\
outputTableSpec=alpinashop-datos:alpinashop_analitica.pedidos_streaming,\
outputDeadletterTable=alpinashop-datos:alpinashop_analitica.pedidos_streaming_errores

Fixa't en outputDeadletterTable: els missatges que no encaixin amb l'esquema van a una taula d'errors en comptes de bloquejar el pipeline. És el mateix patró de quarantena que vam programar a mà, ja resolt.

Hi ha dos sabors:

  • Plantilles clàssiques: el graf es compila en crear-la; els paràmetres només poden ser valors en temps d'execució (ValueProvider).
  • Plantilles flexibles (Flex Templates): el pipeline s'empaqueta com a imatge de contenidor a Artifact Registry i el graf es construeix en llançar. Són més flexibles i són l'opció recomanada avui per a plantilles pròpies.

Crear una plantilla flexible del pipeline per lots d'AlpinaShop, perquè Composer o Cloud Scheduler la invoquin a 04-06:

# 1) Construir la imatge i publicar la plantilla
gcloud dataflow flex-template build \
  gs://alpinashop-dataflow/plantillas/pedidos-lote.json \
  --image-gcr-path europe-west1-docker.pkg.dev/alpinashop-prod/alpinashop/dataflow-pedidos:1.0.0 \
  --sdk-language PYTHON \
  --flex-template-base-image PYTHON3 \
  --py-path . \
  --env FLEX_TEMPLATE_PYTHON_PY_FILE=pipeline_pedidos.py \
  --env FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE=requirements.txt

# 2) Executar-la amb parametres
gcloud dataflow flex-template run "pedidos-$(date +%Y%m%d-%H%M%S)" \
  --template-file-gcs-location gs://alpinashop-dataflow/plantillas/pedidos-lote.json \
  --region europe-west1 \
  --service-account-email "[email protected]" \
  --parameters fecha=2026-03-14

La imatge va al mateix Artifact Registry que ja fa servir el catàleg (europe-west1-docker.pkg.dev/alpinashop-prod/alpinashop/), cosa que manté un únic inventari d'artefactes, coherent amb el que veurem al mòdul 6.

El criteri d'AlpinaShop: per a "Pub/Sub a BigQuery" sense transformació, la plantilla de Google, sense escriure ni una línia. Per al pipeline per lots amb validació de negoci pròpia, codi Beam empaquetat com a plantilla flexible. Escriure codi només quan aporta lògica que la plantilla no té.

  1. Escalat automàtic i Dataflow Prime

Dataflow ajusta el nombre de treballadors durant l'execució. En lot, mira la feina pendent; en temps real, mira el retard de la cua (backlog) i l'ús de CPU.

--num_workers 2          # amb quants comenca
--max_num_workers 20     # sostre dur: el control de cost
--autoscaling_algorithm THROUGHPUT_BASED   # per defecte en streaming

Posar un max_num_workers baix no sempre estalvia: un pipeline que triga deu hores amb dos treballadors pot costar el mateix que un que triga una hora amb vint, perquè es factura per treballador i hora. El que sí que evita el sostre és la sorpresa d'un treball amb un bucle patològic consumint cent màquines tota la nit.

Dataflow Prime és l'evolució del servei amb tres diferències pràctiques:

Aspecte Dataflow clàssic Dataflow Prime
Recursos Tries tipus de màquina per a tot el pipeline Ajust vertical automàtic de memòria per pas
Escalat Horitzontal Horitzontal + vertical
Facturació Per vCPU, memòria i disc per hora Per Unitats de Càlcul de Dades (DCU)
Diagnòstic Mètriques Recomanacions automàtiques de colls d'ampolla
Quan fer-lo servir Pipelines estables i ben dimensionats Pipelines amb passos de consum molt desigual

L'avantatge real de Prime apareix quan un pas del pipeline necessita molta memòria i els altres no: al model clàssic dimensiones totes les màquines per al pas més exigent i malbarates recursos la resta del temps. S'activa amb --dataflow_service_options=enable_prime.

Per a AlpinaShop, amb pipelines modestos i predictibles, el model clàssic amb n2-standard-2 és suficient i més fàcil de raonar a la factura. Prime queda anotat per a quan el volum ho justifiqui.

  1. Monitoratge, paral·lelisme i biaix de dades

La consola de Dataflow mostra el graf d'execució amb cada pas i, a cadascun, elements processats, temps de CPU i estat. És la millor interfície de depuració de dades de la plataforma, i per això insistim a anomenar els passos.

Mètriques que cal mirar:

Mètrica Què indica Llindar d'alarma
Retard del sistema (system lag) Segons que fa que l'element més antic està sense processar En temps real, si creix de manera sostinguda, el pipeline no dona l'abast
Frescor de la dada (data freshness) Antiguitat de la dada més recent ja emesa S'ha de mantenir estable
Elements per segon per pas On és el coll d'ampolla El pas més lent mana
Ús de vCPU Si l'escalat serveix Alt i amb retard creixent: hi ha un límit estructural
Treballadors actuals Comportament de l'escalat automàtic Enganxat al màxim: pujar el sostre o arreglar el pipeline

El problema més freqüent i més difícil de diagnosticar és el biaix de dades (data skew): una clau concentra una proporció enorme dels elements. A AlpinaShop és molt fàcil que passi: si agrupes visites per sku i la motxilla estrella acumula el 40 % del trànsit, un sol treballador processarà el 40 % de la feina mentre els altres dinou esperen. El símptoma és inconfusible: l'escalat no millora res i hi ha un pas amb un treballador al 100 % i la resta ociosos.

Tres remeis, per ordre de preferència:

# 1) EL MILLOR: fer servir CombinePerKey en comptes de GroupByKey.
#    L agregacio parcial passa a cada treballador abans del shuffle.
vendes = linies | "Sumar" >> beam.CombinePerKey(sum)

# 2) Si necessites GroupByKey de debo: afegir sal a la clau
import random

def salar(element, n=20):
    clau, valor = element
    return (f"{clau}#{random.randint(0, n - 1)}", valor)

resultat = (
    linies
    | "Salar"          >> beam.Map(salar)
    | "AgruparParcial" >> beam.CombinePerKey(sum)
    | "TreureSal"      >> beam.Map(lambda kv: (kv[0].split("#")[0], kv[1]))
    | "AgruparFinal"   >> beam.CombinePerKey(sum)
)

# 3) Si un costat del JOIN es petit (el cataleg de productes):
#    fer servir entrada lateral en comptes de CoGroupByKey, i evitar el shuffle
productes = p | "LlegirProductes" >> beam.io.ReadFromBigQuery(query=SQL_PRODUCTOS)
cataleg = beam.pvalue.AsDict(productes | beam.Map(lambda r: (r["sku"], r)))

enriquit = linies | "Enriquir" >> beam.Map(
    lambda linia, cat: {**linia, "categoria": cat.get(linia["sku"], {}).get("categoria")},
    cat=cataleg,
)

La tècnica 2, la sal, mereix explicació: en afegir un sufix aleatori a la clau, la motxilla estrella es converteix en vint claus diferents que es reparteixen entre vint treballadors; després es treu la sal i es fa una segona agregació sobre vint valors, que és trivial. Només funciona amb operacions associatives, però cobreix gairebé tots els casos d'agregació.

La tècnica 3, l'entrada lateral (side input), és la que més faràs servir a AlpinaShop: el catàleg de productes són uns quants milers de files i cap a la memòria de cada treballador. Difondre'l evita completament el shuffle del JOIN. Compte amb el límit: si l'entrada lateral no cap a la memòria, el pipeline es degrada moltíssim.

  1. Cost: què es paga exactament

Dataflow no factura per pipeline ni per dada processada: factura els recursos consumits pels treballadors.

Concepte Com es mesura Ordre de magnitud (verificar a la documentació oficial)
vCPU Per vCPU i hora ~0,05-0,07 $ (lot); una mica més en streaming
Memòria Per GB i hora ~0,003-0,004 $
Disc persistent Per GB i hora ~0,00005 $ estàndard
Shuffle (lot) Per GB processats al servei de shuffle ~0,011 $/GB
Streaming Engine Per GB de dades en streaming processades ~0,018 $/GB
Dataflow Prime Per DCU Model unificat

Un exemple realista per a AlpinaShop: el pipeline per lots nocturn, amb 4 treballadors n2-standard-2 (2 vCPU, 8 GB) durant 20 minuts, surt per cèntims. El pipeline en temps real, en canvi, funciona 24×7: 2 treballadors permanents són unes 1.440 hores de vCPU al mes, de l'ordre de 80-100 € mensuals. Aquest és el número que sorprèn tothom, i és la raó de l'advertiment del principi de la lliçó.

Consells concrets per no pagar de més:

  • Activa Streaming Engine i Shuffle Service (--enable_streaming_engine, --experiments=shuffle_mode=service). Mouen el shuffle i l'estat fora dels treballadors, permetent màquines més petites i discos molt menors. Gairebé sempre surt a compte.
  • Redueix el disc. Amb Streaming Engine, --disk_size_gb=30 n'hi ha prou; el valor per defecte és molt més gran i es paga per hora.
  • Posa --max_num_workers sempre. És el fre de mà.
  • Fes servir VM Spot en lot tolerant a interrupcions: --flexrs_goal=COST_OPTIMIZED endarrereix l'arrencada fins a 6 hores a canvi d'un descompte notable. Perfecte per al bolcat nocturn; inacceptable per a temps real.
  • De debò necessites streaming? Un microlot cada 15 minuts amb la plantilla de Cloud Storage a BigQuery costa una fracció d'un pipeline permanent. Si el negoci tolera 15 minuts de retard, la resposta és no.
  • Vigila els pipelines de streaming oblidats. Un job de proves que ningú no va aturar és la partida fantasma més comuna a la factura de dades. El detectaràs amb gcloud dataflow jobs list --status=active.

  1. Quan Dataflow no és la resposta

L'honestedat sobre els límits és part de saber fer servir una eina.

Situació Millor opció Per què
Transformació expressable en SQL sobre dades que ja són a BigQuery BigQuery (04-01) No moguis dades per transformar-les; fes servir INSERT ... SELECT o una vista materialitzada
Ja tens codi Spark o l'equip sap Spark, no Beam Dataproc (04-03) Reescriure a Beam és un cost sense retorn clar
Càrrega simple de fitxers a BigQuery, sense lògica bq load És gratis; Dataflow costaria diners per fer el mateix
Integrar orígens externs amb un equip sense perfil de programació Data Fusion (04-05) Interfície visual, connectors llestos
Decidir l'ordre en què s'executen diversos processos Composer o Workflows (04-06) Dataflow executa un pipeline; no orquestra els altres
Reaccionar a un esdeveniment puntual amb poca lògica Cloud Functions (06-03) Un pipeline sencer per processar un fitxer és desproporcionat
Copiar una base de dades amb captura de canvis Datastream (04-05) CDC gestionat, sense codi

La confusió més habitual és la de la penúltima fila. Dataflow no és un orquestrador. Pot executar un pipeline complexíssim, però no sap "primer exporta Cloud SQL, després carrega a BigQuery, després llança aquest pipeline, i si alguna cosa falla avisa la Marta". Això és exactament el que resol 04-06.

Errors habituals i consells

Oblidar save_main_session=True. El pipeline funciona en local i falla a Dataflow amb NameError: name 'PROJECTE' is not defined. Els treballadors no reben l'àmbit global del mòdul si no els ho dius.

No fixar la versió del SDK. Els treballadors fan servir la versió amb què vas llançar el pipeline. Un pip install apache-beam[gcp] sense versió fa que el mateix codi funcioni avui i falli demà. Fixa la versió a requirements.txt.

No posar noms als passos. Sense "Nom" >>, el graf de la consola és il·legible i les mètriques no diuen res. A més, canviar el nom d'un pas impedeix actualitzar un pipeline en temps real en marxa (--update), perquè Beam no pot mapar l'estat del pas antic al nou.

Fer servir GroupByKey on hi cap CombinePerKey. És la diferència entre un pipeline que escala i un que cau amb una clau calenta.

Llegir d'un topic en comptes d'una subscripció. Amb topic, Dataflow crea una subscripció temporal i es perd tot el que s'ha publicat mentre el pipeline està aturat. Amb subscripció, els missatges s'acumulen i es recuperen.

Oblidar timestamp_attribute. Les finestres es calculen sobre el moment de publicació en comptes del moment de l'esdeveniment, i els números no quadren amb els del procés per lots. És subtil i greu.

Deixar un pipeline de streaming de proves encès. Factura 24×7. Posa'ls etiquetes i revisa gcloud dataflow jobs list --status=active amb periodicitat.

No fer servir el compte de servei propi. Sense --service_account_email, els treballadors fan servir el compte per defecte de Compute Engine, que en molts projectes és Editor. Un pipeline amb permisos d'Editor sobre producció és un risc innecessari.

Consell: escriu sempre la branca d'errors. Un pipeline sense sortida de rebuigs no és un pipeline de producció. Si no pots veure què s'ha descartat i per què, no pots confiar en els números.

Consell: --update per modificar un pipeline en temps real. Permet substituir el codi conservant l'estat en curs, en comptes de drenar-lo i arrencar de zero. Requereix compatibilitat del graf, cosa que és una altra raó per no canviar el nom dels passos alegrement.

Consell: drena, no cancel·lis. gcloud dataflow jobs drain deixa de llegir entrades noves i acaba el que té en vol. cancel mata el treball i pot perdre dades en procés.

Exercicis

Exercici 1: pipeline de ressenyes amb quarantena

Escriu un pipeline per lots en Beam que llegeixi gs://alpinashop-catalogo/exportaciones/2026/03/opiniones-*.csv amb les columnes opinion_id,sku,fecha,puntuacion,texto,pais, i que:

  • rebutgi les files la puntuacion de les quals no sigui un enter entre 1 i 5, o el pais de les quals no sigui un codi de dues lletres;
  • normalitzi el sku a majúscules i retalli el texto a 500 caràcters;
  • escrigui les vàlides a alpinashop-datos:alpinashop_analitica.opiniones;
  • escrigui les rebutjades, amb el seu motiu, a gs://alpinashop-catalogo/cuarentena/opiniones/;
  • porti un comptador de vàlides i rebutjades visible a Dataflow.

Prova'l amb DirectRunner i un fitxer local amb almenys dues files defectuoses.

Exercici 2: finestres i activadors per al tauler de campanya

Direcció vol un tauler de la campanya de tardor amb les vendes per hora i per país. Requisits: agrupar per temps de l'esdeveniment, mostrar un avanç parcial cada 30 segons perquè el tauler es vegi viu, acceptar esdeveniments amb fins a 90 minuts de retard emetent correccions, i que cada emissió contingui el total acumulat de l'hora. Escriu únicament el fragment de WindowInto i l'agregació, i explica què escriu exactament el destí i per què el mode d'acumulació escollit obliga a un write_disposition concret.

Exercici 3: diagnòstic d'un pipeline que no escala

El pipeline en temps real de visites porta tres dies funcionant. Des d'ahir, el system lag ha passat de 4 segons a 22 minuts i continua creixent. L'escalat automàtic ha arribat als 20 treballadors (el seu màxim), però l'ús mitjà de vCPU del conjunt és del 18 %. Al graf, el pas AgruparPerSKU mostra un treballador amb 6 hores de temps de CPU i els altres amb menys de 10 minuts. Ahir, màrqueting va llançar una campanya de la motxilla MOCH-40L-AZ que ha multiplicat per dotze les seves visites.

Diagnostica la causa, explica per què pujar --max_num_workers a 50 no arreglaria res, i proposa dues solucions de codi amb la seva diferència pràctica.

Solucions

Solució 1

import csv, io, logging, re
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions

PROJECTE, DATASET = "alpinashop-datos", "alpinashop_analitica"
BUCKET = "alpinashop-catalogo"
RE_PAIS = re.compile(r"^[A-Z]{2}$")


class ProcessarRessenya(beam.DoFn):
    REBUIGS = "rebuigs"

    def __init__(self):
        self.ok = beam.metrics.Metrics.counter("ressenyes", "valides")
        self.ko = beam.metrics.Metrics.counter("ressenyes", "rebutjades")

    def process(self, linia):
        try:
            c = next(csv.reader(io.StringIO(linia)))
        except Exception as exc:
            self.ko.inc()
            yield beam.pvalue.TaggedOutput(self.REBUIGS,
                {"linia": linia, "motiu": f"csv illegible: {exc}"})
            return

        if len(c) != 6:
            self.ko.inc()
            yield beam.pvalue.TaggedOutput(self.REBUIGS,
                {"linia": linia, "motiu": f"esperades 6 columnes, n hi ha {len(c)}"})
            return

        opinion_id, sku, data, punt, text, pais = [x.strip() for x in c]
        motius = []

        try:
            puntuacio = int(punt)
            if not 1 <= puntuacio <= 5:
                motius.append("puntuacio fora del rang 1-5")
        except ValueError:
            motius.append("puntuacio no numerica")
            puntuacio = None

        pais = pais.upper()
        if not RE_PAIS.match(pais):
            motius.append(f"pais invalid: {pais}")

        if not sku:
            motius.append("sku buit")

        if motius:
            self.ko.inc()
            yield beam.pvalue.TaggedOutput(self.REBUIGS,
                {"linia": linia, "motiu": "; ".join(motius)})
            return

        self.ok.inc()
        yield {
            "opinion_id": opinion_id,
            "sku": sku.upper(),
            "fecha": data,
            "puntuacion": puntuacio,
            "texto": text[:500],
            "pais": pais,
        }


def main():
    opcions = PipelineOptions()
    opcions.view_as(SetupOptions).save_main_session = True

    with beam.Pipeline(options=opcions) as p:
        sortides = (
            p
            | "Llegir" >> beam.io.ReadFromText(
                f"gs://{BUCKET}/exportaciones/2026/03/opiniones-*.csv",
                skip_header_lines=1)
            | "Processar" >> beam.ParDo(ProcessarRessenya()).with_outputs(
                ProcessarRessenya.REBUIGS, main="valides")
        )

        (sortides.valides
         | "EscriureBQ" >> beam.io.WriteToBigQuery(
             table=f"{PROJECTE}:{DATASET}.opiniones",
             schema="opinion_id:STRING,sku:STRING,fecha:DATE,"
                    "puntuacion:INTEGER,texto:STRING,pais:STRING",
             create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
             write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
             additional_bq_parameters={
                 "timePartitioning": {"type": "DAY", "field": "fecha"},
                 "clustering": {"fields": ["sku"]},
             }))

        (sortides[ProcessarRessenya.REBUIGS]
         | "Serialitzar" >> beam.Map(lambda r: f'{r["motiu"]}\t{r["linia"]}')
         | "EscriureQuarantena" >> beam.io.WriteToText(
             f"gs://{BUCKET}/cuarentena/opiniones/rechazos",
             file_name_suffix=".tsv"))


if __name__ == "__main__":
    logging.getLogger().setLevel(logging.INFO)
    main()

Fitxer de prova amb defectes deliberats:

opinion_id,sku,fecha,puntuacion,texto,pais
OPI-0001,moch-40l-az,2026-03-10,5,Molt comoda per a travessies llargues,es
OPI-0002,FRON-300L,2026-03-10,9,Puntuacio impossible,FR
OPI-0003,CRAM-12P,2026-03-11,quatre,Puntuacio no numerica,ES
OPI-0004,TIEN-2P,2026-03-11,4,Pais mal format,ESPANYA

La primera passa (el sku es normalitza a majúscules i es a ES); les tres següents van a quarantena amb motius diferents. Mètriques esperades: valides=1, rebutjades=3.

Solució 2

from apache_beam import window
from apache_beam.transforms.trigger import (
    AfterWatermark, AfterProcessingTime, AccumulationMode
)

vendes_hora_pais = (
    esdeveniments
    | "FinestraCampanya" >> beam.WindowInto(
        window.FixedWindows(3600),                     # 1 hora, temps de l ESDEVENIMENT
        trigger=AfterWatermark(
            early=AfterProcessingTime(30),             # avanc cada 30 s
            late=AfterProcessingTime(300),             # correccions cada 5 min
        ),
        allowed_lateness=90 * 60,                      # 90 minuts de retard
        accumulation_mode=AccumulationMode.ACCUMULATING,
    )
    | "ClauPais" >> beam.Map(lambda d: (d["pais"], float(d["total_pedido"])))
    | "SumarPerPais" >> beam.CombinePerKey(sum)
    | "Formatar" >> beam.Map(
        lambda kv, w=beam.DoFn.WindowParam: {
            "hora_inicio": w.start.to_utc_datetime().isoformat(),
            "pais": kv[0],
            "ventas_eur": round(kv[1], 2),
        })
)

Què escriu exactament el destí. Amb ACCUMULATING, cada emissió d'una finestra conté el total acumulat des de l'inici de la finestra, no l'increment. Per a la finestra 10:00-11:00 i el país ES, el destí rebrà una seqüència com: 120 € (a les 10:00:30), 345 € (10:01:00), … 4.210 € (en passar la marca d'aigua), 4.235 € (a les 11:20, correcció per un esdeveniment tardà).

Per què obliga a un write_disposition concret. Si es fes WRITE_APPEND sobre la taula final, aquestes sis emissions se sumarien i el tauler mostraria més de 9.000 € on n'hi ha 4.235: la dada quedaria multiplicada. Les opcions correctes són:

  • Escriure en una taula d'estats intermedis amb WRITE_APPEND i que el tauler consulti només l'última emissió per finestra i clau, per exemple amb QUALIFY ROW_NUMBER() OVER (PARTITION BY hora_inicio, pais ORDER BY momento_emision DESC) = 1. És l'opció recomanada en temps real, perquè conserva l'historial de correccions i permet auditar.
  • O bé fer servir un destí que admeti sobreescriptura per clau (MERGE posterior sobre hora_inicio + pais).

Si en canvi es triés AccumulationMode.DISCARDING, cada emissió portaria només l'increment i llavors WRITE_APPEND amb una suma posterior seria el correcte. La combinació mode d'acumulació + disposició d'escriptura s'ha de decidir conjuntament; equivocar-se produeix taulers amb xifres inflades que ningú no detecta fins que algú creua el número amb comptabilitat.

Solució 3

Diagnòstic: biaix de dades (data skew) per clau calenta.

Les tres evidències apunten al mateix i es reforcen entre si:

  1. 20 treballadors amb 18 % de CPU mitjana. Si el pipeline estigués realment saturat, la CPU estaria alta. Un ús baix amb retard creixent significa que els treballadors esperen, no que treballin.
  2. Un treballador amb 6 hores de CPU i la resta amb 10 minuts. Aquesta és la signatura exacta del biaix: la feina no es reparteix.
  3. La campanya de MOCH-40L-AZ. En un GroupByKey per sku, tots els esdeveniments de la motxilla estrella s'envien pel shuffle a un únic treballador, perquè la clau determina el destí. Els altres no el poden ajudar.

Per què pujar a 50 treballadors no arregla res. La unitat de paral·lelisme d'una agrupació per clau és la clau, no l'element. Els esdeveniments de MOCH-40L-AZ continuaran anant tots al mateix lloc, amb 20 treballadors o amb 500. Els 30 nous estarien ociosos i facturant: el retard continuaria creixent i el cost es multiplicaria per 2,5. Escalar horitzontalment no resol un problema de distribució.

Solució A — substituir GroupByKey per CombinePerKey (la bona, si l'operació ho permet):

# ABANS: tots els esdeveniments de la clau calenta viatgen a un treballador
visites_per_sku = (
    esdeveniments
    | "ClauSKU" >> beam.Map(lambda e: (e["sku"], 1))
    | "Agrupar" >> beam.GroupByKey()
    | "Comptar" >> beam.Map(lambda kv: (kv[0], len(list(kv[1]))))
)

# DESPRES: cada treballador suma el seu abans del shuffle
visites_per_sku = (
    esdeveniments
    | "ClauSKU" >> beam.Map(lambda e: (e["sku"], 1))
    | "Comptar" >> beam.CombinePerKey(sum)
)

Per la xarxa viatgen sumes parcials, una per treballador i clau, en comptes de milions d'elements individuals. Amb 20 treballadors, el treballador de MOCH-40L-AZ rep 20 números en comptes de dotze milions d'esdeveniments. És una línia de codi i sol resoldre el problema completament.

Solució B — afegir sal a la clau (quan cal GroupByKey de debò, per exemple per conservar els elements):

import random
N_SAL = 50

visites_per_sku = (
    esdeveniments
    | "SalarClau"     >> beam.Map(
        lambda e: (f'{e["sku"]}#{random.randint(0, N_SAL - 1)}', 1))
    | "ParcialSalat"  >> beam.CombinePerKey(sum)
    | "TreureSal"     >> beam.Map(lambda kv: (kv[0].split("#")[0], kv[1]))
    | "TotalFinal"    >> beam.CombinePerKey(sum)
)

Diferència pràctica entre A i B. L'A és preferible sempre que es pugui: és més simple, més barata i no té paràmetres per ajustar. La B afegeix un shuffle extra i un paràmetre (N_SAL) que cal dimensionar —massa baix no reparteix, massa alt crea sobrecàrrega—, però és l'única via quan l'operació no és associativa (per exemple, si necessites la llista completa d'esdeveniments de la sessió per reconstruir un recorregut) o quan el biaix persisteix fins i tot amb agregació parcial.

Mesura complementària immediata: abans de tocar codi, baixar --max_num_workers a 5 per deixar de pagar 15 màquines ocioses mentre es prepara el desplegament, i fer servir --update per substituir el pipeline conservant l'estat en curs, sense perdre les dades en vol.

Conclusió

AlpinaShop ja té canonades. En aquesta lliçó has vist per què un script en una VM és una solució que funciona fins que deixa de funcionar, i què aporta exactament un servei gestionat: paral·lelisme automàtic, reintents, escalat durant l'execució, semàntica d'escriptura fiable i observabilitat de sèrie.

Has après el model d'Apache Beam —Pipeline, PCollection immutable, PTransform, runner— i la seva idea central: un lot és un flux acotat i un stream és un flux no acotat, així que la mateixa lògica serveix per a tots dos. Coneixes les transformacions que cobreixen gairebé tots els casos i, sobretot, saps per què CombinePerKey és gairebé sempre millor que GroupByKey, i quan un ParDo amb classe DoFn guanya un Map.

Has escrit el pipeline per lots d'AlpinaShop línia a línia: llegeix les exportacions CSV d'alpinashop-catalogo, analitza, valida amb regles de negoci reals —imports amb coma, quantitats negatives, dates trencades—, escriu el detall net a lineas_pedido, calcula l'agregat de vendes per SKU amb agregació parcial, i envia tot el que és defectuós a una quarantena al bucket amb el seu motiu, sense tombar el procés. L'has provat en local amb el DirectRunner, que és lent i estricte expressament, amb un fitxer brut a consciència, i després l'has llançat a Dataflow amb el compte de servei sa-dataflow-pedidos, dins de sn-datos-euw1, sense IP pública, amb sostre de treballadors i etiquetes de facturació.

Has entrat al territori del temps, que és el que de debò distingeix el processament de dades seriós: temps de l'esdeveniment davant de temps de procés, amb el client del metro com a exemple de per què agrupar pel segon produeix números falsos; finestres fixes, lliscants i de sessió; marques d'aigua com a estimació de "ja no espero més"; activadors per emetre resultats parcials sense renunciar a la correcció posterior; i allowed_lateness com a decisió de negoci explícita sobre quant s'espera els endarrerits. I has deixat escrit el pipeline en temps real que consumirà pedidos-nuevos així que existeixi, amb timestamp_attribute —la línia que fa que tot l'anterior sigui cert— i llegint d'una subscripció, no d'un topic.

Has vist que per al que és comú hi ha drecera: les plantilles de Google resolen "Pub/Sub a BigQuery" sense escriure codi, i les plantilles flexibles empaqueten el teu propi pipeline com a imatge a Artifact Registry perquè un altre sistema l'invoqui. Coneixes l'escalat automàtic i els seus límits, Dataflow Prime i quan compensa, i saps llegir el graf, el system lag i la frescor de la dada per diagnosticar. Saps reconèixer el biaix de dades —l'escalat que no millora res— i corregir-lo amb agregació parcial, sal a la clau o entrades laterals. I coneixes la factura: per vCPU, memòria, disc i shuffle, amb l'advertiment gros que un pipeline en temps real està encès sempre i costa desenes d'euros al mes encara que no passi res.

Queda alguna cosa pendent que has notat al llarg de tota la lliçó. El pipeline per lots funciona perquè les dades ja són al bucket, i el de temps real funciona perquè algú publica a pedidos-nuevos. Però aquest topic encara no existeix, i no hem parlat de què passa amb la resta de la casa: quan entra una comanda, el magatzem l'ha de preparar, facturació ha d'emetre la factura, el client ha de rebre el seu correu de confirmació i l'analítica se n'ha d'assabentar. Avui l'aplicació Flask hauria de cridar els quatre, un darrere l'altre, i quedar-se penjada si el quart no respon.

Abans de resoldre-ho, però, hi ha una peça de l'ecosistema de dades que convé conèixer, perquè molta gent arriba a Google Cloud amb ella ja posada. A 04-03, Cloud Dataproc, veurem Spark i Hadoop gestionats: què són, per què continuen important després de vint anys, i quin és el patró que converteix un clúster car i permanent en un d'efímer que viu quatre minuts, fa la seva feina sobre les dades del bucket i s'autodestrueix. Crearem alpinashop-spark, executarem un treball PySpark que calcula quins productes es compren junts —la base del futur recomanador del mòdul 5— i compararem amb honestedat quan convé Dataproc, quan Dataflow i quan no cal cap dels dos.

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