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

  1. Per què Spark: un DAG en memòria davant de fases a disc
  2. Arquitectura: driver, cluster manager, executors, tasks i stages
  3. RDD: col·leccions immutables, transformacions mandroses i llinatge
  4. DataFrames i Spark SQL: Catalyst i formats columnars
  5. Operacions estretes i amples, shuffle i stages
  6. Optimitzacions: cache, broadcast join, particionament i biaix
  7. MapReduce davant de Spark
  8. MLlib: recomanacions amb ALS sobre els clics
  9. Pràctica: serveis/analitica/vendes_diaries.py
  10. Errors Comuns i Consells
  11. Exercicis
  12. Conclusió

  1. 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.

  1. 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é el SparkSession (abans SparkContext), 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.

  1. 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 : 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 res

Fins 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).

  1. 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:

  1. Anàlisi: resoldre noms de columnes i tipus contra l'esquema.
  2. 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.
  3. 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).
  4. 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.

  1. 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.

  1. 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.

  1. 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.

  1. 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.

  1. Pràctica: serveis/analitica/vendes_diaries.py

9.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/esdeveniments

Els 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,Lleida

9.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.json fa una passada completa només per inferir tipus. Amb ell, la lectura és una passada i els tipus són els que volem (data_ms com a LongType, no double).
  • explode converteix l'array linies en files, que és el que feia el bucle for ln in ev["dades"]["linies"] del mapper.
  • Un sol shuffle. El groupBy és l'única operació ampla. filter, select, explode i el join broadcast s'encadenen al mateix stage que la lectura. Catalyst insereix una agregació parcial abans del shuffle (HashAggregate amb partial_sum), així que el que travessa la xarxa són parcials per (dia, mercat, productor) per partició: el combiner, sense escriure'l.
  • cache() a linies perquè hi ha dues accions (control.first() i el write); sense ell, el JSON es llegiria i parsejaria dues vegades.
  • partitionBy("dia") amb partitionOverwriteMode=dynamic. La sortida és un directori per dia (dia=2026-09-14/), i overwrite en 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/_SUCCESS

Amb --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: OutOfMemoryError al driver. Per mirar, show() o take(n); per desar, write.
  • Oblidar que les transformacions són mandroses. Un filter mal escrit no falla en escriure'l, sinó a la primera acció, amb una traça que apunta al write. I un print dins 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. reduceByKey o aggregateByKey combinen localment. En DataFrames, groupBy().agg() ja ho fa.
  • UDF de Python per fila. Una udf a PySpark serialitza cada fila cap a un procés Python i de tornada; anul·la Catalyst i Tungsten. Gairebé tot es pot expressar amb F.*; si no, fes servir pandas_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 sense unpersist() 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; i coalesce abans 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 a explain().
  • Inferir l'esquema del JSON en producció. És una passada extra i un esquema que canvia amb les dades (un dia sense versio i la columna desapareix). Esquema explícit, versionat amb el contracte de l'esdeveniment (02-05).
  • Escriure amb overwrite sense partitionOverwriteMode=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.

  1. SortMergeJoin amb dos Exchange per productor en lloc de BroadcastHashJoin: el catàleg de tres files està provocant un shuffle de les vendes i una ordenació. Correcció: F.broadcast(cataleg).
  2. BatchEvalPython [calcular_import]: una UDF de Python per fila, que a més impedeix que el partial_sum es computi a la JVM sense sortir a Python. Correcció: F.col("ln.quantitat") * F.col("ln.preu").
  3. Exchange ... 200: 200 particions de shuffle per a un agregat de desenes de files; 200 tasks minúscules per stage. Correcció: spark.sql.shuffle.partitions a 8 (o AQE amb coalescència activada).
  4. FileScan json [id_esdeveniment, tipus, versio, data_ms, origen, dades] sense PushedFilters ni poda: es llegeixen totes les columnes, i probablement sense esquema explícit (inferència, una passada extra). Correcció: esquema declarat i select de les columnes necessàries just després de la lectura; el filtre tipus hauria d'aparèixer com a PushedFilters (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

Mòdul 2: Comunicació en Sistemes Distribuïts

Mòdul 3: Consistència i Replicació

Mòdul 4: Emmagatzematge Distribuït

Mòdul 5: Computació Distribuïda

Mòdul 6: Seguretat en Sistemes Distribuïts

Mòdul 7: Monitoratge i Manteniment

Mòdul 8: Casos d'Estudi i Aplicacions

© Copyright 2026. Tots els drets reservats