El job de vendes de la lliçó anterior funcionava, però cada pregunta nova costava un job més: vendes per productor i mercat n'era un, el rànquing per mercat un altre, unir-ho amb el catàleg un tercer, i entre cadascun la sortida completa anava a HDFS i tornava. Apache Spark va néixer a Berkeley (2009, Matei Zaharia) precisament d'aquesta frustració: els algorismes iteratius d'aprenentatge automàtic i les consultes interactives trigaven a MapReduce deu o cent vegades més del que la CPU justificava, perquè el temps se n'anava a escriure i llegir disc entre fases. Spark conserva l'essencial del model (dades en particions, tasques independents i reexecutables, un shuffle per agrupar per clau) però expressa tot el càlcul com un DAG d'operadors que un planificador converteix en fases i que s'executa mantenint les dades intermèdies en memòria. Sobre aquesta base va construir una API de col·leccions (RDD), després una de taules amb optimitzador (DataFrames i Spark SQL), i a sobre llibreries d'aprenentatge automàtic (MLlib), grafs (GraphX) i fluxos (Structured Streaming, que veurem a 05-04). Avui és el motor de lots de referència. En aquesta lliçó analitica reescriu el càlcul de vendes diàries com a serveis/analitica/vendes_diaries.py, amb DataFrames i amb RDD per comparar, el llança en local i en un clúster de docker-compose.yml, llegeix el pla d'execució, i entrena unes recomanacions amb ALS sobre els clics.
Contingut
- Per què Spark: un DAG en memòria davant de fases a disc
- Arquitectura: driver, cluster manager, executors, tasks i stages
- RDD: col·leccions immutables, transformacions mandroses i llinatge
- DataFrames i Spark SQL: Catalyst i formats columnars
- Operacions estretes i amples, shuffle i stages
- Optimitzacions: cache, broadcast join, particionament i biaix
- MapReduce davant de Spark
- MLlib: recomanacions amb ALS sobre els clics
- Pràctica:
serveis/analitica/vendes_diaries.py - Errors Comuns i Consells
- Exercicis
- Conclusió
- Per què Spark: un DAG en memòria davant de fases a disc
A MapReduce (05-02), una cadena de tres agrupacions són tres jobs i sis passos per disc: cada job escriu la seva sortida a HDFS (replicada tres vegades) i el següent la llegeix. Per a un algorisme iteratiu que recorre les mateixes dades cent vegades, són cent lectures completes des de disc d'un conjunt que no ha canviat. Spark canvia dues coses:
- Tot el càlcul és un únic programa, un DAG d'operadors (llegir, filtrar, explotar, agrupar, unir, escriure) que el motor coneix complet abans d'executar res. Amb aquesta visió global pot encadenar en una sola passada els operadors que no necessiten redistribuir dades, i només materialitza dades intermèdies als punts on hi ha un shuffle. És el patró dataflow de 05-01.
- Les dades intermèdies viuen en memòria dels executors, i es poden desar en memòria cau explícitament per reutilitzar-les en diverses passades. Cent iteracions sobre un conjunt en memòria cau llegeixen disc una vegada.
El preu de no escriure a disc és la tolerància a fallades: si un node mor, les seves dades intermèdies en memòria desapareixen. La solució de Spark, i la seva idea original, és el llinatge (apartat 3): en lloc de desar les dades intermèdies, desa la recepta per recalcular-les a partir de l'entrada, i recomputa només les particions perdudes. És la reexecució determinista de 05-01, aplicada a fragments d'un càlcul en comptes de a tasques senceres.
- Arquitectura: driver, cluster manager, executors, tasks i stages
flowchart TB
subgraph D[Driver: vendes_diaries.py]
SC[SparkSession / SparkContext<br/>DAG scheduler + task scheduler]
end
CM[Cluster manager<br/>standalone · YARN · Kubernetes]
subgraph W1[Node worker 1]
E1[Executor<br/>JVM amb 4 nuclis i 8 GB]
T1a[task] --- E1
T1b[task] --- E1
end
subgraph W2[Node worker 2]
E2[Executor]
T2a[task] --- E2
T2b[task] --- E2
end
SC -- "demana executors" --> CM
CM -- "llança" --> E1
CM -- "llança" --> E2
SC -- "envia tasks, rep resultats" --> E1
SC -- "envia tasks, rep resultats" --> E2
E1 <-- "shuffle" --> E2
E1 --> HDFS[(HDFS / MinIO)]
E2 --> HDFS
- Driver. El procés que executa el programa principal (
vendes_diaries.py). Conté elSparkSession(abansSparkContext), construeix el DAG, el divideix en jobs, stages i tasks, les envia als executors i recull resultats. És el mestre de 05-02, un per aplicació: si el driver mor, l'aplicació mor (a YARN en cluster mode, el driver és l'ApplicationMaster i es pot reintentar). - Cluster manager. Reparteix recursos entre aplicacions. Spark en porta un de propi (standalone, un mestre i workers: el que farem servir a
docker-compose.yml), i s'integra amb YARN (05-02: el driver demana contenidors al ResourceManager) i amb Kubernetes (07-05: cada executor és un pod). Al driver tant li fa quin sigui; l'API és la mateixa. - Executor. Un procés JVM en un node worker, amb N nuclis i M GB assignats, que viu tota l'aplicació. Executa tasks en fils (una task per nucli alhora) i desa en memòria les particions en memòria cau i les dades de shuffle. A diferència de MapReduce, no s'arrenca una JVM per tasca: la JVM s'arrenca una vegada i executa milers de tasques, cosa que elimina el cost fix per tasca.
- Task. La unitat de treball: executar una cadena d'operadors sobre una partició. Una fase amb 200 particions són 200 tasks.
- Stage. Un conjunt de tasks que es poden executar sense shuffle. El DAG es talla en stages a cada operació ampla (apartat 5); dins d'un stage, els operadors s'encadenen (pipelining) i una fila passa per tots ells sense tocar disc.
- Job. Tot el que desencadena una acció (apartat 3):
write,collect,count. Un programa pot tenir diversos jobs; cada job té un o més stages.
Amb PySpark, el driver és un procés Python que es comunica amb una JVM (via Py4J), i als executors el codi Python de les funcions d'RDD s'executa en processos Python auxiliars amb serialització d'anada i tornada; amb DataFrames, en canvi, la majoria d'operacions es tradueixen a codi JVM i Python no intervé per fila. És la raó principal per la qual l'API de DataFrames és la recomanada.
- RDD: col·leccions immutables, transformacions mandroses i llinatge
L'RDD (Resilient Distributed Dataset) és l'abstracció original de Spark: una col·lecció d'elements particionada entre els executors, immutable (no es modifica: se'n crea un altre RDD a partir d'ell) i resilient per llinatge. S'hi opera amb dos tipus de mètodes:
| Tipus | Què fa | Exemples | Executa res |
|---|---|---|---|
| Transformació | Retorna un RDD nou definit a partir d'un altre | map, filter, flatMap, reduceByKey, groupByKey, join, distinct, repartition |
No: és mandrosa, només afegeix un node al DAG |
| Acció | Retorna un resultat al driver o escriu en un sink | collect, count, take, reduce, saveAsTextFile, foreach |
Sí: llança un job que executa totes les transformacions pendents |
La mandra és la que permet optimitzar: quan el programa diu filter i després map i després count, Spark no executa tres passades; en arribar a count coneix tota la cadena i l'executa en una sola passada per partició. I el llinatge és el DAG de transformacions que va portar a cada RDD: si es perd una partició de l'RDD vendes perquè va morir el seu executor, Spark mira el llinatge (vendes = linies.reduceByKey(...), linies = esdeveniments.flatMap(...), esdeveniments = sc.textFile(...)) i recalcula només aquella partició des del bloc d'HDFS corresponent. No cal desar res intermedi a disc; el cost és recomputar, que és assumible mentre el llinatge sigui curt (per a llinatges llargs, iteratius, existeix checkpoint(), que sí que materialitza a HDFS i talla el llinatge).
El càlcul de vendes amb RDD, que servirà de comparació a la pràctica:
# Fragment de serveis/analitica/vendes_diaries.py (versió RDD, vegeu l'apartat 9)
import json
from datetime import datetime, timezone
def dia_de(ev): # temps d'esdeveniment (01-05), no de procés
return datetime.fromtimestamp(ev["data_ms"] / 1000, tz=timezone.utc).strftime("%Y-%m-%d")
def vendes_rdd(sc, entrada):
esdeveniments = sc.textFile(entrada) # RDD[str]: una partició per bloc HDFS
comandes = esdeveniments.map(json.loads).filter(lambda e: e["tipus"] == "comanda.creada")
linies = comandes.flatMap(lambda e: [ # una fila per línia de comanda
((dia_de(e), e["dades"]["mercat"], ln["productor"]), ln["quantitat"] * ln["preu"])
for ln in e["dades"]["linies"]])
vendes = linies.reduceByKey(lambda a, b: a + b) # ampla: shuffle per clau; combina localment abans
return vendes # encara no s'ha executat resFins aquí no s'ha llegit ni un byte: textFile, map, filter, flatMap i reduceByKey són transformacions. vendes.collect() o vendes.saveAsTextFile(...) llançarien el job. Dos detalls que distingeixen un usuari de Spark: reduceByKey fa la reducció local a cada partició abans del shuffle (el combiner de 05-02, automàtic), mentre que groupByKey seguit d'una suma mou tots els valors per la xarxa; i les funcions lambda viatgen serialitzades als executors, així que no poden capturar objectes no serialitzables (una connexió a base de dades, per exemple).
- DataFrames i Spark SQL: Catalyst i formats columnars
Un RDD és una col·lecció d'objectes opacs per a Spark: no sap que ln["productor"] és una columna, així que no pot optimitzar res més enllà de l'encadenament. Un DataFrame és una taula amb esquema (columnes amb nom i tipus), distribuïda en particions com un RDD, i sobre la qual s'opera amb una API relacional (select, filter, groupBy, join, agg) o directament amb SQL (spark.sql("SELECT ...")). Totes dues produeixen el mateix pla lògic, que passa per Catalyst, l'optimitzador:
- Anàlisi: resoldre noms de columnes i tipus contra l'esquema.
- Optimització lògica: regles com empènyer els filtres cap a la lectura (predicate pushdown: filtrar
tipus = 'comanda.creada'en llegir, no després), llegir només les columnes necessàries (column pruning), simplificar expressions, reordenar joins. - Planificació física: triar algorismes: un join es fa per broadcast hash si un costat és petit, per sort-merge si no; una agregació en dues fases (parcial a cada partició, final després del shuffle).
- Generació de codi: Tungsten compila el pla a bytecode Java especialitzat (whole-stage codegen) i gestiona la memòria fora del heap de la JVM en format binari compacte. És el que fa que PySpark amb DataFrames sigui tan ràpid com Scala: Python només descriu el pla.
Des de Spark 3, l'Adaptive Query Execution (AQE) reoptimitza el pla durant l'execució amb estadístiques reals: redueix el nombre de particions després d'un shuffle petit, converteix un sort-merge join en broadcast si descobreix que un costat hi cap, i divideix particions esbiaixades (apartat 6).
Els DataFrames es combinen amb els formats columnars. Un fitxer Parquet desa les dades per columnes en lloc de per files: tots els productor junts, tots els import_eur junts, comprimits amb codificacions adequades a cada tipus (diccionari per a cadenes repetides, run-length per a valors consecutius) i amb estadístiques (mínim, màxim, nuls) per bloc de files. Una consulta que només necessita productor i import_eur llegeix aquestes dues columnes i salta la resta del fitxer, i un filtre dia = '2026-09-14' salta blocs sencers el rang dels quals no el conté. Davant de JSONL, que obliga a llegir i parsejar cada byte, Parquet redueix en un ordre de magnitud tant els bytes llegits com la mida a disc. Per això el llac de 04-02 comença en JSONL (el que produeix Kafka) i el primer pas de l'analítica el converteix a Parquet, particionat per dia.
- Operacions estretes i amples, shuffle i stages
La divisió del DAG en stages depèn d'un únic criteri: quina dependència té cada partició de sortida respecte de les d'entrada.
| Dependència estreta (narrow) | Dependència ampla (wide) | |
|---|---|---|
| Cada partició de sortida depèn de | Una partició d'entrada (o unes poques fixes) | Totes (o moltes) les particions d'entrada |
| Operacions | map, filter, flatMap, select, withColumn, union, coalesce, join amb broadcast, join amb particionament idèntic |
groupBy/reduceByKey, distinct, join sort-merge, orderBy, repartition |
| Cost | S'encadena a la mateixa task, sense xarxa | Shuffle: escriure, transferir, llegir; talla el DAG en un stage nou |
| Recuperació després d'una fallada | Recomputar la partició perduda des de la seva única entrada | Recomputar pot exigir rellegir moltes particions (Spark desa els fitxers de shuffle per evitar-ho) |
flowchart LR
subgraph S0[Stage 0: sense shuffle]
R[read json<br/>2 particions] --> F[filter tipus] --> X[explode línies] --> P[agg parcial<br/>per dia, mercat, productor]
end
P == "shuffle<br/>hashpartitioning(dia, mercat, productor)<br/>200 particions" ==> S1
subgraph S1[Stage 1]
A[agg final] --> J[broadcast join<br/>amb catàleg] --> W[write parquet]
end
C[read catàleg<br/>1 partició] -. broadcast a tots els executors .-> J
El DAG de vendes té exactament un shuffle, el groupBy, i per tant dos stages. Tota la resta s'encadena: una fila del JSON es filtra, s'explota i s'agrega parcialment a la mateixa task sense que ningú no l'escrigui. El join amb el catàleg, que a MapReduce seria un tercer job, és estret gràcies al broadcast (apartat 6). A la interfície web de Spark (http://localhost:4040 durant l'execució) cada job apareix amb els seus stages, cada stage amb les seves tasks, i per a cada stage els bytes de shuffle write i shuffle read: són els comptadors de 05-02, i el criteri de diagnòstic és el mateix.
El shuffle de Spark escriu també a disc local (els fitxers de shuffle, servits pels executors o per un external shuffle service), així que "en memòria" no vol dir "sense disc": vol dir que entre operacions estretes no s'escriu, i que les dades en memòria cau se serveixen de memòria. Amb 200 particions de shuffle per defecte (spark.sql.shuffle.partitions), un job petit produeix 200 tasks minúscules al segon stage; AQE les fusiona, però en versions antigues o amb AQE desactivat convé abaixar el nombre.
- Optimitzacions: cache, broadcast join, particionament i biaix
cache() i persist(). Marquen un DataFrame o RDD perquè, la primera vegada que es calculi, les seves particions es desin (en memòria per defecte; persist(StorageLevel.MEMORY_AND_DISK) permet desbordar a disc, DISK_ONLY, o serialitzat per estalviar memòria). Compensa quan el mateix resultat intermedi es fa servir en diverses accions: el DataFrame linies de la pràctica alimenta les vendes per productor, un rànquing i una comprovació de qualitat; sense cache, cada acció torna a llegir i parsejar el JSON d'HDFS. Amb unpersist() s'allibera. Desar en memòria cau una cosa que es fa servir una vegada és pur cost.
Broadcast join. Unir les vendes (milions de files, repartides) amb el catàleg (unes desenes de files) no hauria de moure les vendes. Amb F.broadcast(cataleg), Spark envia una còpia del catàleg a cada executor i el join es resol localment, sense shuffle del costat gran: és el map-side join de 05-02, automàtic. Spark ho fa sol si estima que el costat petit ocupa menys de spark.sql.autoBroadcastJoinThreshold (10 MB per defecte); forçar-ho amb broadcast() convé quan l'estimació falla (un CSV sense estadístiques). Amb dos costats grans, el pla és un sort-merge join: shuffle de tots dos per la clau i barreja ordenada.
repartition(n) i coalesce(n). El nombre de particions governa el paral·lelisme. repartition(n) fa un shuffle complet per obtenir n particions equilibrades (o repartition("dia") per col·locar cada dia a la seva partició, útil abans d'escriure particionat); coalesce(n) redueix el nombre sense shuffle, fusionant particions locals, i és la manera barata de no escriure 200 fitxers Parquet de 3 KB. Com a orientació: particions de 100–200 MB en memòria i entre 2 i 4 tasks per nucli disponible.
Biaix amb salting. La Setmana del Formatge Artesà torna: en agrupar per productor, la partició de formatgeria-montblanc rep la meitat de les files. En DataFrames, l'agregació parcial per partició (que Catalyst insereix sempre) alleuja el cas de les sumes, igual que el combiner; el biaix fa mal quan l'operació no es preredueix (un join sort-merge, collect_list, finestres). La tècnica de 05-01, ara amb columnes:
from pyspark.sql import functions as F
N_SAL = 8
# 1) afegir sal determinista a la clau calenta (derivada de l'id de comanda: reexecutable)
amb_sal = linies.withColumn("sal", F.when(F.col("productor") == "formatgeria-montblanc",
F.pmod(F.hash("comanda_id"), F.lit(N_SAL))).otherwise(F.lit(0)))
# 2) agregar per (clau, sal): la clau calenta es reparteix en 8 particions
parcial = amb_sal.groupBy("dia", "mercat", "productor", "sal").agg(F.sum("import_eur").alias("import_eur"))
# 3) segona agregació, ja petita, traient la sal
vendes = parcial.groupBy("dia", "mercat", "productor").agg(F.sum("import_eur").alias("import_eur"))I des de Spark 3, AQE amb spark.sql.adaptive.skewJoin.enabled=true detecta particions de join esbiaixades (per mida relativa a la mediana) i les divideix automàticament, sense sal manual. Per a agregacions esbiaixades sense prereducció possible continua calent la sal.
- MapReduce davant de Spark
| MapReduce (Hadoop) | Spark | |
|---|---|---|
| Model | Map i reduce; cadenes de jobs | DAG d'operadors en un programa |
| Dades intermèdies | Disc local + HDFS entre jobs | Memòria (i disc local només al shuffle) |
| Tolerància a fallades | Reexecució de tasques des de disc | Recomputació per llinatge; fitxers de shuffle; checkpoint opcional |
| Cost fix | JVM per tasca; desenes de segons per job | Executors persistents; tasks de mil·lisegons |
| API | Java (Streaming per a altres llenguatges); només map/reduce | Scala, Java, Python, R, SQL; desenes d'operadors; DataFrames amb optimitzador |
| Joins, iteracions | A mà, diversos jobs | Natius; cache per iterar |
| Iteratiu (ML, grafs) | 10–100× més lent per relectures | Dissenyat per a això (MLlib, GraphX) |
| Interactiu | No | Sí (spark-shell, notebooks, Spark SQL) |
| Streaming | No | Structured Streaming (05-04) |
| Memòria necessària | Modesta | Més gran: els executors necessiten RAM per a cache i shuffle |
| Estat actual | Base històrica; YARN i HDFS continuen en ús | Motor de lots de referència |
L'avantatge de Spark en el job de vendes de 05-02 es veu en els números: el mateix càlcul que a MapReduce trigava 52 s al clúster de proves triga 6 s a Spark sobre YARN amb els mateixos recursos, i encadenar el rànquing i el join amb el catàleg no afegeix jobs, sinó un stage. La contrapartida és la memòria: un executor mal dimensionat (massa poca memòria per al shuffle o la cache) falla amb OutOfMemoryError on MapReduce simplement escrivia a disc.
- MLlib: recomanacions amb ALS sobre els clics
El segon cas de l'analítica de Quilòmetre Zero, "productes recomanats per a l'Anna", és un problema iteratiu: exactament el tipus de problema que MapReduce resolia malament. MLlib porta algorismes distribuïts a punt, i per a recomanació el clàssic és ALS (Alternating Least Squares), una factorització de la matriu clients × productes que aprèn un vector de factors per client i un altre per producte a partir d'interaccions (compres, clics) i prediu l'interès de cada parella no observada. Amb clics no hi ha "puntuació", sinó senyals implícits (va veure, va afegir al cistell, va comprar), i ALS té un mode per a això.
Els clics del llac, /km0/clics/2026-09-14/hora=13/web-01.jsonl (04-02), tenen aquesta forma:
{"data_ms":1789390812000,"client":"anna","producte":"formatge-curat","accio":"vist"}
{"data_ms":1789390834000,"client":"anna","producte":"vi-crianca","accio":"cistell"}
{"data_ms":1789390901000,"client":"marc","producte":"vi-crianca","accio":"comprat"}
{"data_ms":1789391010000,"client":"llucia","producte":"formatge-fresc","accio":"vist"}
{"data_ms":1789391044000,"client":"llucia","producte":"carbasso","accio":"comprat"}# km0/serveis/analitica/recomanacions_als.py
from pyspark.sql import SparkSession, functions as F
from pyspark.ml.feature import StringIndexer
from pyspark.ml.recommendation import ALS
spark = SparkSession.builder.appName("km0-recomanacions").getOrCreate()
clics = spark.read.json("hdfs://namenode:8020/km0/clics/2026-09-*/") # 7 dies, totes les hores
pes = F.when(F.col("accio") == "comprat", 5).when(F.col("accio") == "cistell", 2).otherwise(1)
interes = clics.withColumn("pes", pes).groupBy("client", "producte").agg(F.sum("pes").alias("interes"))
# ALS necessita ids enters: StringIndexer n'assigna un per valor diferent
idx_cli = StringIndexer(inputCol="client", outputCol="client_id").fit(interes)
idx_pro = StringIndexer(inputCol="producte", outputCol="producte_id").fit(interes)
dades = idx_pro.transform(idx_cli.transform(interes))
als = ALS(userCol="client_id", itemCol="producte_id", ratingCol="interes",
implicitPrefs=True, # els senyals són implícits (clics), no notes
rank=10, maxIter=10, regParam=0.1, coldStartStrategy="drop", seed=42)
model = als.fit(dades) # 10 iteracions: cadascuna alterna factors de clients i productes
recomanacions = model.recommendForAllUsers(3) # 3 productes per client
etiquetes = spark.createDataFrame(enumerate(idx_pro.labels), ["producte_id", "producte"])
(recomanacions.select("client_id", F.explode("recommendations").alias("r"))
.join(etiquetes, F.col("r.producte_id") == etiquetes.producte_id)
.join(spark.createDataFrame(enumerate(idx_cli.labels), ["client_id", "client"]), "client_id")
.select("client", "producte", F.round("r.rating", 3).alias("afinitat"))
.write.mode("overwrite").parquet("hdfs://namenode:8020/km0/recomanacions/2026-09-14"))Cada iteració de fit és un job amb diversos stages sobre les mateixes dades, que MLlib desa en memòria cau internament; amb 10 iteracions a MapReduce serien 20 jobs i 20 lectures d'HDFS. El resultat (anna → vi-crianca, llucia → formatge-curat...) el carrega el pipeline de 05-05 a la base de cataleg perquè el web el mostri. El que MLlib no decideix és la qualitat: triar rank, regParam i els pesos de les accions exigeix avaluar (RegressionEvaluator o mètriques de rànquing sobre un conjunt reservat), i això pertany a un curs d'aprenentatge automàtic, no a aquest.
- Pràctica:
serveis/analitica/vendes_diaries.py
serveis/analitica/vendes_diaries.py9.1 Entorn: Spark a docker-compose.yml
Al docker-compose.yml de 04-02 (HDFS) i 05-02 (YARN) s'hi afegeix un clúster Spark standalone, que és més lleuger que YARN per desenvolupar. El driver correrà a la nostra màquina o al contenidor del mestre:
# km0/docker-compose.yml (fragment)
services:
spark-master:
image: bitnami/spark:3.5
environment:
- SPARK_MODE=master
ports:
- "8080:8080" # interfície web del mestre standalone
- "7077:7077" # port al qual es connecten drivers i workers
volumes:
- ./serveis/analitica:/app/analitica
- ./esdeveniments:/app/esdeveniments
spark-worker:
image: bitnami/spark:3.5
environment:
- SPARK_MODE=worker
- SPARK_MASTER_URL=spark://spark-master:7077
- SPARK_WORKER_CORES=2
- SPARK_WORKER_MEMORY=2G
deploy:
replicas: 2 # docker compose up --scale spark-worker=2
volumes:
- ./serveis/analitica:/app/analitica
- ./esdeveniments:/app/esdevenimentsEls dos workers ofereixen 4 nuclis en total, i HDFS és abastable com a hdfs://namenode:8020 des de la mateixa xarxa de Compose. El catàleg és un CSV petit, serveis/analitica/cataleg.csv, que en producció vindria d'una exportació diària de km0_cataleg:
productor,nom_productor,provincia
horta-la-vega,Horta La Vega,Girona
formatgeria-montblanc,Formatgeria Montblanc,Tarragona
celler-roure-alt,Celler Roure Alt,Lleida9.2 El programa amb DataFrames
# km0/serveis/analitica/vendes_diaries.py
"""Vendes per productor, mercat i dia a partir dels esdeveniments comanda.creada del llac.
Ús: spark-submit vendes_diaries.py <entrada> <cataleg.csv> <sortida> [--rdd]
entrada: hdfs://namenode:8020/km0/esdeveniments/2026-09-14/comandes.jsonl (o una ruta local, o un glob)
sortida: hdfs://namenode:8020/km0/agregats/vendes_diaries (Parquet particionat per dia)
"""
import sys
from pyspark.sql import SparkSession, functions as F, types as T
ESQUEMA = T.StructType([ # declarar l'esquema evita una passada d'inferència
T.StructField("id_esdeveniment", T.StringType()),
T.StructField("tipus", T.StringType()),
T.StructField("versio", T.IntegerType()),
T.StructField("data_ms", T.LongType()),
T.StructField("origen", T.StringType()),
T.StructField("dades", T.StructType([
T.StructField("comanda_id", T.StringType()),
T.StructField("client", T.StringType()),
T.StructField("mercat", T.StringType()),
T.StructField("linies", T.ArrayType(T.StructType([
T.StructField("producte", T.StringType()),
T.StructField("productor", T.StringType()),
T.StructField("quantitat", T.IntegerType()),
T.StructField("preu", T.DoubleType()),
]))),
])),
])
def linies_de_comanda(spark, entrada):
"""DataFrame amb una fila per línia de comanda: dia, mercat, comanda_id, producte, productor, quantitat, import_eur."""
esdeveniments = spark.read.schema(ESQUEMA).json(entrada)
comandes = esdeveniments.filter(F.col("tipus") == "comanda.creada") # Catalyst ho empeny a la lectura
return (comandes
.select(
F.to_date(F.from_unixtime(F.col("data_ms") / 1000)).alias("dia"), # temps d'esdeveniment
F.col("dades.mercat").alias("mercat"),
F.col("dades.comanda_id").alias("comanda_id"),
F.explode("dades.linies").alias("ln")) # una fila per element de l'array
.select("dia", "mercat", "comanda_id",
F.col("ln.producte").alias("producte"),
F.col("ln.productor").alias("productor"),
F.col("ln.quantitat").alias("quantitat"),
(F.col("ln.quantitat") * F.col("ln.preu")).alias("import_eur")))
def vendes_dataframe(spark, entrada, ruta_cataleg):
linies = linies_de_comanda(spark, entrada).cache() # es fa servir en dues accions (vendes i control)
cataleg = spark.read.option("header", True).csv(ruta_cataleg) # 3 files: candidat a broadcast
vendes = (linies
.groupBy("dia", "mercat", "productor") # UNA operació ampla: un shuffle
.agg(F.round(F.sum("import_eur"), 2).alias("import_eur"),
F.sum("quantitat").alias("unitats"),
F.countDistinct("comanda_id").alias("comandes"))
.join(F.broadcast(cataleg), "productor", "left") # estreta: el catàleg viatja a cada executor
.select("dia", "mercat", "productor", "nom_productor", "provincia", "import_eur", "unitats", "comandes"))
control = linies.agg(F.countDistinct("comanda_id").alias("comandes"), F.round(F.sum("import_eur"), 2).alias("total"))
print("Control:", control.first().asDict()) # acció 1: fa servir la cache
return vendes # l'acció 2 serà el write
if __name__ == "__main__":
entrada, ruta_cataleg, sortida = sys.argv[1:4]
spark = (SparkSession.builder.appName("km0-vendes-diaries")
.config("spark.sql.shuffle.partitions", "8") # job petit: 200 seria absurd
.config("spark.sql.sources.partitionOverwriteMode", "dynamic") # sobreescriure NOMÉS les particions escrites
.getOrCreate())
if "--rdd" in sys.argv:
from vendes_rdd import vendes_rdd # versió RDD de l'apartat 3
for (dia, mercat, productor), import_eur in sorted(vendes_rdd(spark.sparkContext, entrada).collect()):
print(f"{dia} {mercat:10s} {productor:22s} {import_eur:12,.2f}")
else:
vendes = vendes_dataframe(spark, entrada, ruta_cataleg)
vendes.explain() # pla físic (apartat 9.3)
(vendes.coalesce(1) # un fitxer per partició de sortida
.write.mode("overwrite")
.partitionBy("dia") # sortida/dia=2026-09-14/part-....parquet
.parquet(sortida))
vendes.orderBy("dia", "mercat", "productor").show(truncate=False)
spark.stop()Punts clau del programa:
- Esquema explícit. Sense ell,
spark.read.jsonfa una passada completa només per inferir tipus. Amb ell, la lectura és una passada i els tipus són els que volem (data_mscom aLongType, nodouble). explodeconverteix l'arrayliniesen files, que és el que feia el buclefor ln in ev["dades"]["linies"]del mapper.- Un sol shuffle. El
groupByés l'única operació ampla.filter,select,explodei el join broadcast s'encadenen al mateix stage que la lectura. Catalyst insereix una agregació parcial abans del shuffle (HashAggregateambpartial_sum), així que el que travessa la xarxa són parcials per(dia, mercat, productor)per partició: el combiner, sense escriure'l. cache()aliniesperquè hi ha dues accions (control.first()i elwrite); sense ell, el JSON es llegiria i parsejaria dues vegades.partitionBy("dia")ambpartitionOverwriteMode=dynamic. La sortida és un directori per dia (dia=2026-09-14/), ioverwriteen mode dinàmic reemplaça només els dies que aquest job escriu, deixant intactes els altres. Executar dues vegades el job del 14 de setembre produeix exactament la mateixa sortida: és la idempotència per partició que el pipeline de 05-05 necessita per reprocessar i fer backfill.coalesce(1)abans d'escriure, perquè l'agregat són unes 12 files per dia i no volem 8 fitxers Parquet de 2 KB. Amb agregats grans es trauria o s'ajustaria.
9.3 Llançament i pla d'execució
En local, amb quatre fils com a executors simulats (local[4]), sobre el fitxer generat a 05-01:
$ pip install pyspark==3.5.1
$ spark-submit --master 'local[4]' serveis/analitica/vendes_diaries.py \
esdeveniments/2026-09-14/comandes.jsonl serveis/analitica/cataleg.csv sortida/vendes_diaries
Control: {'comandes': 400000, 'total': 3590252.0}
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [dia, mercat, productor, nom_productor, provincia, import_eur, unitats, comandes]
+- BroadcastHashJoin [productor], [productor], LeftOuter, BuildRight
:- HashAggregate(keys=[dia, mercat, productor], functions=[sum(import_eur), sum(quantitat), count(distinct comanda_id)])
: +- Exchange hashpartitioning(dia, mercat, productor, 8)
: +- HashAggregate(keys=[dia, mercat, productor], functions=[partial_sum(import_eur), partial_sum(quantitat), ...])
: +- InMemoryTableScan [dia, mercat, comanda_id, productor, quantitat, import_eur]
: +- InMemoryRelation ... (linies en memòria cau)
: +- Generate explode(dades.linies) ...
: +- Filter (tipus = comanda.creada)
: +- FileScan json [tipus, data_ms, dades] PushedFilters: [IsNotNull(tipus), EqualTo(tipus,comanda.creada)]
+- BroadcastExchange HashedRelationBroadcastMode
+- FileScan csv [productor, nom_productor, provincia]
+----------+---------+---------------------+---------------------+---------+----------+-------+--------+
|dia |mercat |productor |nom_productor |provincia|import_eur|unitats|comandes|
+----------+---------+---------------------+---------------------+---------+----------+-------+--------+
|2026-09-14|girona |celler-roure-alt |Celler Roure Alt |Lleida |246187.40 |25121 |23967 |
|2026-09-14|girona |formatgeria-montblanc|Formatgeria Montblanc|Tarragona|473402.10 |46330 |41218 |
|2026-09-14|girona |horta-la-vega |Horta La Vega |Girona |178013.70 |58940 |31502 |
...El pla es llegeix de baix a dalt, i cada línia confirma una decisió dels apartats anteriors: FileScan json amb PushedFilters (el filtre empès a la lectura), Generate explode, InMemoryRelation (la cache), el HashAggregate amb partial_sum abans de l'Exchange hashpartitioning(..., 8) (agregació parcial, després l'únic shuffle amb 8 particions) i el HashAggregate final, i BroadcastHashJoin amb BroadcastExchange del CSV (el catàleg viatja, les vendes no). Si en lloc de BroadcastHashJoin hi aparegués SortMergeJoin amb dos Exchange, sabríem que el broadcast no s'ha aplicat i que estem pagant un shuffle de més.
Contra el clúster standalone de Compose, amb el fitxer a HDFS:
$ docker compose exec spark-master spark-submit \
--master spark://spark-master:7077 \
--executor-memory 1G --executor-cores 2 --num-executors 2 \
/app/analitica/vendes_diaries.py \
hdfs://namenode:8020/km0/esdeveniments/2026-09-14/comandes.jsonl \
/app/analitica/cataleg.csv \
hdfs://namenode:8020/km0/agregats/vendes_diaries
$ docker compose exec namenode hdfs dfs -ls /km0/agregats/vendes_diaries/
drwxr-xr-x - spark supergroup 0 /km0/agregats/vendes_diaries/dia=2026-09-14
-rw-r--r-- 2 spark supergroup 0 /km0/agregats/vendes_diaries/_SUCCESSAmb --master yarn i la configuració de Hadoop a HADOOP_CONF_DIR, el mateix fitxer es llançaria sobre el YARN de 05-02, amb el driver com a ApplicationMaster (--deploy-mode cluster), sense canviar una línia de codi. A la interfície del mestre (localhost:8080) es veuen els workers i les aplicacions; a la del driver (localhost:4040, mentre corre) els jobs, stages i tasks, amb els seus bytes de shuffle.
9.4 La versió RDD, per comparar
vendes_rdd.py conté la funció de l'apartat 3. Llançat amb --rdd, produeix les mateixes xifres, i a la interfície es veuen dues diferències: l'stage de lectura és més lent (cada línia passa per un procés Python amb json.loads, en lloc del parser JSON de la JVM) i no hi ha pla a llegir: Spark executa les lambdes tal qual, sense empènyer filtres ni podar columnes, perquè no sap què fan. Al fitxer de 130 MB la diferència és de 9 s davant de 4 s a local[4]; en terabytes, d'hores. L'RDD continua sent l'eina correcta quan les dades no tenen esquema o la lògica no cap en expressions de columna, i és el que hi ha sota de tot DataFrame; però per a analítica sobre esdeveniments amb esquema, DataFrames és l'API.
Errors Comuns i Consells
collect()sobre un DataFrame gran. Porta totes les files al driver, que té uns GB de memòria:OutOfMemoryErroral driver. Per mirar,show()otake(n); per desar,write.- Oblidar que les transformacions són mandroses. Un
filtermal escrit no falla en escriure'l, sinó a la primera acció, amb una traça que apunta alwrite. I unprintdins d'una lambda no apareix al driver: s'executa als executors (mira'n els logs). groupByKey+ suma en RDD. Mou tots els valors per la xarxa.reduceByKeyoaggregateByKeycombinen localment. En DataFrames,groupBy().agg()ja ho fa.- UDF de Python per fila. Una
udfa PySpark serialitza cada fila cap a un procés Python i de tornada; anul·la Catalyst i Tungsten. Gairebé tot es pot expressar ambF.*; si no, fes servirpandas_udf(vectoritzada per lots amb Arrow). - Desar en memòria cau sense fer-ho servir o no alliberar.
cache()sobre una cosa que es fa servir una vegada és un cost sense benefici; desar-ne moltes senseunpersist()desborda la memòria dels executors i provoca desallotjaments i recomputacions silencioses. - 200 particions de shuffle per a 10 MB. El valor per defecte de
spark.sql.shuffle.partitionsés per a clústers grans. Ajusta'l o confia en AQE; icoalesceabans d'escriure per no produir centenars de fitxers d'un quilobyte que ofegaran el NameNode (04-02). - Confiar en el broadcast automàtic amb un CSV. Sense estadístiques, Spark pot no estimar la mida i fer un sort-merge join.
F.broadcast()explícit i comprovar-ho aexplain(). - Inferir l'esquema del JSON en producció. És una passada extra i un esquema que canvia amb les dades (un dia sense
versioi la columna desapareix). Esquema explícit, versionat amb el contracte de l'esdeveniment (02-05). - Escriure amb
overwritesensepartitionOverwriteMode=dynamic. Esborra tot el directori de sortida, inclosos els dies que no s'estaven recalculant. És l'error més car que es comet amb un pipeline de backfill.
Exercicis
Exercici 1: Rànquing per mercat en un sol programa
Amplia vendes_dataframe perquè, a més de les vendes, produeixi per a cada (dia, mercat) el productor amb més vendes i la seva quota sobre el total del mercat (percentatge), fent servir funcions de finestra (pyspark.sql.Window). Quants shuffles afegeix la teva solució? Compara-ho amb els dos jobs de l'exercici 1 de 05-02.
Exercici 2: Llegir un pla
Aquest és el pla d'una versió modificada del programa, escrita per un company. Identifica quatre problemes de rendiment a partir del pla i digues com corregir cadascun.
== Physical Plan ==
+- SortMergeJoin [productor], [productor], LeftOuter
:- Sort [productor ASC]
: +- Exchange hashpartitioning(productor, 200)
: +- HashAggregate(keys=[dia, mercat, productor], functions=[sum(import_eur)])
: +- Exchange hashpartitioning(dia, mercat, productor, 200)
: +- HashAggregate(keys=[dia, mercat, productor], functions=[partial_sum(import_eur)])
: +- BatchEvalPython [calcular_import(quantitat, preu)]
: +- Generate explode(dades.linies)
: +- Filter (tipus = comanda.creada)
: +- FileScan json [id_esdeveniment, tipus, versio, data_ms, origen, dades]
+- Sort [productor ASC]
+- Exchange hashpartitioning(productor, 200)
+- FileScan csv [productor, nom_productor, provincia]Exercici 3: Reprocessar la Setmana de la Verema
Els esdeveniments del 8 al 14 de setembre de 2026 (la Setmana de la Verema) tenien un error en el preu de vi-crianca que comandes ha corregit regenerant els fitxers comandes.jsonl d'aquests set dies a HDFS. Escriu la invocació (o les invocacions) de vendes_diaries.py que recalculi exactament aquests set dies sense tocar els altres, i explica quina combinació d'opcions del programa garanteix que (a) no quedin dades antigues d'aquests dies, (b) no s'esborrin els altres dies, i (c) executar-ho dues vegades doni el mateix resultat. Què canviaries si l'entrada fos un glob comandes-*.jsonl amb diversos fitxers per dia?
Solucions
Exercici 1.
from pyspark.sql import Window
per_mercat = Window.partitionBy("dia", "mercat")
ranquing = (vendes
.withColumn("total_mercat", F.sum("import_eur").over(per_mercat))
.withColumn("quota", F.round(F.col("import_eur") / F.col("total_mercat") * 100, 1))
.withColumn("posicio", F.row_number().over(per_mercat.orderBy(F.desc("import_eur"))))
.filter(F.col("posicio") == 1)
.select("dia", "mercat", "productor", "import_eur", "quota"))Les dues finestres comparteixen la partició (dia, mercat), així que Spark afegeix un shuffle (Exchange hashpartitioning(dia, mercat)) i un Sort dins de cada partició per a row_number; amb AQE i les mateixes claus, de vegades reutilitza l'intercanvi. Total: dos shuffles al programa (el groupBy i la finestra), en un sol job amb tres stages, davant dels dos jobs de MapReduce amb els seus quatre passos per HDFS. I el join amb el catàleg continua sense costar shuffle perquè vendes ja el portava resolt.
Exercici 2.
SortMergeJoinamb dosExchangeperproductoren lloc deBroadcastHashJoin: el catàleg de tres files està provocant un shuffle de les vendes i una ordenació. Correcció:F.broadcast(cataleg).BatchEvalPython [calcular_import]: una UDF de Python per fila, que a més impedeix que elpartial_sumes computi a la JVM sense sortir a Python. Correcció:F.col("ln.quantitat") * F.col("ln.preu").Exchange ... 200: 200 particions de shuffle per a un agregat de desenes de files; 200 tasks minúscules per stage. Correcció:spark.sql.shuffle.partitionsa 8 (o AQE amb coalescència activada).FileScan json [id_esdeveniment, tipus, versio, data_ms, origen, dades]sensePushedFiltersni poda: es llegeixen totes les columnes, i probablement sense esquema explícit (inferència, una passada extra). Correcció: esquema declarat iselectde les columnes necessàries just després de la lectura; el filtretipushauria d'aparèixer com aPushedFilters(ho fa quan la columna es compara amb un literal i l'esquema és conegut).
Un cinquè detall: no hi ha InMemoryRelation, així que si el programa fa més d'una acció, rellegirà el JSON cada vegada.
Exercici 3.
Una sola invocació amb un glob per als set dies, o set invocacions (una per dia, que és el que farà Airflow a 05-05 amb el backfill):
spark-submit --master spark://spark-master:7077 /app/analitica/vendes_diaries.py \
'hdfs://namenode:8020/km0/esdeveniments/2026-09-{08,09,10,11,12,13,14}/comandes.jsonl' \
/app/analitica/cataleg.csv hdfs://namenode:8020/km0/agregats/vendes_diaries(a) mode("overwrite") reemplaça el contingut de cada partició dia=2026-09-NN/ que el job escriu, de manera atòmica per partició (escriptura en temporal i commit). (b) partitionOverwriteMode=dynamic limita l'esborrat a les particions presents a la sortida del job; sense ell, overwrite buidaria vendes_diaries/ sencer, inclosos l'agost i la resta de setembre. (c) El càlcul és determinista (mateixa entrada, mateix agregat) i l'escriptura reemplaça la partició completa: dues execucions deixen els mateixos bytes (tret de noms de fitxer interns), sense duplicar ni acumular. És important que dia es derivi de data_ms (temps d'esdeveniment) i no del nom del directori: si un esdeveniment del 14 arribés tard i comandes l'hagués deixat al fitxer del 15, el job del 15 l'escriuria a dia=2026-09-14/, sobreescrivint la partició del 14 amb només aquell esdeveniment. Per evitar-ho hi ha dues opcions: filtrar al job pel rang de dies que s'està processant (F.col("dia").between(...)) i descartar o registrar els de fora de rang, o que la partició d'entrada i sortida coincideixin per construcció. Amb un glob comandes-*.jsonl per dia no canvia res al programa (spark.read.json accepta globs i directoris); només convé comprovar que cap dels fitxers no estigui encara en escriptura, que és exactament el que el sensor de 05-05 vigilarà amb _SUCCESS o amb un fitxer de tancament.
Conclusió
Spark pren el model de MapReduce (particions, tasques reexecutables, shuffle per clau) i li treu el que el feia lent: el càlcul sencer és un DAG d'operadors que el driver coneix complet, talla en stages només on hi ha una dependència ampla, encadena tota la resta a la mateixa task, i manté les dades intermèdies a la memòria d'executors que viuen tota l'aplicació. La tolerància a fallades ve del llinatge, que recomputa la partició perduda en lloc d'haver-la escrit a disc. Sobre els RDD, immutables i mandrosos, els DataFrames afegeixen esquema i un optimitzador, Catalyst, que empeny filtres, poda columnes, insereix agregacions parcials i tria broadcast o sort-merge per a cada join; Tungsten compila el pla, i Parquet fa que llegir dues columnes d'un any d'esdeveniments costi el que ocupen aquestes dues columnes. Hem après a llegir un explain() per verificar que el pla fa el que creiem, a fer servir cache quan hi ha diverses accions, broadcast per al catàleg, coalesce abans d'escriure, la sal per a Formatgeria Montblanc, i partitionBy amb sobreescriptura dinàmica perquè el job de vendes diàries sigui idempotent per dia. serveis/analitica/vendes_diaries.py calcula ara en sis segons, i en un sol programa, el que a 05-02 eren tres jobs i un minut; i ALS a MLlib ha convertit els clics del llac en recomanacions sense que les cent iteracions suposin cent lectures.
Tot el que hem fet parteix d'un conjunt acotat: el fitxer del 14 de setembre, ja tancat, o els clics de set dies. Però comandes.esdeveniments no es tanca mai: els esdeveniments estoc.actualitzat arriben a raó de centenars per segon, i les posicions dels repartidors a 2,4 milions al dia, i ni el panell de repartiment ni l'alerta d'estoc baix no poden esperar el lot de la nit. La lliçó següent tracta el processament de fluxos: què canvia quan el conjunt no acaba, com es defineix "els últims cinc minuts" quan els esdeveniments arriben desordenats (les marques d'aigua), i com Flink i Spark Structured Streaming executen el mateix DAG d'aquesta lliçó sobre un flux que no s'atura.
Curs d'Arquitectures Distribuïdes
Mòdul 1: Introducció als Sistemes Distribuïts
- Conceptes Bàsics de Sistemes Distribuïts
- Models de Sistemes Distribuïts
- Avantatges i Desafiaments dels Sistemes Distribuïts
- Les Fal·làcies de la Computació Distribuïda
- Temps, Rellotges i Ordenació d'Esdeveniments
- Del Monòlit a la Plataforma Distribuïda: el Cas Quilòmetre Zero
Mòdul 2: Comunicació en Sistemes Distribuïts
- Protocols de Comunicació
- RPC i RMI
- gRPC i Serialització de Dades
- Missatgeria i Cues de Missatges
- Patrons de Comunicació Asíncrona
Mòdul 3: Consistència i Replicació
- Models de Consistència
- El Teorema CAP i PACELC
- Algorismes de Consens
- Replicació de Dades
- Transaccions Distribuïdes i Sagues
Mòdul 4: Emmagatzematge Distribuït
- Particionament de Dades i Hashing Consistent
- Sistemes de Fitxers Distribuïts
- Emmagatzematge d'Objectes
- Bases de Dades Distribuïdes
- Memòries Cau Distribuïdes
Mòdul 5: Computació Distribuïda
- Models de Computació Distribuïda
- MapReduce i Hadoop
- Spark i Computació en Memòria
- Processament de Fluxos de Dades
- Planificació de Treballs i Pipelines de Dades
Mòdul 6: Seguretat en Sistemes Distribuïts
- Autenticació i Autorització
- Xifratge i Protecció de Dades
- Gestió d'Identitats
- Seguretat entre Serveis: mTLS i Gestió de Secrets
- Passarel·les d'API, Limitació de Taxa i Auditoria
Mòdul 7: Monitoratge i Manteniment
- Monitoratge de Sistemes Distribuïts
- Logs Centralitzats i Traçabilitat Distribuïda
- Gestió de Fallades i Recuperació
- Patrons de Resiliència: Timeouts, Reintents i Circuit Breaker
- Automatització i Orquestració
- Proves en Sistemes Distribuïts i Enginyeria del Caos
