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
- Què és un pipeline de dades i per què un script no n'hi ha prou
- Apache Beam: el model unificat
- Els quatre conceptes:
Pipeline,PCollection,PTransform, runner - Les transformacions que faràs servir el 90 % del temps
- Primer pipeline per lots, explicat línia a línia
- Execució local amb
DirectRunner - Execució gestionada amb
DataflowRunner - El temps en streaming: esdeveniment davant de procés
- Marques d'aigua, finestres, activadors i dades tardanes
- El pipeline en temps real de
pedidos-nuevos - Plantilles de Dataflow: la via pràctica
- Escalat automàtic i Dataflow Prime
- Monitoratge, paral·lelisme i biaix de dades
- Cost: què es paga exactament
- Quan Dataflow no és la resposta
- 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.
- 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.
- Els quatre conceptes:
Pipeline, PCollection, PTransform, runner
Pipeline, PCollection, PTransform, runnerPipeline é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:
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.
- 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.
- 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:
yielden comptes dereturn. UnDoFnés un generador: pot emetre zero, un o molts elements per entrada. Unreturnamb valor no funciona com esperes.TaggedOutputmarca 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 exporta89,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:
with beam.Pipeline(...) as p: en sortir del blocwith, Beam cridarun()i espera. Sense elwith, cal cridarp.run().wait_until_finish()explícitament.WRITE_APPENDdavant deWRITE_TRUNCATE: el detall s'afegeix (volem històric acumulat); l'agregat es reescriu sencer (volem la foto vigent). Triar malament aquí duplica dades silenciosament.additional_bq_parameterscrea 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.beam.Flatten()uneix diversesPCollectiondel mateix tipus en una. És la unió de branques del graf, l'equivalent a unUNION ALL.save_main_session=Trueserialitza l'àmbit global del mòdul per als treballadors. Sense això, un pipeline que funciona en local falla a Dataflow ambNameErrorsobre les constants o els imports. És l'error de novell número u.
- Execució local amb
DirectRunner
DirectRunnerAbans 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 DirectRunnerAquest 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,90hi 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.
- Execució gestionada amb
DataflowRunner
DataflowRunnerAmb 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=analiticaOpció 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.
- 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.
- 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.
- El pipeline en temps real de
pedidos-nuevos
pedidos-nuevosA 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.
- 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_erroresFixa'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-14La 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é.
- 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 streamingPosar 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.
- 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.
- 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=30n'hi ha prou; el valor per defecte és molt més gran i es paga per hora. - Posa
--max_num_workerssempre. És el fre de mà. - Fes servir VM Spot en lot tolerant a interrupcions:
--flexrs_goal=COST_OPTIMIZEDendarrereix 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.
- 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
puntuacionde les quals no sigui un enter entre 1 i 5, o elpaisde les quals no sigui un codi de dues lletres; - normalitzi el
skua majúscules i retalli eltextoa 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_APPENDi que el tauler consulti només l'última emissió per finestra i clau, per exemple ambQUALIFY 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 (
MERGEposterior sobrehora_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:
- 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.
- 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.
- La campanya de
MOCH-40L-AZ. En unGroupByKeypersku, 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
- Què és Google Cloud Platform?
- Configuració del teu compte de GCP
- Descripció general de la consola de GCP
- Projectes, jerarquia de recursos i facturació
- Regions, zones i model de responsabilitat compartida
- Cloud Shell i la CLI de gcloud
Mòdul 2: Serveis principals de GCP
- Compute Engine: màquines virtuals a Google Cloud
- Cloud Storage: emmagatzematge d'objectes
- Cloud SQL: bases de dades relacionals gestionades
- App Engine: plataforma com a servei
- Google Kubernetes Engine (GKE)
- Bases de dades NoSQL: Firestore, Bigtable i Spanner
- Com triar el servei de còmput adequat
Mòdul 3: Xarxes i seguretat
- Xarxes VPC
- Balanceig de càrrega al núvol
- Cloud CDN
- Gestió d'identitat i accés (IAM)
- Cloud Armor
- Secrets i xifratge: Secret Manager i Cloud KMS
- Cloud DNS, certificats TLS i publicació segura de serveis
Mòdul 4: Dades i anàlisi
- BigQuery: el magatzem de dades analític
- Cloud Dataflow: processament de dades per lots i en temps real
- Cloud Dataproc: Spark i Hadoop gestionats
- Cloud Pub/Sub: missatgeria asíncrona
- Cloud Data Fusion: integració de dades sense codi
- Orquestració de pipelines amb Cloud Composer i Workflows
- Govern de les dades i taulers amb Dataplex i Looker Studio
Mòdul 5: Aprenentatge automàtic i IA
- Vertex AI: la plataforma d'aprenentatge automàtic de GCP
- AutoML: models a mida sense escriure codi
- TensorFlow a GCP: entrenament i servei de models
- API de llenguatge natural
- API de visió
- IA generativa a Vertex AI: models Gemini i incrustacions
- MLOps: del model al producte amb Vertex AI Pipelines
Mòdul 6: DevOps i monitoratge
- Cloud Build: integració contínua a GCP
- Cloud Source Repositories i gestió del codi font
- Cloud Functions: funcions sense servidor
- Cloud Monitoring (abans Stackdriver): mètriques, taulers i alertes
- Cloud Deployment Manager i infraestructura com a codi nativa
- Cloud Logging i Cloud Trace: registres, traces i diagnòstic
- Terraform a GCP: infraestructura com a codi a la pràctica
Mòdul 7: Temes avançats de GCP
- Híbrid i multinúvol amb Anthos
- Computació sense servidor amb Cloud Run
- Xarxes avançades: VPC compartida, aparellament i connectivitat híbrida
- Bones pràctiques de seguretat
- Gestió i optimització de costos
- Fiabilitat: SLO, alta disponibilitat i recuperació de desastres
- Govern a escala: organització, polítiques i auditoria
