La Lucía té una pregunta que no es respon bé amb SQL: quins productes es compren junts. Si un client posa uns grampons a la cistella, quina probabilitat hi ha que també compri un piolet? Amb aquesta matriu, AlpinaShop podria suggerir complements a la fitxa de producte, agrupar articles en packs i col·locar millor les categories. És un càlcul sobre tots els parells possibles dins de cada comanda, repetit sobre dos anys d'història, i és el tipus de problema que comença sent elegant en SQL i acaba sent un JOIN d'una taula amb si mateixa que ningú no vol mantenir.
És, a més, exactament el tipus de problema per al qual es va inventar Spark: càlcul iteratiu i algorítmic sobre grans volums, expressat en codi, no en consultes.
Aquí convé ser honest sobre per què existeix aquesta lliçó. Si AlpinaShop comencés avui des de zero, probablement resoldria això amb BigQuery i Dataflow i no tocaria Spark mai. Però el món real no comença de zero: hi ha desenes de milers d'empreses amb clústers de Hadoop als seus centres de dades, amb anys de codi Spark en producció, amb equips que saben PySpark i no saben Beam, i amb biblioteques —MLlib, GraphX, tot l'ecosistema de Python científic— que no tenen equivalent directe. Cloud Dataproc és la resposta de Google a aquesta realitat: Hadoop i Spark gestionats, amb la particularitat que gira del revés la idea mateixa de clúster.
En aquesta lliçó veuràs què són Hadoop i Spark en l'essencial, crearàs el clúster alpinashop-spark, entendràs el patró que fa de Dataproc una cosa diferent d'un Hadoop llogat —el clúster efímer sobre emmagatzematge a Cloud Storage—, executaràs el treball PySpark que respon a la pregunta de la Lucía, i compararàs amb criteri Dataproc, Dataproc Serverless i Dataflow.
Contingut
- Hadoop i Spark en una pàgina
- Per què continuen important el 2026
- Què aporta Dataproc
- Anatomia d'un clúster de Dataproc
- Crear
alpinashop-spark - El patró clau: clúster efímer sobre Cloud Storage
- Enviar treballs:
gcloud dataproc jobs submit - El treball PySpark d'AlpinaShop: cistella mitjana i productes comprats junts
- Accions d'inicialització i versions d'imatge
- Dataproc Serverless per a Spark
- Dataproc, Serverless i Dataflow: la taula de decisió
- Notebooks i Spark SQL sobre BigQuery
- Migrar un Hadoop on-premise a Google Cloud
- Cost i VM Spot
- Hadoop i Spark en una pàgina
Hadoop va néixer el 2006 per resoldre un problema concret: processar més dades de les que cabien en una màquina, fent servir molts ordinadors barats que fallen sovint. Té tres peces:
- HDFS, un sistema de fitxers distribuït que parteix els fitxers en blocs i els replica (per defecte tres vegades) pels discos de les màquines del clúster.
- YARN, el gestor de recursos que decideix quin procés corre a quina màquina.
- MapReduce, el model de programació original: parteixes la feina en una fase map (transformar cada registre) i una altra reduce (agregar per clau), i el framework s'encarrega del repartiment i de les fallades.
MapReduce funcionava i era lentíssim, per una raó de disseny: escrivia al disc entre cada fase. Un algorisme iteratiu que necessita vint passades sobre les mateixes dades feia vint rondes d'escriptura i lectura al disc.
Spark va aparèixer el 2014 amb la correcció òbvia: mantenir les dades intermèdies en memòria. Per al mateix algorisme iteratiu, la millora era d'un o dos ordres de magnitud. I va afegir una API molt més agradable.
Spark s'organitza avui al voltant de:
| Component | Què és | Ús |
|---|---|---|
| Spark Core / RDD | L'abstracció original: col·lecció distribuïda i resilient | Control fi, operacions no expressables en SQL |
| DataFrame / Spark SQL | Taules amb esquema i un optimitzador (Catalyst) | El 90 % de l'ús actual; SQL sobre dades distribuïdes |
| MLlib | Biblioteca de machine learning distribuït | Agrupació en clústers, recomanació, classificació a escala |
| Structured Streaming | Processament continu amb l'API de DataFrame | Alternativa a Beam dins del món Spark |
| GraphX / GraphFrames | Algorismes sobre grafs | Xarxes, camins, comunitats |
L'arquitectura d'execució, que cal tenir al cap per entendre el que ve:
flowchart TD
D["Driver<br/>el teu programa PySpark<br/>construeix el pla"]
CM["Gestor de recursos (YARN)"]
E1["Executor 1<br/>tasques + memoria cau"]
E2["Executor 2"]
E3["Executor N"]
S["Emmagatzematge<br/>HDFS o Cloud Storage"]
D -->|demana recursos| CM
CM -->|assigna| E1 & E2 & E3
D -->|envia tasques| E1 & E2 & E3
E1 & E2 & E3 <--> S
El driver executa el teu codi, construeix un graf d'operacions i el trosseja en tasques. Els executors executen aquestes tasques sobre particions de les dades. Igual que a Beam, les transformacions són mandroses: no passa res fins que crides una acció (count(), collect(), write()). Aquesta mandra és el que permet a l'optimitzador reordenar i fusionar operacions.
- Per què continuen important el 2026
Amb BigQuery i Dataflow disponibles, per què aprendre això? Quatre raons concretes i una conseqüència.
Codi existent. Una empresa que migra al núvol amb 200.000 línies de PySpark en producció no les reescriurà. Reescriure no aporta valor de negoci, introdueix errors i consumeix mesos. Dataproc permet moure aquestes càrregues sense tocar el codi, i ja s'optimitzarà després.
Persones. El mercat té molts més enginyers que saben Spark que enginyers que saben Beam. Si l'equip de dades d'AlpinaShop contracta demà, és més probable que el candidat porti Spark. Triar la tecnologia que el teu equip sap fer servir és una decisió d'arquitectura legítima, no una concessió.
MLlib i l'ecosistema de Python. Algorismes com ALS per a recomanació, k-means a escala o FP-Growth per a regles d'associació estan implementats, provats i distribuïts. Escriure'ls en Beam seria absurd. I dins d'un job de Spark pots fer servir pandas, NumPy o scikit-learn sobre particions concretes.
Formats i ecosistema oberts. Hive, Presto/Trino, HBase, Kafka, Iceberg, Delta Lake, Hudi: un món sencer d'eines obertes que parla l'idioma de Hadoop. Si AlpinaShop volgués marxar de Google Cloud, un llac de dades en Parquet sobre emmagatzematge d'objectes i processat amb Spark viatja a qualsevol proveïdor sense canvis.
La conseqüència: Dataproc no és un servei de segona ni una relíquia. És la via de migració menys traumàtica i l'eina correcta quan el problema és algorítmic i l'equip sap Spark.
- Què aporta Dataproc
Muntar un clúster de Hadoop a mà —instal·lar, configurar YARN, dimensionar HDFS, ajustar la memòria dels executors, integrar l'autenticació— és una feina de setmanes i una font permanent de manteniment. Dataproc ho redueix a una comanda.
| Aspecte | Hadoop autogestionat | Dataproc |
|---|---|---|
| Temps de creació | Dies o setmanes | Menys de 2 minuts |
| Configuració | Manual, per component | Preconfigurada i coherent |
| Escalat | Comprar i muntar maquinari | Canviar un número; escalat automàtic disponible |
| Actualitzacions | Projecte en si mateix | Versions d'imatge gestionades |
| Cost en repòs | El maquinari, sempre | Zero si el clúster no existeix |
| Integració amb el núvol | Cal construir-la | Cloud Storage, BigQuery, Logging, IAM de sèrie |
Els dos minuts de creació no són una dada de màrqueting: són el que canvia el model mental. Si crear un clúster costa setmanes, el clúster és una instal·lació permanent que cal cuidar, compartir entre equips i mantenir sempre encesa per si de cas. Si costa noranta segons, el clúster passa a ser d'un sol ús: es crea per a un treball, es destrueix en acabar, i cada treball pot tenir el seu amb la versió i les biblioteques que necessita.
Això és el que veurem a l'apartat 6, i és la idea més important de la lliçó.
- Anatomia d'un clúster de Dataproc
Un clúster de Dataproc té tres tipus de node:
| Node | Funció | Quantitat | Notes |
|---|---|---|---|
| Mestre | Executa el driver, YARN ResourceManager, HDFS NameNode | 1 (o 3 en alta disponibilitat) | Si cau amb 1 node, el treball mor |
| Workers primaris | Executen tasques; aporten disc a HDFS | Mínim 2 (o 0 en mode node únic) | VM estàndard, estables |
| Workers secundaris | Només càlcul; no emmagatzemen HDFS | 0 a N | Poden ser Spot: fins a ~80 % més barats |
Els workers secundaris mereixen atenció. Com que no participen a HDFS, poden desaparèixer sense que es perdi cap dada. Per això poden ser VM Spot (les mateixes de 02-01, que Google pot reclamar amb 30 segons d'avís). Si una desapareix a mitja tasca, YARN reassigna aquella tasca a un altre node i el treball continua, més lent però correcte.
Aquesta és la combinació guanyadora per a AlpinaShop: pocs workers primaris estàndard per donar estabilitat, i molts secundaris Spot per donar potència barata.
flowchart TD
subgraph C["Cluster alpinashop-spark"]
M["Mestre<br/>n2-standard-4<br/>driver + YARN RM"]
W1["Worker primari 1<br/>n2-standard-4"]
W2["Worker primari 2<br/>n2-standard-4"]
S1["Secundari Spot 1"]
S2["Secundari Spot 2"]
S3["Secundari Spot N<br/>escalat automatic"]
end
GCS["Cloud Storage<br/>gs://alpinashop-catalogo<br/>gs://alpinashop-datalake"]
BQ["BigQuery<br/>alpinashop_analitica"]
M --- W1 & W2 & S1 & S2 & S3
W1 & W2 & S1 & S2 & S3 <--> GCS
W1 & W2 <--> BQ
Fixa't que l'emmagatzematge és fora del clúster. Aquest és el pròxim apartat.
- Crear
alpinashop-spark
alpinashop-sparkgcloud config set project alpinashop-datos
# Bucket propi per al llac de dades i els artefactes de Spark
gcloud storage buckets create gs://alpinashop-datalake \
--project=alpinashop-datos --location=europe-west1 \
--uniform-bucket-level-access
# Compte de servei dedicat, amb minim privilegi
gcloud iam service-accounts create sa-dataproc-analitica \
--display-name="Clusters de Dataproc d analitica"
SA="[email protected]"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/dataproc.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-datalake \
--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 clúster:
gcloud dataproc clusters create alpinashop-spark \
--region=europe-west1 \
--zone=europe-west1-b \
--service-account="$SA" \
--subnet=sn-datos-euw1 \
--no-address \
--master-machine-type=n2-standard-4 \
--master-boot-disk-size=100GB \
--master-boot-disk-type=pd-balanced \
--num-workers=2 \
--worker-machine-type=n2-standard-4 \
--worker-boot-disk-size=200GB \
--num-secondary-workers=2 \
--secondary-worker-type=spot \
--image-version=2.2-debian12 \
--optional-components=JUPYTER \
--enable-component-gateway \
--bucket=alpinashop-datalake \
--max-idle=30m \
--properties="spark:spark.sql.adaptive.enabled=true,\
spark:spark.dynamicAllocation.enabled=true,\
spark:spark.sql.sources.partitionOverwriteMode=dynamic" \
--labels=entorno=produccion,equipo=datos,centro-coste=analitica,aplicacion=analiticaRepàs del que importa:
| Opció | Per què hi és |
|---|---|
--subnet=sn-datos-euw1 + --no-address |
El clúster viu a alpinashop-vpc, sense IP pública. Surt a internet pel Cloud NAT de 03-01 i accedeix a les API per Private Google Access |
--service-account |
Identitat pròpia amb permisos mínims, no el compte per defecte de Compute |
--num-secondary-workers=2 --secondary-worker-type=spot |
Potència barata que pot desaparèixer sense trencar res |
--image-version=2.2-debian12 |
Versió fixada. Sense això, Google triaria la més recent i un treball que funcionava ahir podria fallar demà |
--optional-components=JUPYTER + --enable-component-gateway |
Notebooks accessibles des de la consola amb autenticació d'IAM, sense obrir ports |
--bucket=alpinashop-datalake |
Bucket de treball per a registres i fitxers temporals del clúster |
--max-idle=30m |
El clúster s'autodestrueix després de 30 minuts sense treballs. L'opció més rendible de tota la comanda |
spark.sql.adaptive.enabled |
Execució adaptativa: Spark reajusta particions i estratègies de JOIN en temps real. Mitiga bastant el biaix de dades |
spark.dynamicAllocation.enabled |
Spark demana i allibera executors segons necessita |
Afegir escalat automàtic al clúster requereix una política a part:
# politica-autoescalado.yaml
workerConfig:
minInstances: 2
maxInstances: 2 # els primaris NO escalen: aporten HDFS
secondaryWorkerConfig:
minInstances: 0
maxInstances: 20 # els Spot si que escalen, fins a 20
basicAlgorithm:
cooldownPeriod: 2m
yarnConfig:
scaleUpFactor: 1.0 # afegeix el 100 % del que YARN demana
scaleDownFactor: 0.5 # retira la meitat del que sobra: prudent
gracefulDecommissionTimeout: 10m # espera que acabin les tasques en cursgcloud dataproc autoscaling-policies import pol-autoescalado-analitica \
--region=europe-west1 --source=politica-autoescalado.yaml
gcloud dataproc clusters update alpinashop-spark --region=europe-west1 \
--autoscaling-policy=pol-autoescalado-analiticaEl gracefulDecommissionTimeout és important: sense ell, en reduir el clúster es maten nodes amb tasques en curs i cal refer-les. Amb ell, s'espera que acabin. I scaleDownFactor: 0.5 evita l'efecte acordió de pujar i baixar constantment.
Verificació:
gcloud dataproc clusters list --region=europe-west1 \
--format="table(clusterName, status.state, config.workerConfig.numInstances)"
gcloud dataproc clusters describe alpinashop-spark --region=europe-west1
- El patró clau: clúster efímer sobre Cloud Storage
Aquí hi ha la idea que ho canvia tot.
En un Hadoop tradicional, l'emmagatzematge i el càlcul són a les mateixes màquines. HDFS viu als discos dels workers. Això tenia una raó excel·lent el 2006: moure dades per la xarxa era caríssim comparat amb llegir-les del disc local, així que el principi era "porta el càlcul a la dada".
Però té una conseqüència devastadora: si apagues el clúster, perds les dades. Per això els clústers Hadoop estan sempre encesos, encara que només treballin tres hores al dia. Es paga el maquinari les vint-i-quatre.
A Google Cloud, la xarxa interna és tan ràpida que la premissa ja no s'aguanta: llegir de Cloud Storage no és significativament més lent que llegir d'un disc local. Això permet invertir el disseny:
flowchart LR
subgraph Abans["Hadoop tradicional"]
H["Cluster permanent<br/>calcul + HDFS<br/>ences 24x7"]
end
subgraph Ara["Patro Dataproc"]
GCS["Cloud Storage<br/>gs://alpinashop-datalake<br/>PERMANENT, barat"]
C1["Cluster efimer A<br/>Spark 3.5<br/>viu 20 min"]
C2["Cluster efimer B<br/>Spark 3.3 + biblioteca X<br/>viu 5 min"]
end
GCS <--> C1
GCS <--> C2
El connector de Cloud Storage ve preinstal·lat a Dataproc i fa que les rutes gs:// es comportin com a rutes d'HDFS per a qualsevol codi de Spark o Hadoop:
# El mateix codi que llegia d HDFS...
df = spark.read.parquet("hdfs:///datos/pedidos/")
# ...llegeix de Cloud Storage canviant el prefix. Res mes.
df = spark.read.parquet("gs://alpinashop-datalake/pedidos/")Els avantatges de separar emmagatzematge i càlcul són concrets:
- Pagues càlcul només quan calcules. Un clúster que viu 20 minuts al dia costa l'1,4 % d'un de permanent.
- Les dades sobreviuen al clúster. El pots destruir amb total tranquil·litat.
- Diversos clústers sobre les mateixes dades. El de la Lucía amb Spark 3.5 i el d'un proveïdor extern amb una versió antiga, simultàniament, sense interferir-se.
- Durabilitat molt superior. Cloud Storage replica amb garanties d'11 nous; HDFS amb tres còpies en tres discos del mateix rack no s'hi acosta.
- Les dades són accessibles des de fora de Spark. BigQuery les llegeix com a taula externa, Dataflow les processa, l'aplicació les baixa.
- Sense cerimònia d'actualització. Per passar a una versió nova de Spark, crees un clúster nou. No migres res.
Els inconvenients, per ser justos: latència una mica més gran en operacions de molts fitxers petits, i que Cloud Storage no té canvi de nom atòmic de directoris, cosa que afecta certs patrons d'escriptura. Es mitiga amb formats columnars i fitxers de mida raonable (128-512 MB), que és el que cal fer de totes maneres.
La regla d'AlpinaShop: HDFS només com a espai de treball temporal dins d'un treball. Res que hagi de sobreviure al clúster s'escriu a HDFS. I --max-idle a tots els clústers, sense excepció. Un clúster oblidat un cap de setmana costa més que tota la resta de l'analítica del mes.
Per a treballs programats, ni tan sols es manté un clúster: es crea, s'executa i es destrueix en un sol pas amb els workflow templates:
# 1) Plantilla de flux de treball
gcloud dataproc workflow-templates create wf-cesta-media --region=europe-west1
# 2) Cluster gestionat: neix i mor amb el flux
gcloud dataproc workflow-templates set-managed-cluster wf-cesta-media \
--region=europe-west1 \
--cluster-name=cluster-efimero-cesta \
--service-account="$SA" --subnet=sn-datos-euw1 --no-address \
--master-machine-type=n2-standard-4 \
--worker-machine-type=n2-standard-4 --num-workers=2 \
--num-secondary-workers=4 --secondary-worker-type=spot \
--image-version=2.2-debian12
# 3) El treball que s executara
gcloud dataproc workflow-templates add-job pyspark \
gs://alpinashop-datalake/jobs/cesta_media.py \
--step-id=cesta-media --workflow-template=wf-cesta-media \
--region=europe-west1 \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31
# 4) Executar: crea el cluster, llanca el job, destrueix el cluster
gcloud dataproc workflow-templates instantiate wf-cesta-media --region=europe-west1Aquesta comanda final és la que invocarà l'orquestrador de 04-06. Cost total: els minuts que duri el treball. Zero la resta del mes.
- Enviar treballs:
gcloud dataproc jobs submit
gcloud dataproc jobs submitDataproc admet diversos tipus de treball:
# PySpark
gcloud dataproc jobs submit pyspark gs://alpinashop-datalake/jobs/mi_job.py \
--cluster=alpinashop-spark --region=europe-west1 \
-- arg1 arg2
# Spark (JAR d Scala o Java)
gcloud dataproc jobs submit spark --cluster=alpinashop-spark --region=europe-west1 \
--class=com.alpinashop.Informe --jars=gs://alpinashop-datalake/jars/informes.jar
# Spark SQL des d un fitxer
gcloud dataproc jobs submit spark-sql --cluster=alpinashop-spark \
--region=europe-west1 --file=gs://alpinashop-datalake/sql/ventas.sql
# Hive, Pig, Presto/Trino tambe estan disponiblesTot el que va després de -- són arguments per al teu programa, no per a gcloud. És una confusió molt freqüent.
Opcions útils en enviar:
gcloud dataproc jobs submit pyspark gs://alpinashop-datalake/jobs/cesta_media.py \
--cluster=alpinashop-spark \
--region=europe-west1 \
--py-files=gs://alpinashop-datalake/jobs/utilidades.zip \
--files=gs://alpinashop-datalake/config/categorias.json \
--jars=gs://spark-lib/bigquery/spark-3.5-bigquery-0.42.0.jar \
--properties="spark.executor.memory=6g,spark.executor.cores=2,spark.sql.shuffle.partitions=200" \
--labels=proceso=cesta-media \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31 \
--salida=gs://alpinashop-datalake/resultados/cesta/--py-files: mòduls Python propis, empaquetats en.zipo.egg, distribuïts a tots els executors.--files: fitxers de dades que el treball necessita, accessibles per nom al directori de treball.--jars: dependències Java, com el connector de BigQuery.spark.sql.shuffle.partitions: nombre de particions després d'un shuffle. El valor per defecte (200) és un mal ajust per a gairebé tothom: massa per a dades petites, insuficient per a grans. Una regla raonable és 2-3 vegades el nombre total de cores del clúster.
Seguiment:
gcloud dataproc jobs list --region=europe-west1 --cluster=alpinashop-spark \
--format="table(reference.jobId, status.state, statusHistory[0].stateStartTime)"
gcloud dataproc jobs wait JOB_ID --region=europe-west1 # segueix els registres en viuEls registres van automàticament a Cloud Logging (06-06) i la interfície de Spark History Server queda accessible pel component gateway, fins i tot després de destruir el clúster si es configura un servidor d'historial persistent.
- El treball PySpark d'AlpinaShop: cistella mitjana i productes comprats junts
Ara el treball real. Llegeix les línies de comanda del llac, calcula la cistella mitjana per mes i país, i construeix la matriu de coocurrència de productes.
"""
cesta_media.py -- Analisi de cistella d AlpinaShop amb PySpark.
Entrada : gs://alpinashop-datalake/pedidos/ (Parquet, particionat per data)
Sortides: gs://alpinashop-datalake/resultados/cesta/ (Parquet)
alpinashop-datos.alpinashop_analitica.productos_juntos (BigQuery)
Enviament:
gcloud dataproc jobs submit pyspark gs://.../cesta_media.py \
--cluster=alpinashop-spark --region=europe-west1 \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31
"""
import argparse
from pyspark.sql import SparkSession, functions as F, Window
PROJECTE = "alpinashop-datos"
DATASET = "alpinashop_analitica"
LLAC = "gs://alpinashop-datalake"
def crear_sessio():
"""La SparkSession es el punt d entrada. A Dataproc, la configuracio
de recursos i el mestre els aporta YARN: no cal indicar-los."""
return (
SparkSession.builder
.appName("alpinashop-cesta-media")
# Bucket temporal que necessita el connector de BigQuery per escriure
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate()
)def carregar_linies(spark, des_de, fins_a):
"""Llegeix les linies de comanda del llac i les filtra per data.
Parquet desa l esquema i les estadistiques per bloc, aixi que Spark
pot saltar-se fitxers sencers: es el 'predicate pushdown'.
"""
linies = (
spark.read.parquet(f"{LLAC}/pedidos/lineas/")
.filter((F.col("fecha_pedido") >= des_de) & (F.col("fecha_pedido") <= fins_a))
# Seleccionar columnes aviat redueix memoria i shuffle
.select("pedido_id", "fecha_pedido", "sku", "cantidad",
"precio_unitario", "importe_linea")
)
capcaleres = (
spark.read.parquet(f"{LLAC}/pedidos/cabeceras/")
.filter((F.col("fecha_pedido") >= des_de) & (F.col("fecha_pedido") <= fins_a))
.filter(~F.col("estado").isin("cancelado", "devuelto"))
.select("pedido_id", "pais", "canal", "cliente_id")
)
# broadcast(): les capcaleres del periode caben a la memoria de cada executor.
# Evita el shuffle del JOIN, igual que l entrada lateral de Beam a 04-02.
return linies.join(F.broadcast(capcaleres), on="pedido_id", how="inner")
def calcular_cistella_mitjana(df):
"""Cistella mitjana per mes i pais: import total i nombre d articles."""
per_comanda = (
df.groupBy("pedido_id", "pais", F.trunc("fecha_pedido", "month").alias("mes"))
.agg(
F.sum("importe_linea").alias("importe_pedido"),
F.sum("cantidad").alias("articulos_pedido"),
F.countDistinct("sku").alias("skus_distintos"),
)
)
return (
per_comanda.groupBy("mes", "pais")
.agg(
F.count("*").alias("num_pedidos"),
F.round(F.avg("importe_pedido"), 2).alias("cesta_media_eur"),
F.round(F.expr("percentile_approx(importe_pedido, 0.5)"), 2)
.alias("cesta_mediana_eur"),
F.round(F.avg("articulos_pedido"), 2).alias("articulos_medios"),
F.round(F.avg("skus_distintos"), 2).alias("skus_medios"),
)
.orderBy("mes", "pais")
)F.broadcast() mereix atenció: li diu a Spark que repliqui el DataFrame petit a tots els executors en comptes de repartir tots dos costats per la xarxa. És la mateixa optimització que l'entrada lateral de Beam a 04-02, i és la diferència entre un JOIN de segons i un de minuts. Només funciona si el costat petit cap a la memòria (uns quants centenars de MB com a màxim).
def calcular_productes_junts(df, suport_minim=20):
"""Matriu de coocurrencia: quins parells de SKU apareixen a la mateixa comanda.
L algorisme es un self-join del conjunt de comandes amb si mateix,
amb dues precaucions fonamentals de rendiment.
"""
# 1) Una comanda pot tenir el mateix SKU en diverses linies: ens quedem
# amb parells (comanda, sku) unics per no comptar dues vegades.
comanda_sku = df.select("pedido_id", "sku").distinct()
# 2) PRECAUCIO CRITICA: descartar comandes amb massa linies.
# Una comanda de 200 SKU genera 200*199/2 = 19.900 parells ella sola,
# i aquestes comandes corporatives rares dominarien el calcul i la memoria.
mida = comanda_sku.groupBy("pedido_id").agg(F.count("*").alias("n_skus"))
comandes_valides = mida.filter((F.col("n_skus") >= 2) & (F.col("n_skus") <= 30))
base = comanda_sku.join(F.broadcast(comandes_valides.select("pedido_id")),
on="pedido_id", how="inner")
esq = base.withColumnRenamed("sku", "sku_a")
dre = base.withColumnRenamed("sku", "sku_b")
parells = (
esq.join(dre, on="pedido_id")
# 3) sku_a < sku_b elimina els parells amb si mateix I els duplicats
# invertits: (A,B) es compta una vegada, no dues com (A,B) i (B,A).
.filter(F.col("sku_a") < F.col("sku_b"))
.groupBy("sku_a", "sku_b")
.agg(F.count("*").alias("veces_juntos"))
.filter(F.col("veces_juntos") >= suport_minim)
)
# 4) Metriques de regles d associacio: suport, confianca i lift.
recompte_sku = (base.groupBy("sku").agg(F.count("*").alias("veces_total")))
total_comandes = base.select("pedido_id").distinct().count()
resultat = (
parells
.join(F.broadcast(recompte_sku.withColumnRenamed("sku", "sku_a")
.withColumnRenamed("veces_total", "total_a")),
on="sku_a")
.join(F.broadcast(recompte_sku.withColumnRenamed("sku", "sku_b")
.withColumnRenamed("veces_total", "total_b")),
on="sku_b")
.withColumn("soporte", F.col("veces_juntos") / F.lit(total_comandes))
.withColumn("confianza_a_b", F.col("veces_juntos") / F.col("total_a"))
.withColumn("confianza_b_a", F.col("veces_juntos") / F.col("total_b"))
.withColumn(
"lift",
(F.col("veces_juntos") * F.lit(total_comandes))
/ (F.col("total_a") * F.col("total_b")),
)
)
# 5) Top 5 acompanyants de cada producte, amb funcio de finestra
finestra = Window.partitionBy("sku_a").orderBy(F.desc("lift"))
return (
resultat
.withColumn("puesto", F.row_number().over(finestra))
.filter(F.col("puesto") <= 5)
.select("sku_a", "sku_b", "veces_juntos",
F.round("soporte", 5).alias("soporte"),
F.round("confianza_a_b", 4).alias("confianza"),
F.round("lift", 3).alias("lift"),
"puesto")
)Com s'interpreten les tres mètriques, perquè són les que la Lucía portarà a la reunió:
| Mètrica | Què significa | Exemple AlpinaShop |
|---|---|---|
| Suport | Proporció de comandes que contenen tots dos productes | 0,012 → l'1,2 % de les comandes porten grampons i piolet |
| Confiança | Si compra A, probabilitat que compri B | 0,34 → un terç de qui compra grampons compra piolet |
| Lift | Quant més probable és que vagin junts respecte de l'atzar | 8,5 → vuit vegades i mitja més del que s'esperaria: associació fortíssima |
El lift és el que cal mirar. La confiança enganya amb els productes supervendes: si el 60 % de les comandes porten mitjons tècnics, qualsevol producte tindrà alta confiança cap a ells sense que hi hagi relació real. El lift corregeix per la popularitat de cada producte. Lift més gran que 1 indica associació real; lift proper a 1, independència.
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--fecha-desde", required=True)
parser.add_argument("--fecha-hasta", required=True)
parser.add_argument("--salida", default=f"{LLAC}/resultados/cesta/")
args = parser.parse_args()
spark = crear_sessio()
spark.sparkContext.setLogLevel("WARN")
df = carregar_linies(spark, args.fecha_desde, args.fecha_hasta)
# cache(): el DataFrame es fa servir en DOS calculs diferents. Sense cache,
# Spark rellegiria i refiltraria tot el llac dues vegades.
df.cache()
cistella = calcular_cistella_mitjana(df)
(cistella.coalesce(1) # un sol fitxer: el resultat es petit
.write.mode("overwrite")
.parquet(f"{args.salida}/cesta_media"))
junts = calcular_productes_junts(df)
(junts.write.format("bigquery")
.option("table", f"{PROJECTE}.{DATASET}.productos_juntos")
.option("writeMethod", "direct")
.mode("overwrite")
.save())
df.unpersist()
print(f"Cistella mitjana: {cistella.count()} files")
print(f"Parells de productes: {junts.count()} files")
spark.stop()
if __name__ == "__main__":
main()Tres detalls de rendiment que cal interioritzar:
cache()materialitza el DataFrame a la memòria dels executors. Sense ell, com que Spark és mandrós, cada acció posterior recalcularia tota la cadena des de la lectura del llac. Amb dos consumidors, estalvia la meitat de la feina. Iunpersist()en acabar, per alliberar memòria.coalesce(1)redueix a una sola partició abans d'escriure. És correcte només per a resultats petits; amb dades grans, concentrar-ho tot en un executor el tombaria. Per a volums grans es fa servirrepartition(n).writeMethod=directfa servir la Storage Write API de BigQuery en comptes de passar per fitxers temporals al bucket. És més ràpid i evita gestionar-ne la neteja.
I l'advertiment que correspon: aquesta anàlisi fa servir cliente_id i dades de comanda. Si en algun moment s'hi incorporen dades personals identificables, el tractament s'ha d'ajustar al RGPD i l'ha de revisar un professional de compliment normatiu. Per a l'anàlisi de cistella no cal saber qui és el client, només què hi havia a la comanda; mantenir-ho així és minimització de dades per disseny. Totes les dades d'aquest curs són fictícies.
- Accions d'inicialització i versions d'imatge
Una acció d'inicialització és un script que s'executa a cada node en crear-se el clúster. Serveix per instal·lar dependències que no vénen a la imatge.
# Script propi al bucket
cat > init-alpinashop.sh <<'EOF'
#!/bin/bash
set -euxo pipefail
# Biblioteques de Python per a l analisi de cistella
pip install --no-cache-dir mlxtend==0.23.1 pyarrow==16.1.0
# Nomes al mestre: utilitats de diagnostic
ROL=$(/usr/share/google/get_metadata_value attributes/dataproc-role)
if [[ "$ROL" == "Master" ]]; then
pip install --no-cache-dir jupyterlab-git
fi
EOF
gcloud storage cp init-alpinashop.sh gs://alpinashop-datalake/init/
gcloud dataproc clusters create alpinashop-spark-ml \
--region=europe-west1 --subnet=sn-datos-euw1 --no-address \
--service-account="$SA" \
--initialization-actions=gs://alpinashop-datalake/init/init-alpinashop.sh \
--initialization-action-timeout=10m \
--image-version=2.2-debian12 \
--max-idle=30mConsells sobre les accions d'inicialització:
- L'script s'executa a tots els nodes, inclosos els que afegeixi l'escalat automàtic. Ha de ser idempotent i ràpid: un script de cinc minuts multiplica per cinc el temps d'arrencada de cada node nou.
- Fes servir
get_metadata_value attributes/dataproc-roleper distingir mestre de worker. set -euxo pipefailfa que l'script falli sorollosament en comptes de deixar el node a mitges.- Si les dependències són moltes, és millor construir una imatge personalitzada que instal·lar-les a cada arrencada.
Sobre les versions d'imatge: cada versió de Dataproc empaqueta un conjunt concret de Spark, Hadoop, Python i el sistema operatiu. La sèrie 2.2-debian12, per exemple, porta Spark 3.5 i Python 3.11.
| Pràctica | Conseqüència |
|---|---|
| No indicar versió | Google tria la més recent: el teu treball es pot trencar tot sol |
Indicar la sèrie (2.2-debian12) |
Actualitzacions menors automàtiques dins de la sèrie. Equilibri raonable |
Indicar la versió exacta (2.2.28-debian12) |
Reproductibilitat total. Recomanat en producció crítica |
Per a AlpinaShop: sèrie fixada en desenvolupament, versió exacta als fluxos de treball programats. I al fitxer del flux, versionat a Git.
- Dataproc Serverless per a Spark
Encara que un clúster efímer és molt millor que un de permanent, continua calent dimensionar-lo: quants workers, quina màquina, quanta memòria per executor. Dataproc Serverless elimina aquesta decisió: envies el treball de Spark i Google s'encarrega de tot.
gcloud dataproc batches submit pyspark \
gs://alpinashop-datalake/jobs/cesta_media.py \
--batch=cesta-media-$(date +%Y%m%d-%H%M%S) \
--region=europe-west1 \
--version=2.2 \
--service-account="$SA" \
--subnet=sn-datos-euw1 \
--deps-bucket=gs://alpinashop-datalake \
--properties="spark.executor.instances=4,\
spark.dynamicAllocation.enabled=true,\
spark.dynamicAllocation.maxExecutors=20" \
--labels=entorno=produccion,equipo=datos,centro-coste=analitica \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31No hi ha clusters create. No hi ha --max-idle perquè no hi ha res per apagar. No hi ha --num-workers.
| Aspecte | Dataproc amb clúster | Dataproc Serverless |
|---|---|---|
| Gestió | Crees i destrueixes clústers | Cap |
| Arrencada | 90-120 s | 30-60 s |
| Dimensionament | El tries tu | Automàtic |
| Facturació | Per VM i hora | Per DCU (unitats de càlcul de dades) mentre dura |
| Components | Tot l'ecosistema: Hive, HBase, Presto, Jupyter | Només Spark |
| Accions d'inicialització | Sí | No; es fan servir imatges de contenidor personalitzades |
| VM Spot | Sí, molt barat | No aplicable |
| Sessió interactiva | Notebook al clúster | Sessions interactives Serverless |
El criteri: si el teu treball és Spark pur i no necessites Hive ni HBase ni un clúster de llarga vida, comença per Serverless. És menys per administrar i menys per oblidar-se encès. Fes servir clúster quan necessitis components de l'ecosistema, sessions interactives llargues, o quan el descompte de les VM Spot en càrregues molt grans compensi la gestió.
Per a AlpinaShop, la decisió raonada és: l'anàlisi de cistella va a Serverless, perquè és Spark pur, s'executa mensualment i ningú no es vol recordar d'apagar res. El clúster alpinashop-spark es manté únicament com a entorn d'exploració amb Jupyter, amb --max-idle=30m, i es destrueix quan no es faci servir durant un mes.
- Dataproc, Serverless i Dataflow: la taula de decisió
| Criteri | Dataproc (clúster) | Dataproc Serverless | Dataflow |
|---|---|---|---|
| Model de programació | Spark / Hadoop / Hive | Spark | Apache Beam |
| Lot | Sí | Sí | Sí |
| Streaming | Structured Streaming | Limitat | El seu punt fort |
| Arrencada | ~2 min | ~40 s | ~2 min (lot) |
| Infraestructura | La defineixes tu | Cap | Cap |
| Escalat automàtic | Amb política | Automàtic | Automàtic |
| Cost en repòs | El clúster, si el deixes | Zero | Zero (llevat de streaming actiu) |
| Semàntica exactament una vegada | Cal construir-la | Igual | De sèrie |
| Finestres i temps de l'esdeveniment | Manual | Manual | Model complet |
| Ecosistema de biblioteques | Enorme (MLlib, pandas…) | Gran | Limitat a Beam |
| Portabilitat | Alta (Spark corre a tot arreu) | Mitjana | Alta (Beam té diversos runners) |
| Triar si… | Tens codi Spark, necessites Hive/HBase, l'equip sap Spark | Spark pur sense voler gestionar res | Streaming, temps de l'esdeveniment, pipelines nous |
La decisió final per a AlpinaShop queda així:
- Streaming de comandes i visites → Dataflow. El model de finestres i marques d'aigua de Beam no té equivalent còmode a Spark, i les plantilles de Pub/Sub a BigQuery resolen el cas base sense codi.
- Càrregues i transformacions per lots noves → Dataflow, per coherència amb l'anterior i perquè l'equip ja ho té muntat.
- Anàlisi algorítmica: cistella, coocurrència, futurs models amb MLlib → Dataproc Serverless.
- Exploració interactiva sobre el llac → Notebook a
alpinashop-spark, amb--max-idle. - Transformació expressable en SQL sobre dades ja a BigQuery → BigQuery, sense moure res.
- Notebooks i Spark SQL sobre BigQuery
Amb --optional-components=JUPYTER i --enable-component-gateway, la consola de Dataproc mostra un enllaç a JupyterLab protegit per IAM. Res d'obrir ports ni túnels SSH: qui tingui el rol adequat entra, i qui no, no.
El connector de BigQuery per a Spark permet llegir i escriure taules directament:
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder
.appName("exploracion-lucia")
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate())
# Llegir una taula completa: el connector fa servir la Storage Read API,
# que llegeix en paral·lel i en format columnar. No exporta a fitxers.
productes = (spark.read.format("bigquery")
.option("table", "alpinashop-datos.alpinashop_analitica.productos")
.load())
# MILLOR: delegar el filtre a BigQuery per portar menys dades
linies = (spark.read.format("bigquery")
.option("table", "alpinashop-datos.alpinashop_analitica.lineas_pedido")
.option("filter", "fecha_pedido >= '2026-01-01'") # s executa a BigQuery
.load())
# A partir d aqui, Spark SQL normal
linies.createOrReplaceTempView("lineas")
productes.createOrReplaceTempView("productos")
resum = spark.sql("""
SELECT p.categoria,
COUNT(DISTINCT l.pedido_id) AS pedidos,
ROUND(SUM(l.importe_linea), 2) AS ventas_eur
FROM lineas l
JOIN productos p ON p.sku = l.sku
GROUP BY p.categoria
ORDER BY ventas_eur DESC
""")
resum.show(truncate=False)L'opció filter és la que marca la diferència: s'executa a BigQuery, no a Spark. Sense ella, el connector portaria la taula sencera per la xarxa perquè Spark la filtrés, pagant la lectura completa a BigQuery i perdent temps. És el mateix principi que vam veure amb EXTERNAL_QUERY a 04-01: filtra el més a prop possible de l'origen.
I la pregunta obligada: si pots fer això amb Spark SQL, per què no fer-ho a BigQuery directament? Gairebé sempre ho hauries de fer. El connector té sentit quan el resultat del SQL alimenta un algorisme de MLlib, quan creues taules de BigQuery amb fitxers del llac que no estan carregats, o quan el codi Spark ja existeix. Per a una agregació pura, BigQuery és més ràpid i més barat.
- Migrar un Hadoop on-premise a Google Cloud
Aquest és l'escenari per al qual Dataproc està més justificat. Suposa que AlpinaShop absorbeix un competidor amb un clúster Hadoop de 30 nodes.
Què es conserva:
- El codi Spark, Hive i PySpark: en la seva immensa majoria funciona sense canvis.
- Les consultes de Hive: Dataproc inclou Hive, i el metastore es pot migrar.
- Els formats de dades: Parquet, ORC i Avro són idèntics.
- Els fluxos d'Oozie, encara que convé reemplaçar-los per Composer (04-06).
Què es replanteja, obligatòriament:
| Element on-premise | A Google Cloud | Per què canvia |
|---|---|---|
| HDFS permanent | Cloud Storage | És el canvi fonamental: sense ell no hi ha clústers efímers ni estalvi |
| Un clúster gegant compartit | Diversos clústers petits per càrrega | Cada equip amb la seva versió i el seu pressupost; sense cues ni veïns sorollosos |
| Dimensionat per al pic | Escalat automàtic + Spot | Es pagava el pic les 24 hores |
| Kerberos | IAM i comptes de servei | Model d'identitat del núvol (03-04) |
| Metastore de Hive local | Dataproc Metastore gestionat | Sobreviu als clústers efímers: és la peça que ho fa possible |
| Oozie / cron | Cloud Composer o Workflows | 04-06 |
| Impala / Presto per a consultes | BigQuery | Sol ser el salt més gran de rendiment i de simplicitat |
| Flume / Kafka d'ingesta | Pub/Sub (04-04) o Managed Kafka | Gestionat |
L'estratègia recomanada, i l'única que sol sortir bé, és per fases:
- Copiar les dades a Cloud Storage amb el Storage Transfer Service o
hadoop distcp, sense tocar res més. El clúster on-premise continua funcionant. - Aixecar Dataproc Metastore i registrar les taules apuntant a
gs://en comptes dehdfs://. - Executar els treballs existents a Dataproc contra les dades ja a Cloud Storage, comparant resultats amb els del clúster antic. Aquesta fase de doble execució és innegociable: és l'única manera de demostrar que els números coincideixen.
- Apagar el clúster on-premise quan la comparació quadri durant diverses setmanes.
- Només llavors, optimitzar: passar les consultes de Hive a BigQuery, els treballs de streaming a Dataflow, adoptar Serverless.
L'error clàssic és intentar el pas 5 alhora que el 3, és a dir, migrar i modernitzar simultàniament. Quan els números no quadren, no se sap si és per la migració o per la reescriptura, i el projecte s'encalla durant mesos.
- Cost i VM Spot
Dataproc factura dues coses:
- Les VM subjacents (Compute Engine, discos i xarxa), a tarifa normal.
- Una tarifa de gestió de Dataproc, de l'ordre de 0,01 $ per vCPU i hora, verificable a la documentació oficial.
És a dir, el sobrecost de Dataproc respecte de muntar Hadoop tu mateix en VM és petit: pagues poc per no administrar res.
Exemple amb alpinashop-spark (1 mestre + 2 workers + 2 secundaris, tots n2-standard-4, 20 vCPU en total), com a ordre de magnitud:
| Escenari | Hores al mes | Cost aproximat |
|---|---|---|
| Clúster permanent 24×7 | 720 | De l'ordre de 1.400 € |
Clúster amb --max-idle=30m, 2 h d'ús al dia |
~75 | De l'ordre de 150 € |
| Clúster efímer per flux, 20 min al dia | ~10 | De l'ordre de 20 € |
| Amb secundaris Spot en comptes d'estàndard | ~10 | De l'ordre de 12 € |
| Dataproc Serverless, mateix treball mensual | ~1 | Cèntims |
Verifica els preus vigents a la documentació oficial; el que importa aquí és el factor 100 entre la primera fila i la tercera. Aquest factor és la lliçó sencera.
Palanques d'estalvi, per ordre d'impacte:
--max-idlesempre. Un clúster oblidat un pont de quatre dies costa més que un any de Serverless.- Clústers efímers per flux de treball. Elimina el problema d'arrel.
- Workers secundaris Spot. Descomptes de fins al 80 %, amb el matís que no aporten HDFS i poden desaparèixer.
- Serverless per a treballs ocasionals.
- Discos ajustats. Amb les dades a Cloud Storage, HDFS és només espai de shuffle: 200 GB per worker sobren per a gairebé tot.
- Regió coherent. Clúster i buckets a
europe-west1: llegir dades d'una altra regió costa sortida i latència. - Descomptes per ús compromès només si acabes tenint un clúster permanent, cosa que aquest apartat suggereix evitar.
- Formats columnars. Parquet amb fitxers de 128-512 MB llegeix molt menys i evita el problema dels fitxers petits, que és el màxim assassí de rendiment a Spark sobre emmagatzematge d'objectes.
Errors habituals i consells
Crear un clúster permanent per costum. És l'herència mental del Hadoop on-premise i és l'error més car. Si el clúster no té --max-idle, no hauria d'existir.
Escriure dades importants a HDFS. Desapareixen en destruir el clúster. HDFS a Dataproc és memòria de treball, no emmagatzematge.
No fixar la versió d'imatge. Google actualitza la versió per defecte i un treball que funcionava deixa de funcionar sense que ningú hagi tocat el codi.
Fer servir collect() sobre un DataFrame gran. Porta totes les dades al driver, que és una sola màquina. És la causa número u d'OutOfMemoryError a Spark. Fes servir show(), take(n) o escriu a un fitxer.
Oblidar cache() quan un DataFrame es fa servir diverses vegades. Spark recalcula tota la cadena a cada acció. I l'error invers: posar a la memòria cau tot el que es mou, fins a omplir la memòria i provocar bolcat a disc.
Deixar spark.sql.shuffle.partitions a 200. Amb dades petites genera 200 tasques minúscules amb més sobrecàrrega que feina; amb dades grans, particions enormes que no caben a la memòria.
Molts fitxers petits al llac. Deu mil fitxers d'1 MB són molt més lents de llegir que vint de 500 MB, perquè cada obertura té latència. Compacta.
Posar les dades en una regió i el clúster en una altra. Sortida de dades entre regions facturada i latència afegida a cada lectura.
Consell: fes servir Serverless per defecte. Comença per aquí i crea clúster només quan descobreixis que necessites alguna cosa que Serverless no dona. És el camí amb menys deute operatiu.
Consell: --dry-run no existeix, però el subconjunt sí. Abans de llançar un treball sobre dos anys, executa'l sobre una setmana. Els errors de lògica apareixen igual i costen cent vegades menys.
Consell: mira sempre la interfície de Spark. El component gateway dona accés a la UI de Spark, on es veuen les etapes, les tasques i —el més útil— la distribució de temps entre tasques. Si una tasca triga cent vegades més que la mediana, tens biaix de dades, exactament igual que a Dataflow.
Exercicis
Exercici 1: clúster efímer amb autodestrucció
Crea un clúster anomenat alpinashop-spark-pruebas a europe-west1 que: visqui a sn-datos-euw1 sense IP pública, faci servir el compte de servei sa-dataproc-analitica, tingui 1 mestre n2-standard-2 i 2 workers n2-standard-2, afegeixi 2 workers secundaris Spot, fixi la imatge 2.2-debian12, s'autodestrueixi després de 15 minuts d'inactivitat i en qualsevol cas a les 2 hores de vida, i porti les etiquetes estàndard d'AlpinaShop. Després comprova'n l'estat i esborra'l explícitament.
Exercici 2: treball PySpark de devolucions
Escriu un treball PySpark que llegeixi gs://alpinashop-datalake/pedidos/cabeceras/ i gs://alpinashop-datalake/pedidos/lineas/ en Parquet, i calculi per categoria de producte: nombre de comandes amb almenys una devolució, taxa de devolució sobre el total de comandes d'aquella categoria, i import mitjà retornat. El resultat s'ha d'escriure a alpinashop-datos.alpinashop_analitica.devoluciones_categoria. Aplica almenys dues optimitzacions vistes a la lliçó i explica per què les apliques.
Exercici 3: decidir l'eina
Per a cadascun d'aquests cinc encàrrecs d'AlpinaShop, tria entre BigQuery, Dataflow, Dataproc amb clúster, Dataproc Serverless o bq load, i justifica l'elecció en dues o tres frases:
- Carregar cada nit un fitxer Parquet de 4 GB de l'ERP en una taula de BigQuery, sense cap transformació.
- Calcular el rànquing mensual de vendes per categoria a partir de taules que ja són a
alpinashop_analitica. - Processar els esdeveniments de
pedidos-nuevosa mesura que arriben, agrupant-los per hora de l'esdeveniment i tolerant 90 minuts de retard. - Entrenar cada setmana un model de recomanació amb ALS de MLlib sobre dos anys d'historial de compres.
- Executar el codi PySpark heretat del competidor absorbit, 15.000 línies que fan servir Hive i funcions definides per l'usuari, mentre es decideix què fer-ne.
Solucions
Solució 1
SA="[email protected]"
gcloud dataproc clusters create alpinashop-spark-pruebas \
--region=europe-west1 \
--zone=europe-west1-b \
--subnet=sn-datos-euw1 \
--no-address \
--service-account="$SA" \
--master-machine-type=n2-standard-2 \
--master-boot-disk-size=100GB \
--num-workers=2 \
--worker-machine-type=n2-standard-2 \
--worker-boot-disk-size=100GB \
--num-secondary-workers=2 \
--secondary-worker-type=spot \
--image-version=2.2-debian12 \
--max-idle=15m \
--max-age=2h \
--labels=entorno=desarrollo,equipo=datos,centro-coste=analitica,aplicacion=analitica# Comprovar l estat
gcloud dataproc clusters describe alpinashop-spark-pruebas --region=europe-west1 \
--format="yaml(status.state, config.lifecycleConfig, config.gceClusterConfig.internalIpOnly)"
gcloud dataproc clusters list --region=europe-west1 \
--format="table(clusterName, status.state, config.softwareConfig.imageVersion)"
# Esborrat explicit, sense esperar el max-idle
gcloud dataproc clusters delete alpinashop-spark-pruebas --region=europe-west1 --quietLes dues opcions de cicle de vida són complementàries i convé posar-hi totes dues:
--max-idle=15m: es destrueix si ningú no el fa servir durant 15 minuts. Cobreix el cas "he anat a dinar".--max-age=2h: es destrueix a les 2 hores passi el que passi. Cobreix el cas "un treball en bucle infinit manté el clúster ocupat i--max-idleno es dispara mai". És la fallada que produeix les factures de cap de setmana.
--no-address amb --subnet=sn-datos-euw1 exigeix que Private Google Access estigui activat en aquella subxarxa (03-01) o els nodes no podran baixar paquets ni parlar amb les API. Si el clúster es queda en CREATING i després falla, aquesta és la primera causa que cal revisar.
Solució 2
"""devoluciones.py -- Taxa de devolucio per categoria."""
from pyspark.sql import SparkSession, functions as F
PROJECTE, DATASET = "alpinashop-datos", "alpinashop_analitica"
LLAC = "gs://alpinashop-datalake"
spark = (SparkSession.builder
.appName("alpinashop-devoluciones")
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate())
# OPTIMITZACIO 1: seleccionar nomes les columnes necessaries a la lectura.
# Parquet es columnar: les columnes no demanades no es llegeixen del bucket.
capcaleres = (spark.read.parquet(f"{LLAC}/pedidos/cabeceras/")
.select("pedido_id", "estado", "fecha_pedido", "total_pedido"))
linies = (spark.read.parquet(f"{LLAC}/pedidos/lineas/")
.select("pedido_id", "sku", "importe_linea"))
# El cataleg es llegeix de BigQuery, no del llac
productes = (spark.read.format("bigquery")
.option("table", f"{PROJECTE}.{DATASET}.productos")
.option("filter", "activo = true")
.load()
.select("sku", "categoria"))
# OPTIMITZACIO 2: broadcast del cataleg (milers de files, cap a la memoria).
# Evita completament el shuffle del JOIN amb la taula de linies, que es
# la gran. Sense aixo, tots dos costats es reparticionarien per 'sku'.
linies_cat = linies.join(F.broadcast(productes), on="sku", how="inner")
# Una comanda retornada ho esta sencera: la marquem a nivell de capcalera
cap = capcaleres.withColumn(
"es_devuelto", F.when(F.col("estado") == "devuelto", 1).otherwise(0)
)
# OPTIMITZACIO 3: cache, perque el DataFrame unit es fa servir dues vegades
detall = linies_cat.join(F.broadcast(cap), on="pedido_id", how="inner")
detall.cache()
# Comandes diferents per categoria (una comanda pot tocar diverses categories)
per_categoria = (
detall.groupBy("categoria")
.agg(
F.countDistinct("pedido_id").alias("pedidos_totales"),
F.countDistinct(
F.when(F.col("es_devuelto") == 1, F.col("pedido_id"))
).alias("pedidos_devueltos"),
F.round(
F.avg(F.when(F.col("es_devuelto") == 1, F.col("importe_linea"))), 2
).alias("importe_medio_devuelto"),
)
.withColumn(
"tasa_devolucion_pct",
F.round(100 * F.col("pedidos_devueltos") / F.col("pedidos_totales"), 2),
)
.orderBy(F.desc("tasa_devolucion_pct"))
)
(per_categoria.write.format("bigquery")
.option("table", f"{PROJECTE}.{DATASET}.devoluciones_categoria")
.option("writeMethod", "direct")
.mode("overwrite")
.save())
per_categoria.show(truncate=False)
detall.unpersist()
spark.stop()Les optimitzacions i la seva justificació:
- Projecció primerenca de columnes. Parquet és columnar; demanar quatre columnes en comptes de vint redueix proporcionalment els bytes llegits del bucket i la memòria dels executors. És l'equivalent exacte de no fer
SELECT *a BigQuery. broadcast()als dosJOIN. El catàleg de productes són milers de files i les capçaleres del període també són petites comparades amb les línies. Difondre-les evita repartir per la xarxa la taula gran, que és el cost dominant.cache()sobredetall. Encara que en aquesta versió final només hi ha una agregació, tan bon punt s'afegeixi un segon càlcul (per país, per mes) Spark recalcularia tota la cadena. És la preparació correcta; ambunpersist()en acabar per no retenir memòria.- Filtre delegat a BigQuery (
option("filter", "activo = true")). S'executa allà i arriben menys files.
Nota metodològica que cal explicitar en presentar el resultat: una comanda amb productes de tres categories compta com a comanda a les tres, així que les xifres per categoria no sumen el total de comandes. És correcte per mesurar taxa per categoria, però cal dir-ho a l'informe o algú restarà i no li quadrarà.
Solució 3
1. Carregar un Parquet de 4 GB cada nit sense transformació → bq load.
No hi ha transformació, per tant no cal motor de processament. La càrrega per lots a BigQuery és gratuïta, Parquet porta l'esquema incorporat i no cal dimensionar res. Fer servir Dataflow o Spark aquí seria pagar càlcul per fer una còpia. Una sola comanda, orquestrada a 04-06.
2. Rànquing mensual sobre taules ja a alpinashop_analitica → BigQuery.
Les dades ja hi són i l'operació és una agregació amb funció de finestra, exactament el que vam fer a 04-01 amb RANK() OVER i QUALIFY. Treure les dades de BigQuery per processar-les fora i tornar-les a posar és l'antipatró clàssic: cost de lectura, cost de càlcul, cost d'escriptura i latència, per obtenir un resultat pitjor. Si cal refrescar-ho sovint, vista materialitzada.
3. Esdeveniments de pedidos-nuevos per hora de l'esdeveniment amb 90 minuts de tolerància → Dataflow.
És literalment el cas d'ús per al qual existeix el model de Beam: streaming no acotat, agrupació per temps de l'esdeveniment, finestres fixes, marca d'aigua i allowed_lateness. Structured Streaming de Spark podria, però amb un model de temps menys expressiu i sense la integració nativa amb Pub/Sub. A més, Dataflow dona semàntica d'exactament una vegada cap a BigQuery sense feina addicional.
4. Entrenar ALS de MLlib setmanalment → Dataproc Serverless. ALS és un algorisme iteratiu distribuït implementat a MLlib, sense equivalent a Beam ni en SQL pur. És Spark del principi a la fi. Serverless en comptes de clúster perquè s'executa una vegada per setmana: ningú no vol mantenir ni recordar-se d'apagar un clúster que treballa una hora cada set dies. (El pas següent, servir aquest model en producció, és territori de Vertex AI al mòdul 5.)
5. 15.000 línies de PySpark heretades amb Hive i UDF → Dataproc amb clúster. Aquí mana la restricció pràctica: el codi existeix, funciona i fa servir Hive, que Serverless no inclou. La prioritat és que continuï funcionant amb el mínim canvi, així que clúster de Dataproc amb Dataproc Metastore, executant el codi tal qual contra les dades ja copiades a Cloud Storage. La modernització —passar consultes a BigQuery, streaming a Dataflow— és una fase posterior i separada, mai simultània a la migració: si els números no quadren, cal poder saber si és pel trasllat o per la reescriptura.
Conclusió
Has recorregut el món de Hadoop i Spark amb la perspectiva justa: què són, quin problema van resoldre —processar més dades de les que caben en una màquina, amb màquines que fallen—, per què Spark va desplaçar MapReduce mantenint les dades intermèdies en memòria, i per què el 2026 continuen important encara que existeixin BigQuery i Dataflow: hi ha codi escrit, hi ha persones que el saben fer servir, hi ha biblioteques com MLlib sense equivalent, i hi ha un ecosistema obert que dona portabilitat real.
Has vist què aporta Dataproc: clústers en menys de dos minuts, configurats i coherents, integrats amb Cloud Storage, BigQuery, IAM i Logging. Coneixes l'anatomia —mestre, workers primaris que sostenen HDFS, workers secundaris que només aporten càlcul i per això poden ser Spot— i has creat alpinashop-spark dins de sn-datos-euw1, sense IP pública, amb el compte sa-dataproc-analitica, amb la imatge fixada, amb Jupyter accessible pel component gateway i amb una política d'escalat automàtic que només escala els secundaris i els retira amb elegància.
Però el que és important d'aquesta lliçó no és una comanda, és una inversió conceptual: el clúster efímer sobre emmagatzematge a Cloud Storage. Com que la xarxa interna fa que llegir de gs:// sigui comparable a llegir del disc local, ja no cal que la dada visqui al clúster. I si la dada no viu al clúster, el clúster pot morir. D'aquí surten --max-idle, --max-age, els fluxos de treball amb clúster gestionat que neix i mor amb el job, i el factor cent de diferència a la factura entre un clúster permanent i un d'efímer.
Has escrit el treball PySpark que respon a la pregunta de la Lucía: la cistella mitjana per mes i país amb la seva mediana al costat, i la matriu de productes comprats junts amb suport, confiança i lift —la mètrica que corregeix per la popularitat i evita concloure que tothom compra mitjons amb tot—, aplicant broadcast per evitar el shuffle, cache per no recalcular, un topall de línies per comanda perquè les comandes corporatives rares no dominin el càlcul, i el truc de sku_a < sku_b per comptar cada parell una sola vegada. Coneixes les accions d'inicialització, els seus riscos i per què fixar la versió d'imatge no és una mania.
Has conegut Dataproc Serverless, que elimina fins i tot la decisió de dimensionar, i has fixat el criteri d'AlpinaShop: l'anàlisi de cistella a Serverless, el clúster només com a entorn d'exploració amb Jupyter i autodestrucció. I tens la taula de decisió completa entre Dataproc, Serverless i Dataflow, amb streaming i temps de l'esdeveniment del cantó de Beam, algorismes i ecosistema del cantó de Spark, i SQL sobre dades ja carregades del cantó de BigQuery, sense moure res. Saps com es migra un Hadoop on-premise per fases, amb la regla d'or de no modernitzar i migrar alhora. I saps on són els diners: --max-idle, clústers efímers, Spot, Serverless, discos ajustats i fitxers grans en format columnar.
Queda una promesa pendent des de fa dues lliçons. Tot el que has construït —el pipeline de Dataflow, el treball de Spark, les taules de BigQuery— funciona sobre dades que ja són en algun lloc. Però la botiga continua sent una illa: quan un client confirma una comanda, l'aplicació Flask ha d'avisar el magatzem perquè la prepari, facturació perquè emeti la factura, el servei de correu per a la confirmació, i ara també l'analítica. Si ho fa cridant els quatre un darrere l'altre, la venda es queda penjada esperant el més lent, i si un falla, no queda clar què ha passat amb els altres tres. És un disseny fràgil que es trenca justament el dia de més vendes de l'any.
A 04-04, Cloud Pub/Sub, trencarem aquest acoblament. Crearem per fi el topic pedidos-nuevos i les subscripcions sub-almacen, sub-facturacion i sub-analitica; entendràs les garanties reals de la missatgeria —lliurament almenys una vegada, ordre no garantit— i per què això obliga que els teus consumidors siguin idempotents; veuràs els temes de missatges fallits, els reintents amb retrocés exponencial, la reproducció de missatges amb seek, els filtres per atribut i les subscripcions directes a BigQuery que ingereixen sense escriure ni una línia de codi. I per fi connectarem de debò les notificacions del bucket alpinashop-catalogo que van quedar promeses a 02-02.
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
