El Mòdul 3 va acabar amb una frase que val la pena rellegir: tot el que sabem sobre replicar, coordinar i compensar "tracta les dades com si ja fossin en algun lloc". Aquesta lliçó obre el Mòdul 4 responent la pregunta prèvia: quan les comandes de Quilòmetre Zero, les posicions de furgoneta-3 o els esdeveniments de comandes.esdeveniments són massa per a un sol node, com es reparteixen entre molts? Replicar copia les mateixes dades en diversos nodes per sobreviure a fallades; particionar (o sharding) reparteix dades diferents entre nodes perquè cap no hagi de desar-ho o servir-ho tot. Veurem les estratègies de repartiment (per rang, per hash, compostes), el problema que apareix quan el nombre de nodes canvia, la solució clàssica del hashing consistent amb nodes virtuals, el seu parent rendezvous hashing, i com una petició troba el node correcte. Tot plegat és la base sobre la qual es construeixen HDFS, Cassandra i Redis Cluster, que veurem a la resta del mòdul, i és també una decisió de disseny que Quilòmetre Zero ha de prendre avui: la clau de partició de les seves comandes i del seu estoc.
Contingut
- Particionar enfront de replicar, i com es combinen
- Particionament per rang de clau i els punts calents
- Particionament per hash de clau i particionament compost
- Índexs secundaris particionats: local enfront de global
- El problema del hash mòdul N
- Hashing consistent: l'anell i els nodes virtuals
- Rendezvous hashing
- Reequilibratge de particions
- Encaminament de peticions i descobriment de particions
- Triar la clau de partició a Quilòmetre Zero
- Errors Comuns i Consells
- Exercicis
- Conclusió
- Particionar enfront de replicar, i com es combinen
Els dos mecanismes responen a problemes diferents i es confonen sovint perquè a la pràctica apareixen junts:
| Aspecte | Replicació (03-04) | Particionament (aquesta lliçó) |
|---|---|---|
| Què copia | Les mateixes dades en diversos nodes | Dades diferents a cada node |
| Problema que resol | Disponibilitat, tolerància a fallades, lectures properes | Volum i rendiment: cap node no aguanta tot el conjunt |
| Si cau un node | Un altre node en té una còpia | Es perd la seva partició… llevat que estigui replicada |
| Cost principal | Consistència entre còpies | Repartiment desigual, consultes que creuen particions |
| Pregunta clau | Quantes còpies i amb quines garanties? | Quina dada va a quin node? |
En un sistema real cada partició es replica: el conjunt de dades es divideix en particions P1…Pn, i cada partició té un líder i diversos seguidors (o un quòrum sense líder). Un node físic sol ser líder d'algunes particions i seguidor d'altres, de manera que la càrrega i el risc es reparteixen:
flowchart LR
subgraph Node A
A1[P1 líder]
A2[P2 seguidor]
A3[P4 seguidor]
end
subgraph Node B
B1[P2 líder]
B2[P3 seguidor]
B3[P1 seguidor]
end
subgraph Node C
C1[P3 líder]
C2[P4 líder]
C3[P2 seguidor]
end
subgraph Node D
D1[P4 seguidor]
D2[P1 seguidor]
D3[P3 seguidor]
end
L'objectiu del particionament és que les dades i la càrrega es reparteixin de manera uniforme. Si el repartiment és desigual, algunes particions reben molta més càrrega que d'altres: són els punts calents (hot spots), i una partició calenta converteix un sistema de 10 nodes en un sistema amb la capacitat d'1. Tot el que segueix gira al voltant de dues preguntes: com assignar claus a particions per evitar punts calents, i com assignar particions a nodes perquè afegir o treure màquines no obligui a moure totes les dades.
- Particionament per rang de clau i els punts calents
L'estratègia més intuïtiva és ordenar les claus i assignar a cada partició un rang contigu, com els volums d'una enciclopèdia (A–C, D–F…). Els límits no han de ser necessàriament regulars: s'adapten a la densitat de dades perquè cada partició tingui una mida semblant.
Avantatge: les consultes per rang són eficients, perquè les claus veïnes són al mateix node. Si la clau de les comandes és la data, "totes les comandes de la setmana passada" es resol en una o dues particions.
Inconvenient: els rangs concentren càrrega quan el patró d'accés es concentra en una zona de claus. Quilòmetre Zero ho va viure a la "Setmana de la Verema": amb les comandes particionades per data_creacio, totes les escriptures del dia van a la mateixa partició (la del rang que conté "avui"), mentre les particions de dies anteriors resten ocioses. Deu nodes i només un escrivint.
| Clau de rang | Consulta que afavoreix | Punt calent |
|---|---|---|
data_creacio |
Comandes per període | Escriptures sempre a la partició d'"avui" |
client_id alfabètic |
Comandes d'un client | Clients amb nom que comenci per lletres freqüents |
producte_slug |
Comandes d'un producte | vi-crianca a la Setmana de la Verema |
Un remei habitual és anteposar a la clau alguna cosa que dispersi: per exemple, mercat + data (Girona, Lleida, Tarragona, València reparteixen les escriptures d'avui en quatre particions). Continua havent-hi quatre punts calents, però repartits; i la consulta "comandes d'aquesta setmana a València" continua sent un rang. Això anticipa el particionament compost de l'apartat següent.
- Particionament per hash de clau i particionament compost
Per eliminar els punts calents per proximitat de claus, s'aplica una funció hash a la clau i es particiona pel resultat. Un bon hash (MD5, SHA-1, MurmurHash, xxHash; no cal que sigui criptogràfic, sí que sigui uniforme i estable entre llenguatges i versions) converteix claus semblants (P-2026-000123 i P-2026-000124) en valors molt allunyats, de manera que les comandes consecutives de la Setmana de la Verema cauen en particions diferents.
import hashlib
def hash_clau(clau: str) -> int:
"""Hash estable de 64 bits d'una clau de text (els primers 8 bytes de MD5)."""
return int.from_bytes(hashlib.md5(clau.encode()).digest()[:8], "big")
for comanda in ["P-2026-000123", "P-2026-000124", "P-2026-000125", "P-2026-000126"]:
print(comanda, hash_clau(comanda) % 4)Cada línia mostra la comanda i la partició (0–3) que li toca: comandes correlatives acaben disperses. És important no fer servir hash() de Python per a això: des de Python 3.3 està aleatoritzat per procés per a cadenes, així que dos serveis diferents obtindrien particions diferents per a la mateixa clau. Per això fem servir hashlib.
El preu del hash és perdre les consultes per rang: "comandes entre el 10 i el 14 de setembre" ja no és a cap partició concreta, cal preguntar a totes (scatter/gather, apartat 4).
El particionament compost combina tots dos: una part de la clau es hasheja per triar la partició, i la resta es fa servir per ordenar dins de la partició. Cassandra ho fa explícit amb la partition key i les clustering columns (04-04), però la idea és general:
- Clau de partició:
hash(client_id)→ totes les comandes de l'Anna són a la mateixa partició. - Clau d'ordenació:
data_creacio DESC→ dins d'aquesta partició estan ordenades, així que "les últimes 20 comandes de l'Anna" és una lectura seqüencial en un sol node.
La consulta "comandes de l'Anna entre dues dates" és eficient; "comandes de tots els clients d'ahir" no ho és. Això és l'essencial del modelatge orientat a consultes: la clau de partició es tria segons la consulta que més importa, no segons el model de domini.
- Índexs secundaris particionats: local enfront de global
Les comandes es particionen per client_id, però l'equip de repartiment necessita "totes les comandes pendents de lliurament a Girona", que no esmenta cap client. Cal un índex secundari (per mercat i estat), i un índex en un sistema particionat també cal particionar-lo. Hi ha dues maneres:
| Índex local (per document) | Índex global (per terme) | |
|---|---|---|
| On viu | Cada partició indexa les seves pròpies dades | L'índex es particiona pel valor indexat (mercat = Girona viu en una partició concreta) |
| Escriptura | Només toca la partició de la dada: ràpida, atòmica amb la dada | Toca la partició de la dada i la de l'índex: més lenta, sovint asíncrona (eventual) |
| Lectura per l'índex | Cal preguntar a totes les particions i unir resultats: scatter/gather | Una sola partició respon |
| Latència de lectura | La de la partició més lenta (tail latency) | Baixa i predictible |
| Exemples | Cassandra (índexs secundaris locals), Elasticsearch per defecte | DynamoDB GSI, vistes materialitzades de Cassandra |
El scatter/gather mereix atenció: si hi ha 12 particions, la consulta llança 12 subconsultes en paral·lel i espera la més lenta. Amb una probabilitat de l'1 % que una partició trigui més de 500 ms, la probabilitat que la consulta completa superi aquest temps és 1 − 0,99¹² ≈ 11 %. És el mateix fenomen que va fer tan importants els percentils alts a 01-03.
Per a repartiment, Quilòmetre Zero tria un índex global asíncron, materialitzat a partir dels esdeveniments comanda.creada i pagament.confirmat del tòpic comandes.esdeveniments (02-05): una taula comandes_per_mercat particionada per mercat, que pot anar uns mil·lisegons endarrerida. L'assignació de repartidors tolera aquest retard; la creació de la comanda no toleraria una escriptura síncrona en dues particions.
- El problema del hash mòdul N
Amb el hash a la mà, l'assignació més òbvia de claus a nodes és node = hash(clau) % N. Funciona perfectament fins al dia que N canvia. En passar de 4 a 5 nodes, hash % 4 i hash % 5 coincideixen només per a les claus el hash de les quals dona el mateix residu en tots dos casos, i això passa aproximadament en 1 de cada 5 claus: el 80 % de les dades canvien de node. En un magatzem amb terabytes això vol dir hores de trànsit de xarxa, memòries cau fredes i, si és una memòria cau, una allau sobre la base de dades.
Ho mesurem amb simulacions/hash_modul.py, que genera 100 000 identificadors de comanda amb el format de Quilòmetre Zero:
# km0/simulacions/hash_modul.py
"""Quantes claus de comanda canvien de node en passar de 4 a 5 nodes amb hash % N?"""
import hashlib
def hash_clau(clau: str) -> int:
return int.from_bytes(hashlib.md5(clau.encode()).digest()[:8], "big")
def node_modul(clau: str, n_nodes: int) -> int:
return hash_clau(clau) % n_nodes
def mesurar_moviment(claus, abans: int, despres: int) -> float:
"""Fracció de claus el node de les quals canvia en passar d'`abans` a `despres` nodes."""
mogudes = sum(1 for c in claus if node_modul(c, abans) != node_modul(c, despres))
return mogudes / len(claus)
if __name__ == "__main__":
claus = [f"P-2026-{i:06d}" for i in range(1, 100_001)]
for abans, despres in [(4, 5), (5, 6), (10, 11), (4, 8)]:
pct = mesurar_moviment(claus, abans, despres) * 100
print(f"{abans:>2} -> {despres:>2} nodes: es mou el {pct:5.1f} % de les claus")Explicació pas a pas:
hash_claués la mateixa funció estable de l'apartat 3.node_modulés l'assignació ingènua: el residu de dividir el hash entre el nombre de nodes.mesurar_movimentcompara, clau a clau, el node abans i després, i compta les que canvien.- El bloc principal prova diversos canvis de mida del clúster.
Sortida típica:
4 -> 5 nodes: es mou el 80.0 % de les claus 5 -> 6 nodes: es mou el 83.3 % de les claus 10 -> 11 nodes: es mou el 91.0 % de les claus 4 -> 8 nodes: es mou el 49.7 % de les claus
L'ideal seria moure únicament l'imprescindible: en afegir el cinquè node a quatre, només 1/5 de les claus (les que van a parar al nou). Fixa't en el cas 4 → 8: duplicar el nombre de nodes mou "només" la meitat perquè hash % 8 conserva el bit baix de hash % 4; és un truc que fan servir alguns sistemes (créixer per duplicació), però és rígid i continua movent més del que cal.
- Hashing consistent: l'anell i els nodes virtuals
El hashing consistent (Karger et al., 1997, pensat originalment per a memòries cau web distribuïdes) resol el problema desacoblant l'assignació del nombre de nodes. La idea:
- L'espai de sortida del hash (per exemple 0…2⁶⁴−1) s'imagina com un anell: el valor màxim va seguit del 0.
- Cada node també es hasheja (pel seu nom o adreça) i ocupa una posició a l'anell.
- Una clau s'assigna al primer node que es troba recorrent l'anell en el sentit de les agulles del rellotge des de la posició de la clau.
flowchart TB
subgraph Anell["Anell de hash (0 … 2^64−1, en sentit horari)"]
direction LR
N1(("inv-bcn<br/>pos 0x1A…"))
N2(("inv-vlc<br/>pos 0x6F…"))
N3(("inv-gir<br/>pos 0xB3…"))
N1 --> N2 --> N3 --> N1
end
K1["formatge-curat<br/>hash 0x4C… → inv-vlc"] -.-> N2
K2["tomaquet-rosa<br/>hash 0x9E… → inv-gir"] -.-> N3
K3["vi-crianca<br/>hash 0xE1… → inv-bcn (fa la volta)"] -.-> N1
En afegir un node en una posició de l'anell, només les claus entre el seu predecessor i ell canvien de propietari (passen del successor al node nou). En treure un node, les seves claus passen al seu successor i res més no es mou. Amb N nodes, afegir-ne un mou de mitjana 1/(N+1) de les claus: exactament el mínim.
El problema de l'anell bàsic és que, amb pocs nodes, les seves posicions aleatòries reparteixen molt malament l'espai: un node pot quedar-se amb el 50 % de l'anell i un altre amb el 5 %. I en treure un node, tota la seva càrrega cau sobre un únic successor. La solució són els nodes virtuals (vnodes o tokens): cada node físic es col·loca a l'anell K vegades (inv-bcn#0, inv-bcn#1, …, inv-bcn#149), amb posicions diferents. Amb 100–200 vnodes per node, els arcs es promitgen, el repartiment s'acosta a l'uniforme, i la càrrega d'un node que cau es reparteix entre molts successors diferents. A més, permeten donar més vnodes a les màquines més potents (pesos).
simulacions/anell_consistent.py ho implementa i ho mesura:
# km0/simulacions/anell_consistent.py
"""Anell de hashing consistent amb nodes virtuals."""
import bisect
import hashlib
import statistics
from collections import Counter
def hash_clau(clau: str) -> int:
return int.from_bytes(hashlib.md5(clau.encode()).digest()[:8], "big")
class AnellConsistent:
def __init__(self, nodes=(), vnodes: int = 150):
self.vnodes = vnodes
self._posicions: list[int] = [] # posicions ordenades a l'anell
self._propietari: dict[int, str] = {} # posició -> nom de node físic
for n in nodes:
self.afegir_node(n)
def afegir_node(self, node: str) -> None:
for i in range(self.vnodes):
pos = hash_clau(f"{node}#{i}")
if pos in self._propietari: # col·lisió raríssima: s'ignora el vnode
continue
bisect.insort(self._posicions, pos)
self._propietari[pos] = node
def treure_node(self, node: str) -> None:
for i in range(self.vnodes):
pos = hash_clau(f"{node}#{i}")
if self._propietari.get(pos) == node:
del self._propietari[pos]
self._posicions.remove(pos)
def node_per(self, clau: str) -> str:
if not self._posicions:
raise RuntimeError("anell buit")
h = hash_clau(clau)
idx = bisect.bisect_right(self._posicions, h) # primer vnode a la dreta
if idx == len(self._posicions): # passat el final: fa la volta
idx = 0
return self._propietari[self._posicions[idx]]
def distribucio(anell: AnellConsistent, claus) -> Counter:
return Counter(anell.node_per(c) for c in claus)
def desviacio_relativa(recompte: Counter) -> float:
"""Desviació típica del nombre de claus per node, relativa a la mitjana (en %)."""
valors = list(recompte.values())
return statistics.pstdev(valors) / statistics.mean(valors) * 100
def fraccio_moguda(abans: AnellConsistent, despres: AnellConsistent, claus) -> float:
return sum(1 for c in claus if abans.node_per(c) != despres.node_per(c)) / len(claus)
if __name__ == "__main__":
claus = [f"P-2026-{i:06d}" for i in range(1, 100_001)]
nodes = ["inv-bcn", "inv-vlc", "inv-gir", "inv-lle"]
for vn in (1, 10, 150):
anell = AnellConsistent(nodes, vnodes=vn)
recompte = distribucio(anell, claus)
print(f"vnodes={vn:>3}: {dict(recompte)} desviació={desviacio_relativa(recompte):.1f} %")
abans = AnellConsistent(nodes, vnodes=150)
despres = AnellConsistent(nodes + ["inv-tar"], vnodes=150)
print(f"afegir inv-tar (4->5): es mou el {fraccio_moguda(abans, despres, claus)*100:.1f} %")
sense_vlc = AnellConsistent(nodes, vnodes=150)
sense_vlc.treure_node("inv-vlc")
print(f"treure inv-vlc (4->3): es mou el {fraccio_moguda(abans, sense_vlc, claus)*100:.1f} %")
print("repartiment després de treure inv-vlc:", dict(distribucio(sense_vlc, claus)))Com funciona, línia a línia:
_posicionsés una llista ordenada d'enters (les posicions de tots els vnodes) i_propietaridiu a quin node físic pertany cada posició. Mantenir la llista ordenada ambbisect.insortpermet cercar el successor en O(log V).afegir_nodegenera K posicions hashejantnom#i; els vnodes d'un mateix node queden dispersos per tot l'anell.treure_nodeelimina exactament aquestes posicions; no toca res més, així que les claus dels altres nodes no es veuen afectades.node_perhasheja la clau, cerca ambbisect_rightla primera posició més gran que el hash (el veí "a la dreta") i si se surt pel final torna a la posició 0: això és l'anell.desviacio_relativaresumeix com d'uniforme és el repartiment: 0 % seria perfecte.
Sortida representativa (els valors exactes depenen dels hashos):
vnodes= 1: {'inv-lle': 70171, 'inv-gir': 13875, 'inv-bcn': 14486, 'inv-vlc': 1468} desviació=106.4 %
vnodes= 10: {'inv-bcn': 32315, 'inv-gir': 32959, 'inv-vlc': 10845, 'inv-lle': 23881} desviació=35.7 %
vnodes=150: {'inv-bcn': 24275, 'inv-gir': 27571, 'inv-vlc': 24835, 'inv-lle': 23319} desviació=6.3 %
afegir inv-tar (4->5): es mou el 19.2 %
treure inv-vlc (4->3): es mou el 24.8 %
repartiment després de treure inv-vlc: {'inv-bcn': 32671, 'inv-gir': 37988, 'inv-lle': 29341}Tres conclusions. Primera: sense vnodes, inv-lle s'endú el 70 % de les claus i inv-vlc l'1,5 % (un repartiment inservible, i cada execució amb altres noms de node donaria un altre repartiment igual d'arbitrari); amb 10 vnodes millora i amb 150 la desviació baixa a un 6 %, que continua reduint-se amb més vnodes (amb 500, un 4,6 %; amb 1000, un 2,3 %: la desviació cau aproximadament amb l'arrel quadrada del nombre de vnodes, i per això Cassandra en va fer servir durant anys 256 per node). Segona: afegir el cinquè node mou el 19 % de les claus (enfront del 80 % del mòdul), el mínim teòric d'1/5. Tercera: treure'n un de quatre mou exactament les seves claus (25 %), i aquestes claus es reparteixen entre els tres supervivents, no cauen sobre un de sol, gràcies al fet que els 150 vnodes d'inv-vlc tenien successors diferents.
Aquest anell, amb vnodes i amb rèpliques als nodes següents de l'anell, és el que fan servir Dynamo, Cassandra, Riak i l'encaminament de molts clients de Memcached. El veurem aplicat a 04-04.
- Rendezvous hashing
Hi ha una alternativa a l'anell, més senzilla d'implementar i sense problemes de repartiment: el rendezvous hashing o highest random weight (HRW). Per a cada clau es calcula pes(node, clau) = hash(node + clau) per a tots els nodes i es tria el de més pes. Propietats:
- En treure un node, només les claus que el tenien com a guanyador canvien (van al seu segon millor): moviment mínim, igual que l'anell.
- En afegir-ne un, només es mouen les claus per a les quals el nou guanya: 1/(N+1) de mitjana.
- Sense vnodes: el repartiment és uniforme per construcció, perquè cada clau "sorteja" entre tots els nodes.
- Dona gratis una llista ordenada de nodes per preferència, útil per triar les R rèpliques (els R millors).
- Cost O(N) per cerca, enfront d'O(log V) de l'anell: perfecte per a desenes de nodes, pitjor per a milers.
def node_rendezvous(clau: str, nodes: list[str]) -> str:
return max(nodes, key=lambda n: hash_clau(f"{n}|{clau}"))El fan servir, entre d'altres, el particionament d'alguns balancejadors i sistemes de memòria cau. Per a Quilòmetre Zero, amb menys d'una vintena de nodes per servei, seria una opció perfectament vàlida; l'anell és més comú perquè és el que porten les bases de dades que farem servir.
- Reequilibratge de particions
Fins ara hem assignat claus directament a nodes. Els sistemes reals solen introduir un nivell intermedi: les claus s'assignen a particions i les particions a nodes. Moure particions senceres entre nodes (reequilibratge, rebalancing) és més manejable que moure claus soltes. Hi ha tres esquemes:
| Esquema | Com funciona | Avantatges | Inconvenients | Qui el fa servir |
|---|---|---|---|---|
| Nombre fix de particions | Es creen moltes més particions que nodes (p. ex. 1024 per a 10 nodes); en afegir un node, aquest "roba" unes quantes particions senceres de cada node existent | Senzill; només es mouen particions completes; permet pesos | Cal encertar el nombre al principi: massa = sobrecàrrega; poques = límit de creixement | Riak, Elasticsearch, Couchbase, Redis Cluster (16384 slots, 04-05) |
| Particionament dinàmic | Una partició que supera una mida (p. ex. 10 GB) es divideix en dues; una que es buida es fusiona amb la seva veïna | S'adapta al volum; funciona amb rang i amb hash | Un conjunt de dades nou comença amb 1 partició (un sol node treballa): es mitiga amb pre-splitting | HBase, MongoDB, CockroachDB |
| Proporcional als nodes | Cada node té un nombre fix de particions (vnodes); afegir un node divideix aleatòriament particions existents | La mida de partició es manté estable en créixer | Divisions aleatòries, requereix hash | Cassandra, Ketama |
Dues regles d'operació:
- El reequilibratge ha de ser gradual i limitat en amplada de banda: moure particions satura la xarxa i els discos dels nodes implicats, just quan a més continuen servint trànsit.
- El reequilibratge automàtic és còmode però perillós combinat amb la detecció de fallades: un node lent (no mort) pot ser declarat caigut, el sistema comença a moure les seves particions, la càrrega extra fa que altres nodes semblin lents… una cascada. Molts operadors prefereixen que el sistema proposi el reequilibratge i que un humà l'aprovi.
- Encaminament de peticions i descobriment de particions
Si cataleg vol llegir l'estoc de formatge-curat, a quin node es connecta? És el problema de descobriment de servei aplicat a particions, i té tres respostes:
flowchart LR
subgraph a["(a) Qualsevol node"]
C1[Client] --> N1a[Node 2]
N1a -- reenvia --> N2a[Node 4<br/>propietari]
end
subgraph b["(b) Capa d'encaminament"]
C2[Client] --> R[Router / proxy]
R --> N2b[Node 4<br/>propietari]
end
subgraph c["(c) Client informat"]
C3[Client<br/>coneix el mapa] --> N2c[Node 4<br/>propietari]
end
| Opció | Qui coneix el mapa de particions | Exemple | Comentari |
|---|---|---|---|
| (a) Qualsevol node | Tots els nodes (protocol gossip) | Cassandra, Riak | El client és simple; un salt extra de xarxa en el pitjor cas |
| (b) Capa d'encaminament | El router (sovint amb ajuda d'un coordinador) | mongos a MongoDB, moxi a Couchbase, proxies de Redis |
El router pot convertir-se en coll d'ampolla; cal replicar-lo |
| (c) Client informat | La biblioteca client (descarrega el mapa i el guarda en memòria cau) | Redis Cluster (MOVED), HBase, clients de Kafka |
Màxim rendiment; el client ha de gestionar mapes obsolets |
En tots els casos el problema de fons és el mateix: tots els participants han d'estar d'acord sobre quina partició viu a quin node, i aquest acord ha de sobreviure a fallades. És exactament el problema de consens de 03-03, i per això molts sistemes deleguen el mapa en un servei de coordinació: ZooKeeper (HBase, Kafka clàssic, SolrCloud) o etcd (Kubernetes, CockroachDB). Els nodes es registren a ZooKeeper/etcd, la capa d'encaminament se subscriu als canvis, i quan una partició canvia de node el router se n'assabenta en mil·lisegons. Altres sistemes (Cassandra, Riak) eviten la dependència externa i difonen el mapa per gossip, acceptant que durant uns segons alguns nodes tinguin una vista antiga.
Per a Quilòmetre Zero, que ja fa servir etcd per triar el líder del relay outbox (relay_lider.py de 03-03), l'opció natural per a les seves pròpies particions (les rèpliques d'inventari inv-bcn i inv-vlc, i les que vinguin) és desar el mapa a etcd sota /km0/inventari/particions/ i que els clients gRPC el llegeixin i s'hi subscriguin amb watch. Un esbós:
# km0/serveis/inventari/mapa_particions.py
import etcd3, json
client = etcd3.client(host="etcd", port=2379)
PREFIX = "/km0/inventari/particions/"
def publicar(particio: str, node: str, lease_ttl: int = 15):
lease = client.lease(lease_ttl) # si el node mor, l'entrada caduca
client.put(f"{PREFIX}{particio}", json.dumps({"node": node}), lease=lease)
return lease # el node ha de cridar lease.refresh() periòdicament
def carregar_mapa() -> dict[str, str]:
return {meta.key.decode().removeprefix(PREFIX): json.loads(v)["node"]
for v, meta in client.get_prefix(PREFIX)}
def vigilar(en_canviar):
"""Crida `en_canviar(mapa)` cada vegada que una partició canvia de node."""
for _ in client.watch_prefix(PREFIX)[0]:
en_canviar(carregar_mapa())El lease és el mateix mecanisme de 03-03: si inv-vlc deixa de renovar-lo, la seva entrada desapareix i els clients saben que la partició no té propietari, sense necessitat d'un detector de fallades a part.
- Triar la clau de partició a Quilòmetre Zero
Amb tot l'anterior, prenem les dues decisions que l'equip té pendents. Recorda que la clau de partició es tria per les consultes dominants i per la distribució de la càrrega, i que es pot complementar amb índexs globals asíncrons per a les consultes secundàries.
Comandes (km0_comandes, 40 000 comandes/dia en campanya, lectura dominant: "les meves comandes" a l'app, i "comanda per id" a la confirmació):
| Candidata | A favor | En contra | Veredicte |
|---|---|---|---|
data_creacio (rang) |
Consultes per període per a analitica |
Punt calent permanent a "avui" (Setmana de la Verema) | Descartada |
comanda_id (hash) |
Repartiment perfecte; "comanda per id" en un salt | "Les meves comandes" és scatter/gather sobre totes les particions | Només per a la taula comandes_per_id |
client_id (hash) + data DESC (ordre) |
"Les meves comandes" en un node, ja ordenades; repartiment uniforme (milers de clients) | Un client amb moltíssimes comandes (un restaurant) crea una partició gran; "comanda per id" necessita conèixer el client | Triada com a clau principal |
mercat (hash) |
Consultes de repartiment per ciutat |
Només 4 valors: 4 particions com a màxim, València el doble de gran | Descartada com a clau; sí com a índex global asíncron comandes_per_mercat |
productor_id |
Panell del productor | Molt desigual (Horta La Vega té 30 vegades més comandes que Celler Roure Alt) | Índex global asíncron |
Decisió: la clau de partició de les comandes és client_id, amb clustering per data DESC; es manté una segona taula comandes_per_id (el patró "una taula per consulta" que desenvoluparem a 04-04) alimentada en la mateixa escriptura, i els índexs per mercat i per productor es materialitzen des de comandes.esdeveniments. Per al cas del client amb massa comandes s'afegeix a la clau un cub temporal (client_id + any-mes), de manera que cap partició no creix sense límit: és el particionament compost de l'apartat 3 portat a la clau de partició.
Estoc (km0_inventari, rèpliques inv-bcn/inv-vlc, comptadors que es decrementen amb estoc.reservat):
| Candidata | A favor | En contra | Veredicte |
|---|---|---|---|
producte_slug (hash) |
Cada comptador en un sol node: decrements atòmics sense coordinació entre nodes | formatge-curat a la Setmana del Formatge Artesà és una clau calenta |
Triada; la clau calenta es tracta amb cues i amb memòria cau (04-05), no canviant la partició |
productor_id |
Un productor actualitza tot el seu estoc en un node | Horta La Vega concentra el 40 % del catàleg en una partició | Descartada |
mercat |
Estoc per ciutat per a la web | L'estoc és del productor, no del mercat: duplicaria comptadors i exigiria transaccions entre particions | Descartada |
Aquí la lliçó important és que un comptador que ha de ser consistent (CP, 03-02) ha de viure sencer en una partició: decrementar estoc que està repartit entre dos nodes exigiria 2PC (03-05) a cada reserva. Per això producte_slug guanya encara que tingui claus calentes.
Errors Comuns i Consells
- Fer servir
hash()de Python (o elhashCodeper defecte d'altres llenguatges) com a hash de partició. Està aleatoritzat per procés o depèn de la implementació. Fes servir un hash explícit i documentat (MD5, MurmurHash3, xxHash) i fixa'n la versió. - Confondre particionar amb replicar. Particionar sense replicar redueix la disponibilitat: cada node que cau s'endú la seva part de les dades. Cada partició necessita el seu factor de replicació (04-04).
- Triar la clau de partició pel model de domini i no per les consultes. "Les comandes tenen id, per tant la clau és
comanda_id" produeix un scatter/gather a la consulta més freqüent. Comença per llistar les consultes i la seva freqüència. - Claus amb pocs valors diferents (
mercat,estat,tipus). Limiten el nombre de particions i creen desequilibris. Cardinalitat alta primer. - Ignorar les particions que creixen sense límit (el client amb un milió de comandes, la flota de repartiment amb 2,4 milions de posicions al dia). Afegeix un cub temporal a la clau.
- Hashing consistent sense nodes virtuals. Repartiment desigual i, en caure un node, tota la seva càrrega sobre un sol successor. Fes servir 100–200 vnodes, o rendezvous hashing.
- Reequilibratge automàtic agressiu. Pot desencadenar cascades quan un node només està lent. Limita l'amplada de banda i considera l'aprovació manual.
- Oblidar que el mapa de particions és estat distribuït. Necessita consens (ZooKeeper/etcd) o gossip, i els clients han de tolerar mapes obsolets (reintentar després d'un
MOVEDo equivalent).
Consell final: quan dubtis entre rang i hash, pregunta't si la consulta per rang és realment necessària a la ruta calenta o si la pot servir un índex global o el llac de dades (04-02) fora de línia. Gairebé sempre és el segon.
Exercicis
Exercici 1. Modifica anell_consistent.py perquè suporti pesos: afegir_node(node, pes=1.0) ha de crear int(vnodes * pes) nodes virtuals. Construeix un anell amb inv-bcn (pes 2,0), inv-vlc (1,0) i inv-gir (1,0), reparteix les 100 000 claus i comprova que inv-bcn en rep aproximadament el 50 %. Què passa amb treure_node si no guardes el pes?
Exercici 2. Les posicions dels repartidors (140 furgonetes com furgoneta-3, 2,4 milions de posicions al dia entre totes) es volen desar particionades. Les consultes són: (a) "última posició del repartidor R" (milers per minut, des de la web del client), (b) "recorregut del repartidor R entre dues hores d'avui" (suport), (c) "tots els repartidors ara mateix a València" (panell d'operacions). Proposa la clau de partició (i d'ordenació si escau), indica quina consulta queda com a scatter/gather o índex global, i explica per què repartidor_id a seques genera una partició sense límit i com ho corregeixes.
Exercici 3. Implementa node_rendezvous(clau, nodes) i una funció repliques_rendezvous(clau, nodes, r) que retorni els r nodes de més pes. Amb 5 nodes i 100 000 claus, mesura (a) la desviació relativa del repartiment, (b) la fracció de claus el node principal de les quals canvia en afegir un sisè node, i (c) la fracció de claus el conjunt de 3 rèpliques de les quals canvia. Compara (b) amb el resultat de l'anell.
Solucions
Solució 1:
class AnellConsistentPonderat(AnellConsistent):
def __init__(self, vnodes: int = 150):
super().__init__(vnodes=vnodes)
self._pesos: dict[str, int] = {} # node -> nombre real de vnodes creats
def afegir_node(self, node: str, pes: float = 1.0) -> None:
n = max(1, int(self.vnodes * pes))
self._pesos[node] = n
for i in range(n):
pos = hash_clau(f"{node}#{i}")
if pos not in self._propietari:
bisect.insort(self._posicions, pos)
self._propietari[pos] = node
def treure_node(self, node: str) -> None:
for i in range(self._pesos.pop(node, 0)):
pos = hash_clau(f"{node}#{i}")
if self._propietari.get(pos) == node:
del self._propietari[pos]
self._posicions.remove(pos)
anell = AnellConsistentPonderat()
anell.afegir_node("inv-bcn", 2.0); anell.afegir_node("inv-vlc"); anell.afegir_node("inv-gir")
print(distribucio(anell, claus)) # inv-bcn ≈ 50 000, els altres ≈ 25 000 cadascunSi treure_node recalcula amb self.vnodes en lloc del nombre real, per a inv-bcn només eliminaria 150 dels seus 300 vnodes i deixaria a l'anell 150 posicions que apunten a un node mort: les claus que hi cauen anirien a un node inexistent. Per això cal guardar quants vnodes s'han creat per node (o recórrer-los des de _propietari).
Solució 2:
Clau de partició repartidor_id + dia (per exemple furgoneta-3|2026-09-14), amb clau d'ordenació ts DESC. (a) "Última posició" és la primera fila de la partició d'avui: un node, lectura mínima. (b) "Recorregut entre dues hores" és un rang dins de la mateixa partició, ordenat per temps. (c) "Tots els repartidors a València ara" no esmenta cap repartidor: un scatter/gather sobre 140 particions seria acceptable (140 lectures petites) però és millor un índex global asíncron ultima_posicio_per_mercat, actualitzat amb cada posició (o cada N segons), particionat per mercat, que és el que el panell realment mostra. Amb repartidor_id a seques, furgoneta-3 acumularia unes 17 000 posicions diàries sense límit (més de 6 milions l'any en una sola partició, que a Cassandra començaria a degradar lectures i compactacions); el cub diari acota la partició i a més fa trivial l'esborrat de dades antigues (s'elimina la partició del dia sencera). Aquesta dada era AP a la taula de 03-02, i res del que hem dit no ho canvia: amb W=1 la posició s'accepta encara que faltin rèpliques.
Solució 3:
def node_rendezvous(clau, nodes):
return max(nodes, key=lambda n: hash_clau(f"{n}|{clau}"))
def repliques_rendezvous(clau, nodes, r):
return sorted(nodes, key=lambda n: hash_clau(f"{n}|{clau}"), reverse=True)[:r]
nodes5 = ["n1", "n2", "n3", "n4", "n5"]; nodes6 = nodes5 + ["n6"]
recompte = Counter(node_rendezvous(c, nodes5) for c in claus)
print("desviació:", round(desviacio_relativa(recompte), 2), "%") # ≈ 0,3-0,6 % sense vnodes
mov = sum(node_rendezvous(c, nodes5) != node_rendezvous(c, nodes6) for c in claus) / len(claus)
print("principal canvia:", round(mov * 100, 1), "%") # ≈ 16,7 % (= 1/6)
mov_r = sum(set(repliques_rendezvous(c, nodes5, 3)) != set(repliques_rendezvous(c, nodes6, 3))
for c in claus) / len(claus)
print("conjunt de rèpliques canvia:", round(mov_r * 100, 1), "%") # ≈ 50 % (= 3/6)(a) La desviació queda al voltant del 0,3 %, sense necessitat de nodes virtuals i millor que el 6 % de l'anell amb 150 vnodes, perquè cada clau tria entre tots els nodes de manera independent (l'anell, en canvi, depèn de com de ben repartides caiguin les posicions dels vnodes). (b) El node principal canvia per a 1/6 de les claus, igual que l'anell amb vnodes (el mínim teòric) i molt lluny del 83 % del mòdul. (c) Amb 3 rèpliques, el conjunt canvia en aproximadament la meitat de les claus (3/6): el node nou entra entre els tres millors amb probabilitat 3/6, i en cada cas es mou una rèplica, no les tres; el volum de dades mogut continua sent mínim (cada clau mou com a molt una còpia), encara que la fracció de conjunts afectats sigui més gran.
Conclusió
Particionar és repartir dades diferents entre nodes, i replicar és copiar les mateixes; els sistemes reals fan totes dues coses, replicant cada partició. Hem vist que el repartiment per rang preserva les consultes per interval però concentra la càrrega (les comandes de la Setmana de la Verema escrivint totes a la partició d'"avui"), que el repartiment per hash dispersa les claus a costa d'aquests rangs, i que el particionament compost (hash del client, ordre per data) recupera el millor de tots dos dins de cada partició. Els índexs secundaris obliguen a triar entre índexs locals amb scatter/gather i índexs globals asíncrons, i Quilòmetre Zero ha triat els segons per a repartiment alimentant-los des de comandes.esdeveniments. L'assignació ingènua hash % N mou el 80 % de les dades en passar de 4 a 5 nodes, com ha mesurat hash_modul.py; el hashing consistent ho redueix al 20 % mínim, i els nodes virtuals d'AnellConsistent fan baixar la desviació del repartiment de més del 100 % a al voltant del 6 % (i menys com més vnodes) i reparteixen la càrrega d'un node caigut entre tots els altres. El rendezvous hashing aconsegueix un repartiment encara més uniforme sense anell ni vnodes. Per sobre de les claus, els sistemes mouen particions senceres (fixes, dinàmiques o proporcionals als nodes) i publiquen el mapa de particions a través de gossip o d'un coordinador com ZooKeeper o etcd, que ja coneixíem per l'elecció de líder de 03-03. I hem fixat dues decisions de Quilòmetre Zero: comandes particionades per client_id (amb cub mensual i taula secundària comandes_per_id) i estoc particionat per producte_slug, perquè un comptador CP ha de viure sencer en una partició.
Amb les claus repartides, toca veure els sistemes que les desen. Els tres temes següents són tres maneres diferents d'emmagatzemar bytes a escala: fitxers, objectes i registres de base de dades. Comencem per la més antiga i la que més s'assembla al que ja coneixes, els sistemes de fitxers distribuïts, amb NFS, HDFS i Ceph, i amb HDFS com el llac de dades on Quilòmetre Zero desarà els esdeveniments de comandes.esdeveniments que el Mòdul 5 processarà en massa.
Curs d'Arquitectures Distribuïdes
Mòdul 1: Introducció als Sistemes Distribuïts
- Conceptes Bàsics de Sistemes Distribuïts
- Models de Sistemes Distribuïts
- Avantatges i Desafiaments dels Sistemes Distribuïts
- Les Fal·làcies de la Computació Distribuïda
- Temps, Rellotges i Ordenació d'Esdeveniments
- Del Monòlit a la Plataforma Distribuïda: el Cas Quilòmetre Zero
Mòdul 2: Comunicació en Sistemes Distribuïts
- Protocols de Comunicació
- RPC i RMI
- gRPC i Serialització de Dades
- Missatgeria i Cues de Missatges
- Patrons de Comunicació Asíncrona
Mòdul 3: Consistència i Replicació
- Models de Consistència
- El Teorema CAP i PACELC
- Algorismes de Consens
- Replicació de Dades
- Transaccions Distribuïdes i Sagues
Mòdul 4: Emmagatzematge Distribuït
- Particionament de Dades i Hashing Consistent
- Sistemes de Fitxers Distribuïts
- Emmagatzematge d'Objectes
- Bases de Dades Distribuïdes
- Memòries Cau Distribuïdes
Mòdul 5: Computació Distribuïda
- Models de Computació Distribuïda
- MapReduce i Hadoop
- Spark i Computació en Memòria
- Processament de Fluxos de Dades
- Planificació de Treballs i Pipelines de Dades
Mòdul 6: Seguretat en Sistemes Distribuïts
- Autenticació i Autorització
- Xifratge i Protecció de Dades
- Gestió d'Identitats
- Seguretat entre Serveis: mTLS i Gestió de Secrets
- Passarel·les d'API, Limitació de Taxa i Auditoria
Mòdul 7: Monitoratge i Manteniment
- Monitoratge de Sistemes Distribuïts
- Logs Centralitzats i Traçabilitat Distribuïda
- Gestió de Fallades i Recuperació
- Patrons de Resiliència: Timeouts, Reintents i Circuit Breaker
- Automatització i Orquestració
- Proves en Sistemes Distribuïts i Enginyeria del Caos
