El Mòdul 4 va acabar amb una taula que deia on viu cada dada de Quilòmetre Zero: les comandes a Cassandra, l'estoc a PostgreSQL, el catàleg a Redis, les fotos a MinIO i, al llac de dades d'HDFS, un fitxer per dia amb els centenars de milers d'esdeveniments de comandes.esdeveniments i els milions de clics del web. Repartir les dades resolia el problema de desar-les; ara n'apareix el següent: calcular alguna cosa amb elles. Sumar les vendes de la Setmana de la Verema per productor i mercat sobre els esdeveniments de set dies, o entrenar recomanacions amb els clics d'un any, no cap en una màquina, i encara que hi cabés trigaria hores. Aquesta lliçó explica què canvia quan el que es distribueix és un càlcul i no només una dada: per què convé portar el codi fins on són les dades, quines maneres hi ha de repartir la feina (scatter/gather, divideix i venceràs, cues de treball, BSP, dataflow, actors), en què es diferencien els lots, els fluxos i les consultes interactives, i quins problemes nous apareixen (el cost de moure dades entre nodes, el biaix, els ressagats, la reexecució després d'una fallada, el límit de l'escalat). Tot el que ve després en el mòdul (MapReduce, Spark, Flink, Airflow) és una implementació concreta d'aquestes idees, així que val la pena entendre-les primer amb codi Python d'una sola màquina, que és el que farem a simulacions/.

Contingut

  1. Distribuir un càlcul: portar el codi a les dades
  2. Paral·lelisme de dades i paral·lelisme de tasques
  3. Patrons de computació distribuïda
  4. Lots, fluxos i interactiu
  5. El cost de la comunicació: el shuffle
  6. Biaix de dades
  7. Ressagats i execució especulativa
  8. Tolerància a fallades per reexecució determinista
  9. Escalabilitat: Amdahl i Gustafson
  10. Pràctica: scatter/gather, cua de treball i biaix a simulacions/
  11. Errors Comuns i Consells
  12. Exercicis
  13. Conclusió

  1. Distribuir un càlcul: portar el codi a les dades

Quan el catàleg cabia en una taula de PostgreSQL, "calcular les vendes per productor" era una consulta SQL: el motor llegia les files del seu disc local, les agrupava i retornava vint números. La dada i el càlcul eren a la mateixa màquina. Al llac de dades de 04-02 ja no és així: /km0/esdeveniments/2026-09-14/comandes.jsonl fa 150 MB en dos blocs que viuen en tres DataNodes diferents, i els set dies de la Setmana de la Verema sumen més d'un gigabyte repartit per tot el clúster. Hi ha dues maneres de calcular sobre això:

  • Portar les dades al codi. Un procés a analitica descarrega el gigabyte per la xarxa, el recorre i suma. Funciona, però la xarxa és el recurs més lent que tenim (la fal·làcia 3 de 01-04, "l'ample de banda és infinit"): a 1 Gbit/s, moure 1 GB són uns 10 segons només de transferència, i amb un any de clics (terabytes) l'enfocament és senzillament inviable. A més, el procés que suma és un de sol, així que el càlcul no escala.
  • Portar el codi a les dades (data locality). Enviar a cada DataNode un programa petit (uns quilobytes) que llegeixi el bloc que ja té al seu disc local, calculi un resultat parcial (uns quilobytes: vendes per productor d'aquell bloc) i retorni només això. Es mouen el codi i els resultats, que són minúsculs; les dades, que són enormes, no es mouen.

Aquesta inversió és la idea central de tota la computació distribuïda moderna i la raó per la qual HDFS i MapReduce van néixer junts: el sistema de fitxers exposa on és cada bloc (hdfs fsck -locations ho mostrava a 04-02) precisament perquè el planificador de càlcul pugui executar cada tasca al node que té el seu bloc, o com a mínim al mateix rack. Quan la localitat no és possible (el node està ocupat, o les dades són en un emmagatzematge d'objectes com MinIO que no exposa localitat), es paga la xarxa, i per això les plataformes modernes separen còmput i emmagatzematge però ho compensen amb xarxes de 25–100 Gbit/s i formats columnars que llegeixen només les columnes necessàries (05-03).

El segon canvi és que el càlcul es converteix en moltes tasques independents més una fase de combinació. En comptes d'un programa que ho recorre tot, escrivim una funció que processa un tros i una altra que combina resultats parcials. Aquesta descomposició no és gratuïta: hi ha càlculs que es descomponen de manera natural (sumar vendes per productor) i d'altres que no (ordenar tots els esdeveniments per import, calcular una mediana exacta, recórrer un graf de "clients que van comprar el mateix"). Els patrons de l'apartat 3 són les maneres conegudes de descompondre.

  1. Paral·lelisme de dades i paral·lelisme de tasques

Hi ha dues maneres de repartir una feina entre diversos nodes, i convé distingir-les perquè exigeixen infraestructures diferents:

Paral·lelisme de dades Paral·lelisme de tasques
Què es reparteix Les dades: cada node executa el mateix codi sobre un tros diferent Les tasques: cada node executa codi diferent sobre les mateixes dades o sobre dades relacionades
Exemple a Quilòmetre Zero Sumar vendes per productor: 8 nodes, cadascun amb un vuitè de comandes.jsonl Per a una comanda: un calcula l'import, un altre valida l'estoc, un altre estima el repartiment; els tres alhora
Com escala Amb la mida de les dades: el doble de dades, el doble de nodes, el mateix temps Amb el nombre de tasques diferents, que sol ser petit i fix
Coordinació Al final, per combinar resultats parcials Entre tasques, per dependències (B necessita el resultat d'A)
Dificultat principal Repartir de manera equilibrada; combinar sense perdre informació Sincronitzar i gestionar dependències; el camí crític limita
Models que l'exploten MapReduce, Spark, Flink, SQL distribuït Pipelines de tasques (Airflow, 05-05), microserveis que col·laboren (Mòdul 8)

L'analítica de Quilòmetre Zero és gairebé tota paral·lelisme de dades, i per això aquest mòdul s'hi centra. Però tots dos es combinen: el pipeline diari de 05-05 és paral·lelisme de tasques (validar, agregar, carregar, notificar) en què la tasca "agregar" és, per dins, paral·lelisme de dades a Spark. I la llei que governa l'escalat (apartat 9) és diferent en cada cas: el paral·lelisme de dades s'acosta a l'escalat lineal perquè la fracció seqüencial és petita; el de tasques està limitat per la cadena de dependències més llarga.

  1. Patrons de computació distribuïda

3.1 Scatter/gather (dispersar i reunir)

És el patró més simple i el que ja va aparèixer, sense nom, a 04-01 en parlar d'índexs secundaris a Cassandra: un coordinador reparteix la feina entre N treballadors (scatter), cadascun calcula sobre la seva part, i el coordinador reuneix i combina els resultats parcials (gather).

flowchart LR
    C[Coordinador<br/>analitica] -- tros 1 --> T1[Treballador 1<br/>vendes parcials]
    C -- tros 2 --> T2[Treballador 2<br/>vendes parcials]
    C -- tros 3 --> T3[Treballador 3<br/>vendes parcials]
    C -- tros 4 --> T4[Treballador 4<br/>vendes parcials]
    T1 --> G[Gather:<br/>sumar parcials]
    T2 --> G
    T3 --> G
    T4 --> G
    G --> R[vendes per productor]

Funciona quan el càlcul és associatiu i commutatiu: sumar, comptar, màxim, mínim, unió de conjunts. Tant se val en quin ordre i amb quina agrupació es combinin els parcials, el resultat és el mateix. Una mitjana no és directament combinable (la mitjana de mitjanes és incorrecta si els trossos tenen mides diferents), però es torna combinable portant (suma, compte) en lloc de la mitjana; una mediana exacta no s'hi torna, i cal aproximar-la o pagar una ordenació global. És la primera pregunta que cal fer-se davant de qualsevol càlcul distribuït: quin resultat parcial retorna cada tros, i com es combinen dos parcials?

El límit de l'scatter/gather és que el coordinador és únic: reparteix, espera tothom i combina. Si combinar és costós (milions de claus diferents) o els parcials són grans, el coordinador es converteix en el coll d'ampolla. MapReduce (05-02) resol exactament això distribuint també la fase de combinació.

3.2 Divideix i venceràs distribuït

És l'scatter/gather aplicat recursivament: un problema es parteix en subproblemes, cada subproblema es torna a partir fins que cap en un node, i els resultats es combinen pujant per l'arbre. Una ordenació distribuïda funciona així: cada node ordena el seu tros (mergesort local), i les combinacions successives barregen llistes ordenades. Els frameworks ho fan servir per a operacions de reducció amb molts nodes: en comptes que un coordinador sumi 1 000 parcials, se sumen en arbre (treeReduce a Spark), amb 1 000 → 32 → 1 combinacions en paral·lel.

3.3 Cua de treball amb treballadors (work queue)

A l'scatter/gather el coordinador decideix per endavant quin tros va a cada treballador. Amb una cua de treball no decideix: publica totes les tasques en una cua i cada treballador agafa la següent quan acaba l'anterior. És el patró dels consumidors competidors de 02-04 aplicat al càlcul, i té dos avantatges grans:

  • Equilibri dinàmic. Si un treballador és lent (màquina vella, tasca gran), simplement agafa menys tasques; els ràpids absorbeixen la resta. Ningú no espera ociós.
  • Tolerància a fallades natural. Si un treballador mor a mitja tasca, la tasca torna a la cua (pel lease o l'ack pendent) i un altre l'executa. Perquè això sigui correcte, la tasca ha de ser idempotent: executar-la dues vegades ha de donar el mateix resultat que una (apartat 8).

És el model intern de gairebé tots els planificadors: YARN, Spark i Flink mantenen cues de tasques pendents i les assignen als executors que queden lliures, amb preferència pel node que té les dades. Ho veurem a simulacions/cua_treball.py.

3.4 Bulk Synchronous Parallel (BSP)

Hi ha càlculs que no es fan en una passada, sinó en iteracions en què cada pas depèn de l'anterior: el PageRank d'un graf, l'entrenament d'un model per descens de gradient, la propagació de "clients que van comprar el mateix" pel graf de comandes per a les recomanacions de Quilòmetre Zero. El model BSP (Valiant, 1990) organitza aquests càlculs en supersteps (superpassos):

  1. Càlcul local: cada node treballa només amb les seves dades i els missatges que va rebre al superpas anterior.
  2. Comunicació: cada node envia missatges als altres (per exemple, un vèrtex del graf envia el seu valor als seus veïns).
  3. Barrera: ningú no comença el superpas següent fins que tothom ha acabat l'actual i tots els missatges han arribat.
flowchart TB
    subgraph S1[Superpas 1]
        direction LR
        A1[Node A<br/>calcula] --> M1[missatges]
        B1[Node B<br/>calcula] --> M1
        C1[Node C<br/>calcula] --> M1
    end
    M1 --> BAR1{{Barrera: tothom ha acabat}}
    BAR1 --> S2
    subgraph S2[Superpas 2]
        direction LR
        A2[Node A<br/>calcula] --> M2[missatges]
        B2[Node B<br/>calcula] --> M2
        C2[Node C<br/>calcula] --> M2
    end
    M2 --> BAR2{{Barrera}}
    BAR2 --> FIN[... fins a convergir]

La barrera és el que fa el model fàcil de raonar (dins d'un superpas no hi ha curses: cada node només veu missatges del superpas anterior) i també el que el fa sensible als ressagats (apartat 7): el superpas dura el que duri el node més lent. Pregel de Google, i els seus descendents Apache Giraph i GraphX de Spark, són BSP "pensant com un vèrtex": cada vèrtex del graf rep missatges, actualitza el seu valor i envia missatges als seus veïns, superpas rere superpas fins que cap vèrtex no canvia. A Quilòmetre Zero, "productes que solen comprar-se junts" és un graf on els vèrtexs són productes i les arestes pesen pel nombre de comandes compartides; dos o tres superpassos de propagació basten per trobar que formatge-curat i vi-crianca són més a prop del que suggereixen les seves categories.

3.5 Pipeline / dataflow: el DAG d'operadors

El model de flux de dades (dataflow) descriu el càlcul com un graf dirigit acíclic (DAG) d'operadors: llegir, filtrar, transformar, agrupar, unir, escriure. Cada operador rep dades dels anteriors i n'emet als següents; les dades flueixen per les arestes. El sistema decideix com paral·lelitzar cada operador (quantes instàncies, en quins nodes), com encadenar operadors que no necessiten redistribuir dades (filtrar i transformar poden anar al mateix procés, fila a fila), i on cal redistribuir (agrupar per productor obliga que totes les files d'un productor arribin a la mateixa instància: és el shuffle de l'apartat 5).

flowchart LR
    L[llegir comandes.jsonl] --> F[filtrar tipus = comanda.creada]
    F --> E[explotar línies de la comanda]
    E --> S[[shuffle per productor i mercat]]
    S --> A[sumar import]
    C[llegir catàleg] --> J
    A --> J[unir amb catàleg]
    J --> W[escriure Parquet]

És el model de Spark (05-03) i de Flink (05-04), i també, amb un altre vocabulari, el dels motors SQL distribuïts i el dels planificadors de pipelines (05-05, on els nodes del DAG són treballs sencers en comptes d'operadors). Davant de MapReduce, que obliga a expressar-ho tot com a parelles map/reduce encadenades, el dataflow deixa que el programador escrigui la transformació completa i que l'optimitzador decideixi les fases. Davant de BSP, el dataflow no té barreres globals: un operador processa tan bon punt té dades, i només el shuffle sincronitza.

3.6 Actors

El model d'actors (Hewitt, 1973) pren un altre camí: no reparteix dades ni descriu un graf, sinó que modela el sistema com molts objectes petits (actors) que només es comuniquen per missatges asíncrons, cadascun amb el seu estat privat i la seva bústia, processant un missatge cada vegada. No hi ha memòria compartida ni bloqueigs: un actor rep un missatge, canvia el seu estat, envia missatges a altres actors o crea actors nous. Erlang ho porta al llenguatge des dels anys 80 (les centraletes d'Ericsson; RabbitMQ, de 02-04, està escrit en Erlang) i Akka ho va portar a la JVM. Encaixa de manera natural amb l'estat per entitat: un actor per repartidor de furgoneta-3 que rep les seves posicions i manté la seva ruta, o un actor per comanda que executa la seva saga (03-05). És menys adequat per al càlcul massiu sobre dades històriques, que és el que ocupa aquest mòdul, així que el deixem anomenat: el seu lloc torna a aparèixer quan l'estat per clau és el protagonista, i de fet els operadors amb estat de Flink (05-04) s'assemblen molt a actors particionats per clau.

3.7 Resum de patrons

Patró Repartiment Coordinació Encaixa quan Exemple
Scatter/gather Estàtic, pel coordinador Una vegada, al final Càlcul associatiu, pocs parcials Vendes per productor d'un dia
Divideix i venceràs Recursiu En arbre Combinació costosa, molts nodes Ordenar tots els esdeveniments per import
Cua de treball Dinàmic, per demanda Cap entre treballadors Tasques de durada desigual, fallades freqüents Redimensionar 100 000 fotos de MinIO
BSP Per vèrtex/partició Barrera per superpas Iteratiu, grafs Productes comprats junts
Dataflow (DAG) Per operador i partició Només al shuffle Transformacions encadenades, lots o fluxos Pipeline de vendes diàries, panell de repartiment
Actors Per entitat Missatges asíncrons Estat per entitat, concurrència Un actor per repartidor

  1. Lots, fluxos i interactiu

Independentment del patró, hi ha tres modes de processar segons quan arriben les dades i quan es necessita la resposta:

Per lots (batch) Per fluxos (streaming) Interactiu (ad hoc)
Entrada Un conjunt acotat i complet: "els esdeveniments del 14 de setembre" Un flux no acotat que no acaba: comandes.esdeveniments a Kafka Un conjunt acotat, però la pregunta es decideix al moment
Quan es calcula Programat: cada nit, cada hora Contínuament, esdeveniment a esdeveniment o en microlots Quan algú pregunta
Latència esperada Minuts a hores Mil·lisegons a segons Segons
Mida de dades per execució Gigabytes a petabytes Quilobytes per esdeveniment; milions d'esdeveniments per hora Gigabytes, amb índexs o formats columnars per anar ràpid
Resultat Complet i exacte sobre el conjunt Aproximat o provisional, es refina en arribar més dades (05-04) Exacte sobre el que hi ha
Tolerància a fallades Reexecutar el lot sencer Checkpoints d'estat i offsets Reexecutar la consulta
Quilòmetre Zero Vendes per productor/mercat/dia; recomanacions entrenades cada nit amb els clics Panell de repartiment amb posicions de furgoneta-3; alerta d'estoc baix "Quant va vendre Celler Roure Alt a Lleida durant la Setmana de la Verema?" des de la consola d'analitica
Eines MapReduce (05-02), Spark (05-03) Kafka Streams, Flink, Spark Structured Streaming (05-04) Spark SQL, Presto/Trino, Hive anomenat (05-02)

Els tres modes comparteixen els patrons de l'apartat 3 i els problemes dels apartats 5–8, però cadascun els pateix de manera diferent: el biaix en un lot allarga una nit, en un flux embussa una partició per sempre. I la frontera entre lots i fluxos és més difusa del que sembla: un lot de "els esdeveniments del 14 de setembre" és un flux al qual s'han posat principi i fi; un flux processat en microlots d'un segon són lots molt petits. Aquesta idea, que un motor pot tractar tots dos modes amb el mateix DAG, és la que Spark i Flink exploten i la que 05-04 desenvoluparà.

  1. El cost de la comunicació: el shuffle

Distribuir un càlcul té un cost que no existeix en una màquina: moure dades entre nodes quan el pas següent necessita agrupar-les d'una altra manera. Sumar vendes per productor exigeix que totes les línies de Formatgeria Montblanc acabin al mateix node, i aquestes línies estan repartides per tots els blocs de tots els nodes. L'operació que les ajunta s'anomena shuffle (barrejar), i és amb diferència la fase més cara de qualsevol treball distribuït, perquè:

  • Implica tots-a-tots: cada node envia una part de les seves dades a cadascun dels altres. Amb N nodes són N² fluxos de xarxa.
  • Sol passar per disc: les dades se serialitzen, s'escriuen ordenades per clau de destí, es transfereixen i es tornen a llegir. A MapReduce cada shuffle és una escriptura completa a disc (05-02); Spark el manté en memòria quan pot (05-03).
  • No es pot solapar del tot amb el càlcul: el receptor necessita totes les seves dades abans d'agrupar, així que el shuffle és una barrera implícita.

La regla pràctica és minimitzar el que travessa el shuffle: reduir abans de moure. Si cada node suma localment les seves línies per productor abans d'enviar-les (un combiner, en vocabulari de 05-02; una agregació parcial, a Spark), en lloc de moure 250 000 línies mou 20 parcials per node. Aquesta és la raó d'insistir que el càlcul sigui associatiu: només llavors es pot prereduir. I explica el consell de 04-01 sobre les taules materialitzades per mercat i productor: algú havia pagat el shuffle una vegada, en el moment d'escriure, per no pagar-lo a cada consulta.

Un càlcul que no necessita shuffle (filtrar, transformar fila a fila, sumar un total global que es combina en un sol número) és vergonyosament paral·lel (embarrassingly parallel) i escala gairebé linealment. Un que necessita diversos shuffles encadenats (agrupar per productor, unir amb el catàleg, reagrupar per mercat) està limitat per la xarxa i pel pitjor dels apartats següents.

  1. Biaix de dades

El repartiment ideal dona a cada node la mateixa quantitat de feina. El repartiment real depèn de la distribució de les claus, i les distribucions reals són desiguals. Durant la "Setmana del Formatge Artesà" Formatgeria Montblanc concentra la meitat de les línies de comanda de Quilòmetre Zero; si el shuffle agrupa per productor, el node que rep formatgeria-montblanc processa la meitat de les dades mentre els altres set es reparteixen l'altra meitat. La feina triga el que triga aquell node: amb 8 nodes i una clau que pesa el 50 %, el speedup màxim és 2, no 8. És el biaix de dades (data skew), i és la causa més freqüent de treballs distribuïts que "no escalen encara que afegeixi màquines".

Es detecta mirant la durada de les tasques d'una mateixa fase: si la majoria acaba en 20 s i una triga 3 min, hi ha una clau calenta. Les solucions, totes elles maneres de trencar la clau grossa, es tracten a la pràctica (apartat 10.3) i es reprendran a 05-03 amb el nom que els dona Spark, salting:

  • Prereduir abans del shuffle: si cada node suma les seves línies de Formatgeria Montblanc, el que travessa la xarxa és un parcial per node, no la meitat de les dades.
  • Repartir la clau calenta afegint-hi un sufix aleatori (formatgeria-montblanc#0#7), agrupar per la clau amb sufix i tornar a agrupar els vuit parcials en un segon pas molt més petit.
  • Triar una altra clau de partició quan el càlcul ho permet: agrupar per (productor, mercat) reparteix Formatgeria Montblanc entre quatre mercats.
  • Tractar a part les claus calentes conegudes (filtrar-les, processar-les amb més paral·lelisme) i unir després.

El biaix també existeix a l'entrada: si el fitxer del 14 de setembre són dos blocs i el del 15 en són deu, els treballs per dia tindran durades molt diferents. I existeix en el temps: els esdeveniments de Kafka es concentren al migdia i a última hora de la tarda, així que un flux particionat per hora té hores grosses.

  1. Ressagats i execució especulativa

Encara que el repartiment sigui perfecte, alguna tasca acabarà trigant molt més que les altres sense que les dades ho justifiquin: la màquina té un disc degradat, un altre treball competeix per la seva CPU, la xarxa del seu rack està saturada, la JVM és en una pausa de recollida d'escombraries. Són els ressagats (stragglers). En un treball de 1 000 tasques la probabilitat que alguna caigui en una màquina amb problemes és alta, i com que el treball acaba quan acaba l'última tasca, un sol ressagat allarga tot el treball. A BSP l'efecte es multiplica pel nombre de superpassos.

La solució clàssica, introduïda per MapReduce i present a Spark (spark.speculation) i a Hadoop, és l'execució especulativa: quan una fase està gairebé acabada i una tasca porta molt més temps que la mediana de les altres, el planificador llança una còpia d'aquella tasca en un altre node; la primera que acaba guanya, i l'altra es cancel·la. És una despesa (s'executa feina redundant) que compensa perquè el cost d'una tasca extra és molt menor que el de tot el clúster esperant. Només és possible perquè les tasques són deterministes i idempotents, que és el tema de l'apartat següent: si la còpia i l'original escrivissin totes dues el seu resultat, caldria garantir que se n'escriu un de sol (sortida atòmica, 05-02).

  1. Tolerància a fallades per reexecució determinista

En una màquina, si el programa falla a la meitat, es rellança des del principi. Amb 1 000 nodes i un treball de tres hores, la probabilitat que algun node falli durant el treball és pràcticament 1 (fal·làcia 1 de 01-04), i rellançar-ho tot cada vegada és inacceptable. L'estratègia dels frameworks d'aquest mòdul és la reexecució de gra fi: si falla una tasca, es reexecuta aquella tasca, no el treball. Perquè això sigui correcte calen tres propietats que ja coneixem de 02-05:

  1. Entrada immutable. La tasca llegeix un tros de dades que no canvia (un bloc d'HDFS, un rang d'offsets de Kafka, una partició d'un RDD). Reexecutar llegeix el mateix.
  2. Càlcul determinista. La mateixa entrada produeix la mateixa sortida. Sense números aleatoris sense llavor, sense dependre de l'hora actual ni de l'ordre d'arribada per la xarxa. Si cal aleatorietat, es deriva de la clau o d'una llavor fixa.
  3. Sortida idempotent o atòmica. O bé escriure dues vegades és innocu (un PUT d'un objecte amb el mateix nom a MinIO, un INSERT ... ON CONFLICT DO NOTHING amb l'id de la tasca), o bé la sortida s'escriu en un lloc temporal i es publica de cop en acabar (reanomenar un fitxer, confirmar una transacció). Així una tasca que va morir a mitges no deixa una sortida parcial que confongui la seva reexecució ni l'execució especulativa.

Amb aquestes tres propietats, la fallada d'un node es redueix a "les seves tasques tornen a la cua", que és el que ja feia la cua de treball de l'apartat 3.3. El que canvia amb cada framework és com reconstrueix l'entrada d'una tasca quan aquella entrada era el resultat d'una altra tasca anterior: MapReduce l'escriu sempre a disc (HDFS o local), Spark recorda com es va calcular i la recomputa (el llinatge de 05-03), Flink desa checkpoints periòdics de l'estat (05-04). I el que no canvia és que tot descansa sobre tasques idempotents: el consumidor idempotent de 02-05 i la tasca de Spark són la mateixa idea a escales diferents.

  1. Escalabilitat: Amdahl i Gustafson

A 01-03 vam veure la llei d'Amdahl: si una fracció 1 − p de la feina és seqüencial, el speedup amb N nodes està acotat per 1 / (1 − p) per molts nodes que s'hi afegeixin. En un treball distribuït la part seqüencial és la que no es reparteix: llegir la llista de blocs, planificar tasques, el gather final, escriure el resultat en un sol fitxer, i sobretot el shuffle i l'espera dels ressagats, que encara que s'executin en paral·lel es comporten com una barrera. Amb p = 0,95 el sostre és 20 encara que es facin servir 1 000 nodes, cosa que sembla dir que la computació distribuïda a gran escala no compensa.

La resposta de John Gustafson (1988) és que Amdahl suposa un problema de mida fixa, i no és així com es fan servir els clústers: ningú no compra 1 000 nodes per calcular més ràpid les vendes d'un dia, sinó per calcular les vendes d'un any, o de tots els clics, en el mateix temps. La llei de Gustafson mesura el speedup escalat: si amb N nodes el treball triga un temps T del qual una fracció s és seqüencial i 1 − s és paral·lela, aquest mateix treball en un sol node hauria trigat s + (1 − s) · N vegades T, així que:

Speedup_escalat(N) = N − s · (N − 1)

Amb s = 0,05 i N = 1 000, el speedup escalat és 950: gairebé lineal, perquè en créixer el problema la fracció seqüencial (planificar, reunir 20 números) es manté mentre que la paral·lela (llegir terabytes) creix amb N. Les dues lleis són certes, mesuren coses diferents, i juntes donen la regla de disseny d'aquest mòdul:

Amdahl Gustafson
Suposa Mida del problema fixa Temps fix, problema que creix amb N
Pregunta Quant més ràpid acabo el mateix? Quant més processo en el mateix temps?
Fórmula 1 / ((1 − p) + p / N) N − s · (N − 1)
Lliçó per a Quilòmetre Zero Un dia de vendes no s'accelera amb 100 nodes: el gather i l'arrencada dominen Un any de clics es processa en la mateixa nit amb 100 nodes que un dia amb 1
Com millorar Reduir la part seqüencial: prereduir, evitar el gather únic Mantenir la part seqüencial constant en créixer les dades

La conseqüència pràctica: el paral·lelisme compensa quan la feina per node és gran davant del cost fix d'arrencar i coordinar. Llançar 1 000 tasques de 100 ms cadascuna en un clúster amb 2 s de latència de planificació és pitjor que 10 tasques de 10 s. Ho comprovarem mesurant a la pràctica.

  1. Pràctica: scatter/gather, cua de treball i biaix a simulacions/

Tota la pràctica s'executa en una màquina amb multiprocessing, que reparteix feina entre processos igual que un clúster la reparteix entre nodes, amb la diferència que la "xarxa" és memòria local. És suficient per veure els patrons, mesurar speedups i reproduir el biaix. El fitxer d'entrada és una versió petita de /km0/esdeveniments/2026-09-14/comandes.jsonl, amb l'embolcall d'esdeveniments de 02-05 i un esdeveniment comanda.creada per línia:

{"id_esdeveniment":"e-000123-1","tipus":"comanda.creada","versio":1,"data_ms":1789380000000,"origen":"comandes","dades":{"comanda_id":"P-2026-000123","client":"anna","mercat":"girona","linies":[{"producte":"tomaquet-rosa","productor":"horta-la-vega","quantitat":2,"preu":3.90},{"producte":"formatge-curat","productor":"formatgeria-montblanc","quantitat":1,"preu":12.50}]}}
{"id_esdeveniment":"e-000124-1","tipus":"comanda.creada","versio":1,"data_ms":1789380045000,"origen":"comandes","dades":{"comanda_id":"P-2026-000124","client":"marc","mercat":"lleida","linies":[{"producte":"vi-crianca","productor":"celler-roure-alt","quantitat":6,"preu":9.80}]}}
{"id_esdeveniment":"e-000125-1","tipus":"comanda.creada","versio":1,"data_ms":1789380090000,"origen":"comandes","dades":{"comanda_id":"P-2026-000125","client":"llucia","mercat":"valencia","linies":[{"producte":"formatge-fresc","productor":"formatgeria-montblanc","quantitat":3,"preu":4.20},{"producte":"carbasso","productor":"horta-la-vega","quantitat":4,"preu":1.60}]}}

Aquestes tres línies basten per llegir el codi; per mesurar alguna cosa cal més volum, així que el primer script inclou un generador que crea centenars de milers d'esdeveniments amb la mateixa forma i una distribució esbiaixada cap a Formatgeria Montblanc.

10.1 simulacions/vendes_scatter_gather.py

# km0/simulacions/vendes_scatter_gather.py
"""Scatter/gather local: reparteix comandes.jsonl entre N treballadors i suma vendes per productor.

Ús:   python vendes_scatter_gather.py generar 400000      # crea esdeveniments/2026-09-14/comandes.jsonl
      python vendes_scatter_gather.py calcular 1 2 4 8     # mesura amb 1, 2, 4 i 8 treballadors
"""
import json, os, random, sys, time
from collections import Counter
from multiprocessing import Pool

RUTA = "esdeveniments/2026-09-14/comandes.jsonl"
PRODUCTES = [                      # (producte, productor, preu, pes en la distribució)
    ("tomaquet-rosa",  "horta-la-vega",         3.90, 15),
    ("carbasso",       "horta-la-vega",         1.60, 10),
    ("formatge-curat", "formatgeria-montblanc", 12.50, 35),   # Setmana del Formatge Artesà:
    ("formatge-fresc", "formatgeria-montblanc", 4.20, 15),    # Montblanc concentra el 50 %
    ("vi-crianca",     "celler-roure-alt",      9.80, 25),
]
MERCATS = ["girona", "lleida", "tarragona", "valencia"]
CLIENTS = ["anna", "marc", "llucia"]


def generar(n_comandes: int) -> None:
    """Escriu n_comandes esdeveniments comanda.creada amb una distribució esbiaixada de productors."""
    random.seed(42)                                   # determinista: mateix fitxer a cada execució
    os.makedirs(os.path.dirname(RUTA), exist_ok=True)
    pesos = [p[3] for p in PRODUCTES]
    with open(RUTA, "w", encoding="utf-8") as f:
        for i in range(n_comandes):
            linies = [
                {"producte": prod, "productor": productor, "quantitat": random.randint(1, 6), "preu": preu}
                for prod, productor, preu, _ in random.choices(PRODUCTES, weights=pesos, k=random.randint(1, 3))
            ]
            esdeveniment = {
                "id_esdeveniment": f"e-{i:06d}-1", "tipus": "comanda.creada", "versio": 1,
                "data_ms": 1789344000000 + i * 200, "origen": "comandes",
                "dades": {"comanda_id": f"P-2026-{i:06d}", "client": random.choice(CLIENTS),
                          "mercat": random.choice(MERCATS), "linies": linies},
            }
            f.write(json.dumps(esdeveniment) + "\n")


def trossos_per_bytes(ruta: str, n: int) -> list[tuple[int, int]]:
    """Divideix el fitxer en n rangs [inici, fi) alineats amb salts de línia.

    Imita el que fa HDFS amb els blocs: cada treballador rep un rang de bytes,
    no una llista de línies, per no haver de llegir tot el fitxer al coordinador.
    """
    mida = os.path.getsize(ruta)
    talls = [0]
    with open(ruta, "rb") as f:
        for k in range(1, n):
            f.seek(mida * k // n)                     # salt aproximat
            f.readline()                              # avança fins al final de la línia partida
            talls.append(f.tell())
    talls.append(mida)
    return [(talls[i], talls[i + 1]) for i in range(n)]


def mapar_tros(rang: tuple[int, int]) -> Counter:
    """Fase 'scatter': un treballador llegeix el seu rang i retorna vendes parcials per productor."""
    inici, fi = rang
    parcial = Counter()
    with open(RUTA, "rb") as f:
        f.seek(inici)
        while f.tell() < fi:
            linia = f.readline()
            if not linia:
                break
            ev = json.loads(linia)
            if ev["tipus"] != "comanda.creada":
                continue
            for ln in ev["dades"]["linies"]:
                parcial[ln["productor"]] += ln["quantitat"] * ln["preu"]
    return parcial


def reduir(parcials: list[Counter]) -> Counter:
    """Fase 'gather': combinar parcials. Sumar és associatiu i commutatiu, l'ordre no importa."""
    total = Counter()
    for p in parcials:
        total.update(p)
    return total


def calcular(n_treballadors: int) -> tuple[Counter, float]:
    t0 = time.perf_counter()
    rangs = trossos_per_bytes(RUTA, n_treballadors)
    if n_treballadors == 1:
        parcials = [mapar_tros(rangs[0])]             # sense Pool: evita el cost d'arrencar processos
    else:
        with Pool(n_treballadors) as pool:
            parcials = pool.map(mapar_tros, rangs)
    total = reduir(parcials)
    return total, time.perf_counter() - t0


if __name__ == "__main__":
    if sys.argv[1] == "generar":
        generar(int(sys.argv[2]))
        print(f"Generat {RUTA}: {os.path.getsize(RUTA) / 1e6:.1f} MB")
    else:
        base = None
        for n in map(int, sys.argv[2:]):
            total, seg = calcular(n)
            base = base or seg
            print(f"{n:2d} treballadors: {seg:6.2f} s  speedup {base / seg:4.2f}x  "
                  f"Montblanc = {total['formatgeria-montblanc']:,.2f} €")

Punts que convé entendre del codi:

  • El coordinador no llegeix les dades. trossos_per_bytes només calcula rangs de bytes, amb un seek i un readline per alinear cada tall amb el principi d'una línia (sense això, una línia quedaria partida entre dos treballadors i tots dos la descartarien o la comptarien malament). És exactament el que fa Hadoop amb els input splits i el motiu pel qual JSON Lines o CSV són "divisibles" mentre que un JSON amb un array gegant no ho és.
  • Cada treballador obre el fitxer pel seu compte i llegeix només el seu rang. En un clúster, aquell treballador correria al node que té el bloc (localitat); aquí tots comparteixen disc.
  • El parcial és petit: un Counter amb tres claus, independentment que el tros tingui mil o un milió de línies. El que travessa la "xarxa" (la cua interna de Pool) són tres números per treballador.
  • reduir és trivial perquè sumar és associatiu. Si el càlcul fos "import mitjà per comanda", el parcial hauria de ser (suma, compte) per productor i la reducció dividir al final.

Una execució en un portàtil de 8 nuclis amb 400 000 comandes (uns 130 MB):

$ python vendes_scatter_gather.py generar 400000
Generat esdeveniments/2026-09-14/comandes.jsonl: 131.6 MB
$ python vendes_scatter_gather.py calcular 1 2 4 8 16
 1 treballadors:   6.84 s  speedup 1.00x  Montblanc = 1,893,412.30 €
 2 treballadors:   3.61 s  speedup 1.89x  Montblanc = 1,893,412.30 €
 4 treballadors:   1.96 s  speedup 3.49x  Montblanc = 1,893,412.30 €
 8 treballadors:   1.18 s  speedup 5.80x  Montblanc = 1,893,412.30 €
16 treballadors:   1.14 s  speedup 6.00x  Montblanc = 1,893,412.30 €

El resultat és idèntic amb qualsevol nombre de treballadors (la reducció no depèn del repartiment) i el speedup s'allunya de l'ideal a mesura que creixen els processos: amb 8 n'obtenim 5,8, i amb 16 (més processos que nuclis) res. És Amdahl en acció: arrencar el Pool costa uns 100 ms fixos, el gather i l'alineació de trossos són seqüencials, i el disc és compartit. Amb un fitxer deu vegades més gran la fracció seqüencial es dilueix i el speedup amb 8 s'acosta a 7,5: Gustafson.

10.2 simulacions/cua_treball.py

El segon script converteix el càlcul en tasques idempotents en una cua, amb treballadors que les agafen per demanda i un que mor a mitja tasca. Cada tasca és "calcular les vendes per productor d'un mercat i un dia", i la seva sortida és un fitxer sortida/<dia>-<mercat>.json escrit de manera atòmica.

# km0/simulacions/cua_treball.py
"""Cua de treball amb treballadors que competeixen per tasques idempotents.

Un treballador mor a propòsit a mitja tasca; el coordinador detecta la mort,
retorna la tasca a la cua i llança un treballador substitut. El resultat final és el mateix.
"""
import json, os, sys, time
from collections import Counter
from multiprocessing import Process, Queue

RUTA = "esdeveniments/2026-09-14/comandes.jsonl"
SORTIDA = "sortida"
MERCATS = ["girona", "lleida", "tarragona", "valencia"]


def executar_tasca(tasca: dict) -> Counter:
    """Vendes per productor d'un mercat. Determinista: mateixa entrada, mateix resultat."""
    parcial = Counter()
    with open(RUTA, encoding="utf-8") as f:
        for linia in f:
            ev = json.loads(linia)
            if ev["tipus"] == "comanda.creada" and ev["dades"]["mercat"] == tasca["mercat"]:
                for ln in ev["dades"]["linies"]:
                    parcial[ln["productor"]] += ln["quantitat"] * ln["preu"]
    return parcial


def escriure_atomic(ruta: str, contingut: dict) -> None:
    """Escriu en un temporal i reanomena: ningú no veu mai un fitxer a mitges."""
    tmp = f"{ruta}.tmp-{os.getpid()}"
    with open(tmp, "w", encoding="utf-8") as f:
        json.dump(contingut, f)
    os.replace(tmp, ruta)                             # rename atòmic a POSIX


def treballador(nom: str, pendents: Queue, esdeveniments: Queue, morir_a: str | None) -> None:
    while True:
        tasca = pendents.get()
        if tasca is None:                             # senyal de fi
            return
        esdeveniments.put(("inici", nom, tasca["id"]))
        if morir_a == tasca["id"]:
            time.sleep(0.2)
            os._exit(1)                               # mort sobtada: sense excepció, sense neteja
        resultat = executar_tasca(tasca)
        escriure_atomic(f"{SORTIDA}/{tasca['dia']}-{tasca['mercat']}.json", dict(resultat))
        esdeveniments.put(("fi", nom, tasca["id"]))


def coordinador(n_treballadors: int) -> None:
    os.makedirs(SORTIDA, exist_ok=True)
    pendents, esdeveniments = Queue(), Queue()
    tasques = {f"2026-09-14/{m}": {"id": f"2026-09-14/{m}", "dia": "2026-09-14", "mercat": m} for m in MERCATS}
    for t in tasques.values():
        pendents.put(t)
    en_curs: dict[str, str] = {}                      # treballador -> id de tasca
    acabades: set[str] = set()
    processos: dict[str, Process] = {}

    def llancar(nom, morir_a=None):
        p = Process(target=treballador, args=(nom, pendents, esdeveniments, morir_a), daemon=True)
        p.start()
        processos[nom] = p

    for i in range(n_treballadors):
        llancar(f"w{i}", morir_a="2026-09-14/lleida" if i == 1 else None)   # w1 morirà a lleida

    while len(acabades) < len(tasques):
        while not esdeveniments.empty():
            tipus, nom, id_tasca = esdeveniments.get()
            if tipus == "inici":
                en_curs[nom] = id_tasca
                print(f"[coord] {nom} comença {id_tasca}")
            else:
                acabades.add(id_tasca); en_curs.pop(nom, None)
                print(f"[coord] {nom} acaba {id_tasca}")
        for nom, p in list(processos.items()):        # detecció de fallades: el procés ja no viu
            if not p.is_alive() and nom in en_curs:
                perduda = en_curs.pop(nom)
                print(f"[coord] {nom} ha mort amb {perduda}: la torno a encuar i llanço un substitut")
                pendents.put(tasques[perduda])        # la tasca torna a la cua, intacta
                del processos[nom]
                llancar(nom + "'")
        time.sleep(0.05)

    for _ in processos:
        pendents.put(None)
    total = Counter()
    for m in MERCATS:
        with open(f"{SORTIDA}/2026-09-14-{m}.json", encoding="utf-8") as f:
            total.update(json.load(f))
    print("Total per productor:", {k: round(v, 2) for k, v in total.items()})


if __name__ == "__main__":
    coordinador(int(sys.argv[1]) if len(sys.argv) > 1 else 3)
$ python cua_treball.py 3
[coord] w0 comença 2026-09-14/girona
[coord] w1 comença 2026-09-14/lleida
[coord] w2 comença 2026-09-14/tarragona
[coord] w1 ha mort amb 2026-09-14/lleida: la torno a encuar i llanço un substitut
[coord] w0 acaba 2026-09-14/girona
[coord] w0 comença 2026-09-14/valencia
[coord] w1' comença 2026-09-14/lleida
[coord] w2 acaba 2026-09-14/tarragona
[coord] w0 acaba 2026-09-14/valencia
[coord] w1' acaba 2026-09-14/lleida
Total per productor: {'horta-la-vega': 712338.1, 'formatgeria-montblanc': 1893412.3, 'celler-roure-alt': 984501.6}

El que ensenya l'execució:

  • Equilibri dinàmic: w0 acaba Girona i agafa València sense que ningú l'hi assigni; amb mercats de mida desigual, els treballadors ràpids absorbeixen més tasques.
  • Detecció i reexecució: el coordinador no rep cap missatge d'error de w1 (va morir amb os._exit, com una màquina que s'apaga); ho detecta perquè el procés ja no és viu, igual que YARN detecta un NodeManager per heartbeats perduts. Torna a encuar la tasca tal qual i llança un substitut (en una altra execució pot ser w2, si queda lliure abans, qui agafi Lleida: la cua no assigna, els treballadors competeixen).
  • Idempotència: w1 va morir després de començar i potser amb el fitxer temporal a mitges; w1' executa la mateixa tasca, escriu el seu propi temporal i el reanomena. El temporal orfe de w1 queda a sortida/ sense que ningú no el llegeixi (un treball real el netejaria). Si la sortida hagués estat un INSERT a PostgreSQL sense clau única, la reexecució hauria duplicat files: és el mateix problema del consumidor de 02-05.
  • El resultat és el mateix que el de l'scatter/gather, cosa que és la definició que la fallada ha estat tolerada.

10.3 Biaix: quan una partició triga el doble

El tercer experiment reutilitza vendes_scatter_gather.py però canvia el repartiment: en lloc de trossos de bytes, reparteix per productor, que és el que faria un shuffle ingenu "agrupar per productor". Afegeix a l'script aquesta funció i la crida:

# A nivell de mòdul: Pool ha de poder serialitzar (pickle) la funció, i una funció niada no ho permet.
def mapar_productor(productor):
    t0 = time.perf_counter(); total = 0.0
    with open(RUTA, encoding="utf-8") as f:
        for linia in f:
            ev = json.loads(linia)
            for ln in ev["dades"]["linies"]:
                if ln["productor"] == productor:
                    total += ln["quantitat"] * ln["preu"]
    return productor, total, time.perf_counter() - t0

def calcular_per_productor() -> None:
    """Repartiment per clau: cada treballador processa un productor. Reprodueix el biaix d'un shuffle."""
    productors = ["horta-la-vega", "formatgeria-montblanc", "celler-roure-alt"]
    with Pool(3) as pool:
        for productor, total, seg in pool.map(mapar_productor, productors):
            print(f"{productor:22s} {total:14,.2f} €  {seg:5.2f} s")

(Perquè el biaix es vegi en el temps, i no només en la quantitat de dades, cada treballador en un shuffle real només rebria les seves línies; en aquesta simulació cadascun filtra el fitxer complet, així que afegeix al bucle interior una petita feina proporcional a les línies pròpies, per exemple hashlib.md5(linia).hexdigest() només quan ln["productor"] == productor.) La sortida mostra el problema:

horta-la-vega              712,338.10 €   1.71 s
formatgeria-montblanc    1,893,412.30 €   3.52 s
celler-roure-alt           984,501.60 €   1.93 s

Tres treballadors, i la feina dura 3,52 s: el que triga Formatgeria Montblanc. Els altres dos estan ociosos la meitat del temps. La correcció és repartir la clau calenta en subclaus i fer una segona reducció:

# A nivell de mòdul: Pool ha de poder serialitzar (pickle) la funció, i una funció niada no ho permet.
def mapar_subclau(clau):
    productor, sub, n_sub = clau; total = 0.0
    with open(RUTA, encoding="utf-8") as f:
        for linia in f:
            ev = json.loads(linia)
            # la subclau es deriva de l'id de comanda: determinista, reparteix uniformement
            if productor == "formatgeria-montblanc" and hash(ev["dades"]["comanda_id"]) % n_sub != sub:
                continue
            for ln in ev["dades"]["linies"]:
                if ln["productor"] == productor:
                    total += ln["quantitat"] * ln["preu"]
    return productor, total

def calcular_amb_salting(n_sub: int = 4) -> None:
    """Trenca la clau calenta en n_sub subclaus i torna a combinar. Dues fases de reducció."""
    claus = [("horta-la-vega", 0, 1), ("celler-roure-alt", 0, 1)] + \
            [("formatgeria-montblanc", k, n_sub) for k in range(n_sub)]     # 2 + 4 = 6 tasques
    with Pool(6) as pool:
        parcials = pool.map(mapar_subclau, claus)
    total = Counter()
    for productor, import_eur in parcials:            # segona reducció: 6 números, trivial
        total[productor] += import_eur
    print(dict(total))

Ara la tasca més llarga processa un vuitè de les dades en comptes de la meitat, i la feina sencera baixa a poc més d'1 s amb sis processos. Fixa't que el hash() de Python està aleatoritzat per procés per a cadenes (PYTHONHASHSEED), així que perquè la subclau sigui determinista entre execucions i reexecucions cal fixar la llavor o fer servir zlib.crc32(comanda_id.encode()) % n_sub; és un exemple petit de com s'esmuny el no determinisme de l'apartat 8. Spark farà això mateix amb dos groupBy i una columna de sal a 05-03.

Errors Comuns i Consells

  • Moure les dades al càlcul per costum. El primer instint de qui ve del monòlit és "descarrego el fitxer i el processo". Amb gigabytes funciona malament i amb terabytes no funciona. Pregunta sempre on són les dades i si el càlcul hi pot anar.
  • Dissenyar un resultat parcial que no es combina. Mitjanes, percentils, diferents (count distinct) i medianes no se sumen. Porta (suma, compte), fes servir estructures aproximades (HyperLogLog per a diferents, t-digest per a percentils) o accepta un shuffle complet.
  • Ignorar el shuffle. Un groupBy sobre una clau d'alta cardinalitat és una operació tots-a-tots. Preredueix abans de moure, tria claus amb cardinalitat raonable i mira els bytes de shuffle a la interfície del framework (05-02 i 05-03 els mostren).
  • Culpar el clúster del biaix. Quan "no escala", mira la durada de les tasques de la fase lenta. Si una triga cinc vegades la mediana, no falten màquines: sobra una clau.
  • Tasques no idempotents. Un INSERT sense clau, un comptador incrementat, un append a un fitxer: qualsevol d'ells converteix la reexecució (i l'execució especulativa) en un duplicat. Escriu a un temporal i reanomena, fes servir claus naturals, o escriu per partició completa amb sobreescriptura (05-05).
  • No determinisme amagat. El hash() de Python amb llavor aleatòria, datetime.now(), random sense llavor, l'ordre d'un dict d'una versió antiga, l'ordre d'arribada de missatges. La reexecució d'una tasca ha de produir bytes idèntics.
  • Moltes tasques diminutes. El planificador té un cost per tasca (mil·lisegons a segons). 100 000 tasques de 50 ms són pitjors que 1 000 de 5 s. Com a orientació, entre 2 i 4 tasques per nucli i fase, amb durades de segons a minuts.
  • Extrapolar Amdahl a Gustafson (o a l'inrevés). Si el problema és fix, els nodes extra no ajuden; si el problema creix, sí. Abans de demanar més màquines, decideix quin dels dos casos és el teu.

Exercicis

Exercici 1: Què es combina i què no?

Per a cadascun d'aquests càlculs sobre els esdeveniments comanda.creada de la Setmana de la Verema, indica (a) què retorna cada treballador com a resultat parcial i (b) com es combinen dos parcials, o per què no és possible combinar-los i quina alternativa hi ha:

  1. Import total venut per Celler Roure Alt.
  2. Import mitjà per comanda a cada mercat.
  3. Nombre de clients diferents que van comprar vi-crianca.
  4. Els 10 productes més venuts per quantitat.
  5. La mediana de l'import de les comandes.

Exercici 2: Reexecució i sortida atòmica

A cua_treball.py, substitueix escriure_atomic per una escriptura directa (open(ruta, "w") i json.dump) i fes que el treballador mori durant l'escriptura (per exemple, escriu la meitat del JSON, fes f.flush() i os._exit(1)). Descriu què passa al coordinador al final, com ho detectaries en producció i per què l'os.replace ho evita. Després, proposa com escriuries la sortida si en lloc d'un fitxer fos una taula PostgreSQL vendes_mercat_dia(dia, mercat, productor, import_eur) de manera que la reexecució continués sent segura.

Exercici 3: Amdahl, Gustafson i la mida de la tasca

Amb els temps de l'execució de l'apartat 10.1 (1 treballador: 6,84 s; 8 treballadors: 1,18 s), estima la fracció seqüencial de l'script segons Amdahl. Amb aquesta fracció, quin speedup obtindries amb 64 treballadors sobre el mateix fitxer? I quin seria el speedup escalat de Gustafson amb 64 treballadors si el fitxer creixés 64 vegades? Per acabar: el planificador d'un clúster real afegeix 1,5 s per tasca entre assignació i arrencada; si el fitxer de 130 MB es divideix en 1 000 trossos, quant triga el treball amb 8 nodes, i en quants trossos l'hauries de dividir?

Solucions

Exercici 1.

  1. Parcial: un número (suma de quantitat × preu de les línies de Celler Roure Alt). Combinació: suma. Associatiu i commutatiu; el cas ideal.
  2. Parcial: per mercat, la parella (suma_imports, nombre_comandes). Combinació: sumar component a component; la mitjana es calcula només al final, suma / compte. Combinar mitjanes directament seria incorrecte tret que tots els trossos tinguessin el mateix nombre de comandes.
  3. Parcial: el conjunt d'ids de client que van comprar vi-crianca en aquell tros. Combinació: unió de conjunts, i al final la mida. És combinable però el parcial pot ser gran (centenars de milers d'ids); si això és un problema, un HyperLogLog per tros (uns KB) es combina amb una unió i dona el cardinal amb un error de l'1–2 %.
  4. Parcial: quantitats per producte (un Counter), no "els 10 millors del tros": l'onzè d'un tros pot ser el primer global. Combinació: sumar els Counter i triar els 10 al final. Si el nombre de productes fos enorme, es podria desar el top-K per tros amb K gran com a aproximació, acceptant error.
  5. No és combinable: la mediana de medianes no és la mediana. Alternatives: ordenar globalment (un shuffle per rangs d'import: car però exacte), o un t-digest/percentil aproximat per tros, que sí que es combina i dona la mediana amb error acotat.

Exercici 2.

Amb escriptura directa, w1 deixa sortida/2026-09-14-lleida.json amb mig JSON. El coordinador torna a encuar i w1' torna a obrir el fitxer en mode "w", que el trunca, així que en aquest cas concret el resultat final és correcte; però entre la mort i la reexecució (segons, o minuts en un clúster) el fitxer existeix i està corrupte: qualsevol lector (el pipeline de 05-05, un hdfs dfs -cat) fallaria amb un JSONDecodeError, o pitjor, llegiria un parcial si el format fos CSV. Si el treballador morís després d'escriure però abans d'enviar fi, la reexecució també sobreescriuria, sense dany. El problema real apareix si la reexecució no trunca (mode "a") o si la sortida és un sistema sense truncament. En producció es detecta per lectors que fallen o per fitxers amb mida inesperada; os.replace ho evita perquè el fitxer amb el nom definitiu apareix de cop i complet o no apareix: és la sortida atòmica de l'apartat 8 i la que MapReduce implementa amb directoris _temporary (05-02).

Per a PostgreSQL: una clau primària (dia, mercat, productor) i escriptura amb INSERT ... ON CONFLICT (dia, mercat, productor) DO UPDATE SET import_eur = EXCLUDED.import_eur, tot dins d'una transacció que comença amb DELETE FROM vendes_mercat_dia WHERE dia = %s AND mercat = %s i acaba amb COMMIT: la tasca reemplaça la seva partició sencera de manera atòmica, i executar-la dues vegades deixa exactament les mateixes files. És el "reprocessar un dia sense duplicar" de 05-05.

Exercici 3.

Amdahl: S(8) = 1 / ((1 − p) + p / 8) = 6,84 / 1,18 = 5,80. Aïllant, (1 − p) + p / 8 = 1 / 5,80 = 0,17241 − 0,875 p = 0,1724p = 0,946. Fracció seqüencial 1 − p ≈ 5,4 %. Amb 64 treballadors: S(64) = 1 / (0,054 + 0,946 / 64) = 1 / 0,0688 = 14,5: a la pràctica, menys, perquè el portàtil té 8 nuclis i el disc és un de sol. Gustafson amb s = 0,054 i N = 64: 64 − 0,054 × 63 = 60,6: processar 64 vegades més dades amb 64 nodes en gairebé el mateix temps.

Amb 1 000 trossos i 8 nodes, cada tros són 130 KB (uns 7 ms de procés) més 1,5 s de cost de planificació: 1 000 tasques × 1,507 s / 8 nodes ≈ 188 s. Pitjor que els 6,84 s d'un sol procés. Amb 8 trossos: 8 × (0,86 s + 1,5 s) / 8 ≈ 2,4 s. Amb 16 o 24 trossos (2–3 per node) l'equilibri dinàmic ajuda amb els ressagats sense disparar el cost fix: uns 2,5–3 s. La regla: la feina per tasca ha de ser com a mínim un ordre de magnitud més gran que el cost de planificar-la.

Conclusió

Distribuir un càlcul és més que executar-lo en diverses màquines: és portar el codi on són les dades, descompondre la feina en tasques els resultats parcials de les quals es puguin combinar, i acceptar que la combinació (el shuffle) és la part cara. Hem separat el paral·lelisme de dades, que és el d'aquest mòdul, del paral·lelisme de tasques, que és el dels pipelines; i hem recorregut els patrons amb què s'organitza el repartiment: scatter/gather per a l'associatiu, divideix i venceràs per combinar en arbre, la cua de treball per equilibrar i tolerar fallades per demanda, BSP amb els seus superpassos i barreres per a l'iteratiu i els grafs, el DAG d'operadors del dataflow que Spark i Flink implementen, i els actors per a l'estat per entitat. Els tres modes (lots, fluxos, interactiu) comparteixen patrons i problemes: el biaix que fa que Formatgeria Montblanc allargui tot el treball, els ressagats que l'execució especulativa esquiva, i la tolerància a fallades que només funciona si les tasques són deterministes i idempotents amb sortida atòmica, la mateixa regla que governava els consumidors de 02-05. Amdahl ens va recordar que un problema fix té un sostre, i Gustafson que els clústers existeixen per a problemes que creixen. A simulacions/ hem mesurat un speedup de 5,8 amb 8 processos, hem vist morir i renéixer un treballador sense alterar el resultat, i hem trencat la clau calenta de Formatgeria Montblanc en subclaus.

Tot això ho hem fet a mà, amb multiprocessing i un coordinador de quaranta línies que reparteix, detecta morts i torna a encuar. Un framework de computació distribuïda és exactament això, però per a milers de nodes, amb localitat de dades, shuffle distribuït, execució especulativa i sortida atòmica resolts d'una vegada per a tots els treballs. El primer que ho va aconseguir, i el que va fixar el vocabulari que continuem fent servir, va ser MapReduce, i amb ell Hadoop, que és la lliçó següent: com el mateix càlcul de vendes per productor s'expressa com a map, shuffle & sort i reduce, i com YARN reparteix les tasques pel clúster.

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