A la lliçó anterior vam aprendre a decidir quin node desa cada clau. Ara baixem un nivell i ens ocupem del sistema que emmagatzema físicament els bytes de manera que molts clients, des de moltes màquines, els vegin com a fitxers normals. Un sistema de fitxers distribuït (DFS, Distributed File System) és la forma més antiga d'emmagatzematge compartit en xarxa i continua sent la que més s'assembla al que qualsevol programador coneix: directoris, fitxers, open, read, write. Però sota aquesta interfície familiar hi ha decisions molt diferents: NFS es va dissenyar perquè una oficina compartís el disc d'un servidor; GFS i HDFS perquè milers de màquines barates desessin petabytes de fitxers enormes que s'escriuen un cop i es llegeixen moltes vegades; GlusterFS i Ceph per no dependre de cap node central. Per a Quilòmetre Zero la pregunta és concreta: on acumular els milions d'esdeveniments de comandes.esdeveniments i els logs de clics de la web perquè el Mòdul 5 els pugui processar en massa. La resposta serà HDFS, i el muntarem al docker-compose.yml de km0/.

Contingut

  1. Què és un sistema de fitxers distribuït i quines transparències ofereix
  2. NFS: el model client-servidor
  3. GFS i HDFS: fitxers enormes en màquines barates
  4. Alta disponibilitat del NameNode
  5. El que HDFS no sap fer
  6. GlusterFS i Ceph: sense metadades centralitzades
  7. Taula comparativa i casos d'ús
  8. HDFS a Quilòmetre Zero: el llac de dades
  9. Pràctica: HDFS a docker-compose.yml i l'API WebHDFS
  10. Errors Comuns i Consells
  11. Exercicis
  12. Conclusió

  1. Què és un sistema de fitxers distribuït i quines transparències ofereix

Un DFS presenta als seus clients un espai de noms jeràrquic (directoris i fitxers) el contingut del qual està emmagatzemat en un o diversos servidors remots. El que el distingeix de "copiar fitxers per la xarxa" és que l'accés s'integra al sistema operatiu o en una API equivalent, de manera que els programes no noten (o noten tan poc com sigui possible) on són les dades. Aquest "no notar" és exactament el concepte de transparència que vam definir a 01-01, i un DFS és el millor exemple per repassar-lo:

Transparència Què vol dir en un DFS Qui l'ofereix
D'accés Es fan servir les mateixes crides (open, read) per a fitxers locals i remots NFS (muntatge), Ceph (CephFS), HDFS només parcialment (API pròpia)
D'ubicació El nom del fitxer no revela a quin servidor és (/km0/esdeveniments/…, no //datanode2/…) Tots
De replicació El client no sap quantes còpies hi ha ni quina llegeix HDFS, Ceph, GlusterFS
De fallada Si cau un servidor que desa una còpia, el client continua llegint HDFS, Ceph, GlusterFS; NFS no (un servidor)
De concurrència Diversos clients hi accedeixen sense corrompre's Tots, amb semàntiques molt diferents (apartat 2)
De migració / escalat Es poden afegir discos o nodes sense canviar rutes HDFS, Ceph, GlusterFS

El punt delicat de qualsevol DFS és la semàntica de compartició: què veu un client quan un altre client està escrivint el mateix fitxer. En un sistema local (semàntica UNIX) cada write és visible immediatament per a qualsevol read posterior. Per la xarxa, garantir-ho exigeix que cada operació viatgi al servidor, cosa que mata el rendiment. Cada DFS tria un punt diferent entre rendiment i semàntica, i aquest punt explica gairebé totes les seves diferències.

  1. NFS: el model client-servidor

NFS (Network File System, Sun, 1984; avui a la versió 4.2) és el DFS clàssic: un servidor exporta un directori i els clients el munten al seu arbre local. Cada operació del client es converteix en una crida RPC (l'ONC RPC que vam veure a 02-02, amb XDR com a serialització) al servidor.

sequenceDiagram
    participant App as Aplicació
    participant K as Kernel client<br/>(memòria cau de pàgines i atributs)
    participant S as Servidor NFS
    App->>K: open("/mnt/km0/factures/P-2026-000123.pdf")
    K->>S: LOOKUP + GETATTR (RPC)
    S-->>K: file handle + atributs (mtime, mida)
    App->>K: read(4096 bytes)
    K->>S: READ (si no és a la memòria cau)
    S-->>K: dades
    K-->>App: dades (i es desen a la memòria cau de pàgines)
    App->>K: close()
    K->>S: escriptura de pàgines brutes (close-to-open)

Memòries cau de client i consistència. Sense memòria cau, cada read seria un viatge de xarxa. Per això el client NFS guarda en memòria cau dades i atributs, i aquí apareix el compromís: si el client A té un bloc a la memòria cau i el client B escriu aquest mateix bloc al servidor, A llegirà dades velles fins que revalidi. NFS fa servir la semàntica close-to-open: els canvis d'un client es garanteixen visibles per als altres només quan l'escriptor fa close i el lector fa open després. Entremig, el client revalida els atributs del fitxer cada pocs segons (actimeo, per defecte 3–60 s) i descarta la seva memòria cau si el mtime ha canviat. És consistència eventual amb una finestra de segons, i en el vocabulari de 03-01 no ofereix ni lectura de les pròpies escriptures entre màquines diferents. NFSv4 afegeix delegacions: el servidor pot cedir a un client el dret exclusiu sobre un fitxer mentre ningú més no el demani, cosa que permet guardar en memòria cau amb seguretat i revocar quan apareix un altre client.

Límits. NFS té un sol servidor per exportació: és el límit de capacitat, de rendiment i el punt únic de fallada (es pot muntar en alta disponibilitat amb un parell actiu/passiu i emmagatzematge compartit, però no escala horitzontalment). La versió 4.1 va introduir pNFS (parallel NFS), que separa metadades de dades i permet llegir de diversos servidors, però la seva adopció és limitada.

Quan continua sent vàlid. Moltíssimes vegades. Compartir el directori /home d'un equip de desenvolupament, donar a diverses instàncies d'una aplicació accés als mateixos fitxers de configuració o a un directori d'intercanvi, muntar volums ReadWriteMany a Kubernetes per a càrregues moderades: NFS és simple, és a tots els kernels, i fins a uns quants terabytes i centenars de clients funciona sense més. Quilòmetre Zero el va descartar com a magatzem de fotos de productes per dues raons que veurem a 04-03: vol replicació real entre nodes i una API HTTP directa per a la web, no un muntatge.

  1. GFS i HDFS: fitxers enormes en màquines barates

Google va publicar el 2003 el disseny del seu Google File System (GFS) i Yahoo el va reimplementar en codi obert com a HDFS (Hadoop Distributed File System). Les seves premisses trenquen amb NFS:

  • Els fitxers són enormes (gigabytes o terabytes), i són pocs milions, no milers de milions.
  • S'escriuen un cop (o només s'hi afegeix al final) i es llegeixen moltes vegades, gairebé sempre de manera seqüencial i completa.
  • El maquinari és barat i falla constantment: amb 1000 discos, algun mor cada dia. La tolerància a fallades és la norma, no l'excepció.
  • Importa l'amplada de banda agregada (llegir un terabyte en un minut des de cent màquines), no la latència d'una lectura petita.

Arquitectura. HDFS té dos tipus de nodes:

Rol Quants Què desa Què fa
NameNode 1 actiu (+ 1 en espera, apartat 4) Les metadades: arbre de directoris, permisos, i per a cada fitxer la llista de blocs i a quins DataNodes és cada rèplica. Tot en memòria Atén les operacions d'espai de noms (mkdir, ls, open), decideix on van els blocs nous, ordena re-replicar els que perden còpies
DataNode Desenes a milers Els blocs de dades, com a fitxers normals al seu disc local Serveix lectures i escriptures de blocs directament als clients; envia batecs i l'informe dels seus blocs al NameNode

Els conceptes clau:

  • Blocs grans: per defecte 128 MB (enfront dels 4 KB d'un sistema de fitxers local). Un fitxer d'1 GB són 8 blocs. Els blocs grans redueixen el nombre de metadades que el NameNode desa en memòria i fan que cada lectura sigui una transferència seqüencial llarga, que és el que els discos fan bé. Un fitxer més petit que un bloc ocupa només el que mesura (no malbarata 128 MB), però sí que consumeix una entrada de metadades completa.
  • Factor de replicació: cada bloc es copia en dfs.replication DataNodes, per defecte 3. La replicació és per bloc, no per fitxer, i el NameNode la vigila: si un DataNode deixa d'enviar batecs (10 minuts per defecte), tots els seus blocs queden "infrareplicats" i el NameNode ordena copiar-los des de les rèpliques supervivents a altres nodes.
  • Rack awareness: el NameNode sap a quin armari (rack) és cada DataNode. Amb 3 rèpliques, col·loca la primera al node del client (si és un DataNode), la segona en un node d'un altre rack i la tercera en un altre node del mateix rack que la segona. Així una fallada d'un rack sencer (un switch) no perd cap bloc, i només una de les tres còpies creua entre racks (amplada de banda entre racks, que és l'escassa).
  • Write-once / append: un fitxer es crea, s'escriu (per un sol escriptor) i es tanca; després és immutable, llevat que s'hi afegeixi al final (append). No hi ha escriptures enmig del fitxer. Això simplifica enormement la consistència: no hi ha dos escriptors concurrents sobre els mateixos bytes i les rèpliques d'un bloc tancat són idèntiques per sempre.

Flux d'escriptura i lectura. L'essencial és que les dades mai no passen pel NameNode; només les metadades:

sequenceDiagram
    participant C as Client HDFS
    participant NN as NameNode
    participant D1 as DataNode 1
    participant D2 as DataNode 2
    participant D3 as DataNode 3
    Note over C,D3: ESCRIPTURA de /km0/esdeveniments/2026-09-14/comandes.jsonl
    C->>NN: create(ruta)
    NN-->>C: ok (fitxer en construcció)
    C->>NN: addBlock()
    NN-->>C: bloc B1 → [D1, D2, D3] (rack awareness)
    C->>D1: paquets de B1
    D1->>D2: reenvia (pipeline)
    D2->>D3: reenvia (pipeline)
    D3-->>D2: ack
    D2-->>D1: ack
    D1-->>C: ack
    C->>NN: complete(ruta)
    Note over C,D3: LECTURA
    C->>NN: getBlockLocations(ruta)
    NN-->>C: B1 → [D1, D2, D3], B2 → [D2, D4, D5] … (ordenats per proximitat)
    C->>D1: read B1
    C->>D2: read B2

A l'escriptura, el client demana al NameNode on posar el bloc, i envia les dades al primer DataNode, que les reenvia al segon, que les reenvia al tercer (pipeline de replicació): el client només emet les dades un cop i la replicació consumeix amplada de banda entre DataNodes, no del client. El bloc es considera escrit quan els tres han confirmat: és replicació síncrona en el sentit de 03-04, i per això HDFS és CP (una escriptura no s'accepta si no pot assolir el nombre mínim de rèpliques, dfs.namenode.replication.min, que per defecte és 1). A la lectura, el client obté la llista de blocs amb les seves ubicacions i llegeix cada bloc directament del DataNode més proper; si un falla, passa al següent de la llista. Cada bloc porta sumes de comprovació (CRC32 cada 512 bytes) que el client verifica: un bloc corrupte es descarta, es llegeix d'una altra rèplica i es notifica al NameNode perquè el re-repliqui.

  1. Alta disponibilitat del NameNode

En el disseny original, el NameNode era un punt únic de fallada: si queia, el clúster sencer quedava inaccessible, encara que les dades continuessin intactes als DataNodes. Des de Hadoop 2 existeix la configuració d'alta disponibilitat (HA), que aplica el que ja sabem de 03-03 i 03-04:

  • Dos (o més) NameNodes: un d'actiu i un altre en espera (standby). Tots dos tenen l'espai de noms en memòria.
  • El registre de canvis (edit log) de l'actiu s'escriu en un quòrum de JournalNodes (normalment 3): l'operació de metadades es confirma quan la majoria l'ha persistit. El NameNode en espera llegeix aquest registre contínuament i aplica els canvis, de manera que la seva còpia va només uns mil·lisegons endarrerida. Aquest quòrum és l'aplicació directa del W > N/2 de 03-04.
  • Els DataNodes envien batecs i informes de blocs a tots dos NameNodes, perquè l'standby conegui la ubicació de cada bloc i pugui prendre el comandament sense reconstruir-la.
  • ZooKeeper (03-03) decideix qui és l'actiu: cada NameNode té un procés ZKFailoverController que manté un node efímer a ZooKeeper; si l'actiu deixa de renovar-lo, l'standby adquireix el bloqueig i es promociona.
  • Fencing: abans de promocionar-se, el nou actiu s'assegura que l'antic no pugui continuar escrivint als JournalNodes (que només accepten escriptures del NameNode amb l'època més alta, igual que els termes de Raft), i opcionalment el mata per SSH. És la defensa contra el "cervell dividit" que ja vam discutir a 03-04.
flowchart LR
    ZK[(ZooKeeper<br/>elecció d'actiu)]
    NN1[NameNode actiu] --- ZK
    NN2[NameNode standby] --- ZK
    NN1 -- escriu edits --> J1[(JournalNode 1)]
    NN1 -- escriu edits --> J2[(JournalNode 2)]
    NN1 -- escriu edits --> J3[(JournalNode 3)]
    NN2 -. llegeix edits .-> J1
    NN2 -. llegeix edits .-> J2
    NN2 -. llegeix edits .-> J3
    D1[DataNode] -- batecs i informes --> NN1
    D1 -- batecs i informes --> NN2

Amb HA, la fallada del NameNode actiu es resol en desenes de segons sense intervenció. Sense HA (com al docker-compose de pràctica), un SecondaryNameNode només compacta l'edit log; no és un respatller i no pot prendre el comandament, un malentès tan freqüent que mereix subratllar-se.

  1. El que HDFS no sap fer

Les seves premisses són els seus límits, i convé tenir-los clars abans de posar a HDFS alguna cosa que no hi encaixa:

  • Molts fitxers petits. Cada fitxer, directori i bloc ocupa uns 150 bytes a la memòria del NameNode. Cent milions de fitxers de 10 KB són 15 GB de heap i 1 TB de dades que en blocs de 128 MB haurien estat 8 000 entrades. A més, cada lectura d'un fitxer petit paga una crida al NameNode i una connexió a un DataNode per llegir uns pocs KB. Les fotos dels productes (milers de fitxers de 200 KB) són un mal cas per a HDFS; per això van a un magatzem d'objectes (04-03). Si cal posar fitxers petits a HDFS, s'agrupen (fitxers SequenceFile, HAR o, millor, particions diàries d'un sol fitxer gran, com farem amb els esdeveniments).
  • Accés aleatori de baixa latència. Llegir un registre concret exigeix localitzar el bloc, obrir una connexió i llegir com a mínim un paquet; desenes de mil·lisegons. HDFS no és una base de dades: HBase es va construir a sobre precisament per donar accés aleatori, gestionant els seus propis índexs i fitxers grans.
  • Escriptors concurrents i modificacions enmig del fitxer. Un fitxer té un únic escriptor i només admet append. Un log que molts serveis escriuen alhora ha de passar abans per Kafka (02-04) i abocar-se a HDFS per lots.
  • Semàntica POSIX completa. No hi ha mmap, ni bloqueigs, ni escriptures parcials; el "muntatge" amb NFS Gateway o FUSE és una capa de compatibilitat amb limitacions.

  1. GlusterFS i Ceph: sense metadades centralitzades

El NameNode, fins i tot amb HA, és un límit: tot l'espai de noms cap a la memòria d'una màquina i tota operació de metadades hi passa. Una segona família de DFS elimina el servidor de metadades i localitza les dades per càlcul, amb les mateixes idees del hashing consistent de 04-01.

GlusterFS agrupa directoris exportats per diversos servidors (bricks) en un volum. No hi ha servidor de metadades: la ubicació de cada fitxer es calcula amb un hash del seu nom sobre un rang assignat a cada brick (elastic hashing), i el client, que coneix la configuració del volum, parla directament amb el brick correcte. Els volums poden ser distribuïts (cada fitxer en un brick), replicats (cada fitxer en N bricks) o dispersos (amb codificació d'esborrament, que veurem a 04-03), i es combinen. És senzill, es munta com un sistema de fitxers normal (FUSE o NFS) i escala bé per a fitxers mitjans i grans; el seu punt feble són les operacions de directori (un ls ha de preguntar a tots els bricks) i l'autoreparació després de fallades.

Ceph és més ambiciós: un magatzem d'objectes distribuït (RADOS) sobre el qual es construeixen tres interfícies: blocs (RBD, discos virtuals per a màquines virtuals i Kubernetes), objectes (RGW, compatible amb S3, 04-03) i fitxers (CephFS). Les seves peces:

  • OSD (Object Storage Daemon): un procés per disc, que desa objectes i es replica amb altres OSD entre parells.
  • Monitors (MON): un petit quòrum (Paxos, 03-03) que manté el mapa del clúster (quins OSD existeixen i el seu estat). No desen dades ni metadades de fitxers.
  • CRUSH (Controlled Replication Under Scalable Hashing): l'algorisme que, a partir del nom d'un objecte i del mapa del clúster, calcula a quins OSD són les seves rèpliques. És un parent directe del hashing consistent de 04-01, amb una diferència important: CRUSH entén la topologia (disc → servidor → rack → sala) i unes regles ("tres rèpliques en tres racks diferents"), de manera que les rèpliques no només es reparteixen uniformement sinó que respecten els dominis de fallada. Qualsevol client amb el mapa calcula la ubicació sense preguntar a ningú: no hi ha NameNode, i afegir un OSD mou només la fracció proporcional d'objectes, com a l'anell.
  • MDS (Metadata Server): només per a CephFS, gestiona l'arbre de directoris; n'hi pot haver diversos d'actius, repartint-se subarbres dinàmicament.

Ceph és el motor d'emmagatzematge de molts clouds privats (OpenStack, Proxmox, Kubernetes amb Rook). És també notablement més complex d'operar que HDFS o GlusterFS.

  1. Taula comparativa i casos d'ús

NFS HDFS GlusterFS Ceph
Metadades A l'únic servidor NameNode centralitzat (en memòria) Sense servidor: hash del nom Sense servidor per a objectes (CRUSH); MDS per a CephFS
Dades Un servidor Blocs de 128 MB replicats en DataNodes Fitxers sencers en bricks Objectes en OSD, amb CRUSH
Replicació No (o actiu/passiu extern) Per bloc, síncrona, rack aware Per fitxer, síncrona Per objecte, síncrona, topologia amb regles
Semàntica Close-to-open, POSIX aproximat Write-once + append, un escriptor POSIX aproximat POSIX (CephFS), bloc, objecte
Escala Terabytes, centenars de clients Petabytes, milers de nodes; milions de fitxers Petabytes Petabytes a exabytes
Fitxers petits Malament Regular Bé (com a objectes)
Accés aleatori Malament
Interfície Muntatge del SO API Java/CLI, WebHDFS (HTTP), FUSE limitat Muntatge (FUSE/NFS) Muntatge, S3, disc de blocs
Complexitat operativa Molt baixa Mitjana Mitjana Alta
Cas d'ús típic Directoris compartits, /home, volums ReadWriteMany moderats Llac de dades per a processament per lots (Mòdul 5) Magatzem de fitxers de mida mitjana, contingut web, backups Infraestructura d'emmagatzematge unificada d'un cloud privat

  1. HDFS a Quilòmetre Zero: el llac de dades

A 01-06 vam fixar que analitica processaria per lots les dades històriques i en streaming les recents. Les dades històriques necessiten un lloc on acumular-se durant anys, barat per terabyte, tolerant a fallades i optimitzat perquè el Mòdul 5 les llegeixi senceres de manera paral·lela. Aquest lloc és el llac de dades (data lake) sobre HDFS, amb dues fonts:

  • Els esdeveniments del tòpic comandes.esdeveniments (comanda.creada, estoc.reservat, pagament.confirmat, … amb l'embolcall id_esdeveniment/tipus/versio/data_ms/origen/dades de 02-05). Kafka reté els esdeveniments uns dies; un consumidor d'analitica els aboca a HDFS per lots, en un fitxer per dia i tipus: /km0/esdeveniments/2026-09-14/comandes.jsonl. Cada línia és un esdeveniment en JSON (format JSON Lines). Amb 40 000 comandes i uns 6 esdeveniments per comanda, un dia de campanya són uns 250 000 esdeveniments i uns 150 MB: un o dos blocs, una mida ideal per a HDFS.
  • Els logs de clics de la web (pàgina vista, producte vist, afegit al cistell), que el servidor web escriu en fitxers rotats cada hora i que un procés puja a /km0/clics/2026-09-14/hora=13/web-01.jsonl.

L'organització per directoris de data (data=2026-09-14/) no és casual: és la partició per rang de 04-01 aplicada a fitxers, i permetrà al Mòdul 5 processar "només la Setmana de la Verema" sense llegir la resta del llac. Els esdeveniments són immutables (write-once hi encaixa a la perfecció), l'accés és seqüencial i massiu, i ningú no necessita llegir "l'esdeveniment 123" amb baixa latència: per a això hi ha la base de dades de comandes (04-04). Aquesta és la divisió que tanca el mòdul: HDFS per al massiu i fred, la base de dades per a l'operatiu, els objectes per als fitxers que la web serveix.

  1. Pràctica: HDFS a docker-compose.yml i l'API WebHDFS

Afegim al docker-compose.yml de km0/ un NameNode i dos DataNodes amb la imatge oficial apache/hadoop. La imatge es configura amb variables d'entorn el nom de les quals replica el fitxer XML de Hadoop (CORE-SITE.XML_fs.defaultFS equival a la propietat fs.defaultFS de core-site.xml):

# km0/docker-compose.yml (fragment)
x-hadoop-env: &hadoop-env
  CORE-SITE.XML_fs.defaultFS: hdfs://namenode:8020
  CORE-SITE.XML_hadoop.http.staticuser.user: km0
  HDFS-SITE.XML_dfs.replication: "2"            # només tenim 2 DataNodes
  HDFS-SITE.XML_dfs.namenode.rpc-address: namenode:8020
  HDFS-SITE.XML_dfs.namenode.http-address: 0.0.0.0:9870
  HDFS-SITE.XML_dfs.webhdfs.enabled: "true"
  HDFS-SITE.XML_dfs.permissions.enabled: "false" # simplifica la pràctica; mai en producció

services:
  namenode:
    image: apache/hadoop:3.4.1
    hostname: namenode
    command: ["hdfs", "namenode"]
    environment:
      <<: *hadoop-env
      ENSURE_NAMENODE_DIR: /tmp/hadoop-root/dfs/name   # formata el NameNode la primera vegada
    ports:
      - "9870:9870"    # interfície web i WebHDFS
      - "8020:8020"    # RPC
    volumes:
      - namenode-data:/tmp/hadoop-root/dfs/name

  datanode-1:
    image: apache/hadoop:3.4.1
    hostname: datanode-1
    command: ["hdfs", "datanode"]
    environment: *hadoop-env
    ports:
      - "9864:9864"    # WebHDFS del DataNode (necessari per llegir/escriure dades des de fora)
    volumes:
      - datanode-1-data:/tmp/hadoop-root/dfs/data
    depends_on: [namenode]

  datanode-2:
    image: apache/hadoop:3.4.1
    hostname: datanode-2
    command: ["hdfs", "datanode"]
    environment: *hadoop-env
    ports:
      - "9865:9864"
    volumes:
      - datanode-2-data:/tmp/hadoop-root/dfs/data
    depends_on: [namenode]

volumes:
  namenode-data:
  datanode-1-data:
  datanode-2-data:

Punts a entendre:

  • dfs.replication: 2 perquè només hi ha dos DataNodes; amb el valor per defecte (3) cada bloc quedaria permanentment "infrareplicat" i fsck ho mostraria com a avís.
  • ENSURE_NAMENODE_DIR fa que el contenidor executi hdfs namenode -format si el directori és buit. Formatar un NameNode amb dades esborra l'espai de noms (els blocs queden orfes als DataNodes), així que va en un volum persistent.
  • El port 9870 és la consola web del NameNode (http://localhost:9870), on es veuen els DataNodes vius, la capacitat i l'explorador de fitxers.

Arrenquem i provem amb la CLI, executada dins del contenidor del NameNode:

docker compose up -d namenode datanode-1 datanode-2
docker compose exec namenode hdfs dfsadmin -report | head -20   # 2 DataNodes vius, capacitat

# Espai de noms del llac de dades
docker compose exec namenode hdfs dfs -mkdir -p /km0/esdeveniments/2026-09-14 /km0/clics
docker compose exec namenode hdfs dfs -ls /km0

# Pujar un fitxer local (dins del contenidor) i llegir-lo
docker compose exec namenode bash -c 'printf "%s\n" \
  "{\"id_esdeveniment\":\"e-1\",\"tipus\":\"comanda.creada\",\"versio\":2,\"data_ms\":1789380000000,\"origen\":\"comandes\",\"dades\":{\"comanda_id\":\"P-2026-000123\",\"client\":\"Anna\"}}" \
  "{\"id_esdeveniment\":\"e-2\",\"tipus\":\"estoc.reservat\",\"versio\":1,\"data_ms\":1789380000450,\"origen\":\"inventari\",\"dades\":{\"comanda_id\":\"P-2026-000123\",\"producte\":\"formatge-curat\",\"quantitat\":2}}" \
  > /tmp/comandes.jsonl'
docker compose exec namenode hdfs dfs -put /tmp/comandes.jsonl /km0/esdeveniments/2026-09-14/comandes.jsonl
docker compose exec namenode hdfs dfs -ls -h /km0/esdeveniments/2026-09-14
docker compose exec namenode hdfs dfs -cat /km0/esdeveniments/2026-09-14/comandes.jsonl

# Afegir al final (append) i comprovar
docker compose exec namenode bash -c 'echo "{\"id_esdeveniment\":\"e-3\",\"tipus\":\"pagament.confirmat\",\"versio\":1,\"data_ms\":1789380002000,\"origen\":\"pagaments\",\"dades\":{\"comanda_id\":\"P-2026-000123\"}}" > /tmp/mes.jsonl'
docker compose exec namenode hdfs dfs -appendToFile /tmp/mes.jsonl /km0/esdeveniments/2026-09-14/comandes.jsonl
docker compose exec namenode hdfs dfs -cat /km0/esdeveniments/2026-09-14/comandes.jsonl | wc -l   # 3

hdfs fsck és l'eina per veure l'anatomia d'un fitxer: blocs, rèpliques i a quins DataNodes són:

docker compose exec namenode hdfs fsck /km0/esdeveniments/2026-09-14/comandes.jsonl -files -blocks -locations

Sortida resumida:

/km0/esdeveniments/2026-09-14/comandes.jsonl 512 bytes, replicated: replication=2, 1 block(s):  OK
0. BP-1712...:blk_1073741825_1001 len=512 Live_repl=2  [DatanodeInfoWithStorage[172.20.0.4:9866,DS-...,DISK], DatanodeInfoWithStorage[172.20.0.5:9866,DS-...,DISK]]

Status: HEALTHY
 Number of data-nodes:  2
 Number of racks:       1
 Total blocks (validated): 1 (avg. block size 512 B)
 Minimally replicated blocks: 1 (100.0 %)
 Under-replicated blocks: 0 (0.0 %)
 Default replication factor: 2

Veiem un sol bloc (el fitxer fa menys de 128 MB), amb dues rèpliques vives en dues adreces diferents, i un únic rack (no hem configurat la topologia). Un experiment instructiu: docker compose stop datanode-2, esperar uns 10 minuts (o abaixar dfs.namenode.heartbeat.recheck-interval per accelerar) i repetir l'fsck: el bloc apareix com a under-replicated amb Live_repl=1, i el -cat continua funcionant gràcies a la rèplica de datanode-1. En arrencar de nou datanode-2, torna a Live_repl=2.

WebHDFS des de Python. HDFS exposa la seva API completa per HTTP (WebHDFS), cosa que evita instal·lar un client Java a analitica. Les operacions que toquen dades fan servir una redirecció en dos passos que reprodueix el flux de l'apartat 3: es demana al NameNode CREATE (sense enviar dades), el NameNode respon 307 amb l'URL d'un DataNode, i el client envia les dades a aquest DataNode. L'script serveis/analitica/pujar_esdeveniments_hdfs.py puja el fitxer d'esdeveniments del dia:

# km0/serveis/analitica/pujar_esdeveniments_hdfs.py
"""Puja (o afegeix a) /km0/esdeveniments/<dia>/comandes.jsonl fent servir l'API WebHDFS."""
import sys
import requests

NAMENODE = "http://localhost:9870/webhdfs/v1"
USUARI = "km0"

def _url(ruta: str, op: str, **params) -> str:
    query = "&".join(f"{k}={v}" for k, v in {"op": op, "user.name": USUARI, **params}.items())
    return f"{NAMENODE}{ruta}?{query}"

def mkdirs(ruta: str) -> None:
    r = requests.put(_url(ruta, "MKDIRS"))
    r.raise_for_status()
    assert r.json()["boolean"], f"no s'ha pogut crear {ruta}"

def existeix(ruta: str) -> bool:
    return requests.get(_url(ruta, "GETFILESTATUS")).status_code == 200

def crear(ruta: str, dades: bytes, replicacio: int = 2) -> None:
    # Pas 1: el NameNode ens redirigeix al DataNode que escriurà el primer bloc
    r1 = requests.put(_url(ruta, "CREATE", overwrite="false", replication=replicacio),
                      allow_redirects=False)
    assert r1.status_code == 307, r1.text
    datanode_url = r1.headers["Location"]
    # Pas 2: enviem els bytes al DataNode; ell replica per pipeline
    r2 = requests.put(datanode_url, data=dades, headers={"Content-Type": "application/octet-stream"})
    r2.raise_for_status()               # 201 Created

def append(ruta: str, dades: bytes) -> None:
    r1 = requests.post(_url(ruta, "APPEND"), allow_redirects=False)
    assert r1.status_code == 307, r1.text
    r2 = requests.post(r1.headers["Location"], data=dades,
                       headers={"Content-Type": "application/octet-stream"})
    r2.raise_for_status()               # 200 OK

def llistar(ruta: str) -> list[dict]:
    r = requests.get(_url(ruta, "LISTSTATUS"))
    r.raise_for_status()
    return r.json()["FileStatuses"]["FileStatus"]

def llegir(ruta: str) -> bytes:
    r = requests.get(_url(ruta, "OPEN"))   # aquí sí que seguim la redirecció automàticament
    r.raise_for_status()
    return r.content

if __name__ == "__main__":
    dia = sys.argv[1] if len(sys.argv) > 1 else "2026-09-14"
    local = sys.argv[2] if len(sys.argv) > 2 else f"esdeveniments/{dia}/comandes.jsonl"
    remot = f"/km0/esdeveniments/{dia}/comandes.jsonl"

    with open(local, "rb") as f:
        contingut = f.read()

    mkdirs(f"/km0/esdeveniments/{dia}")
    if existeix(remot):
        append(remot, contingut)
        print(f"afegits {len(contingut)} bytes a {remot}")
    else:
        crear(remot, contingut)
        print(f"creat {remot} amb {len(contingut)} bytes")

    for e in llistar(f"/km0/esdeveniments/{dia}"):
        print(f"{e['pathSuffix']:20} {e['length']:>10} bytes  repl={e['replication']}  bloc={e['blockSize']//2**20} MB")

Explicació:

  • _url construeix l'URL WebHDFS: la ruta HDFS va al path i l'operació a op. user.name identifica l'usuari (sense Kerberos, HDFS es refia del nom: és la raó per la qual en producció s'activa l'autenticació, tema de 06-01).
  • crear fa el ball en dos passos amb allow_redirects=False: volem veure el 307 i llegir la capçalera Location, que apunta a http://datanode-1:9864/webhdfs/v1/.... Des de fora de Docker aquest nom no resol; afegeix 127.0.0.1 datanode-1 datanode-2 a /etc/hosts (i el port 9864 només funciona per a datanode-1, per això datanode-2 publica al 9865: en un entorn real el client seria a la mateixa xarxa que els DataNodes, i aquesta molèstia desapareix).
  • append és el mateix patró amb POST. Només un client pot tenir el fitxer obert per a escriptura; un segon APPEND concurrent rep un error AlreadyBeingCreatedException, que és la semàntica d'un sol escriptor de l'apartat 3 traient el cap per HTTP.
  • llistar retorna, entre d'altres, replication i blockSize, que confirmen el que fsck mostrava.

Amb el consumidor Kafka d'analitica acumulant els esdeveniments del dia a esdeveniments/2026-09-14/comandes.jsonl i aquest script executat cada hora (o en tancar el dia), el llac de dades creix un directori per dia, a punt per al Mòdul 5.

Errors Comuns i Consells

  • Tractar HDFS com un disc de xarxa. No és POSIX, no admet escriptures aleatòries ni escriptors concurrents, i cada fitxer petit li costa memòria al NameNode. Agrupa sempre: un fitxer gran per dia i font.
  • Creure que el SecondaryNameNode és un respatller. Només compacta l'edit log. Sense JournalNodes i ZooKeeper no hi ha alta disponibilitat, i sense còpies del directori de metadades (dfs.namenode.name.dir en diversos discos) una avaria pot deixar els blocs orfes i irrecuperables.
  • Factor de replicació més gran que el nombre de DataNodes. Tot queda infrareplicat per sempre i el NameNode reintenta sense descans. Ajusta dfs.replication al clúster.
  • Ignorar la topologia de racks. Sense net.topology.script.file.name, HDFS creu que tot és en un rack i pot posar les tres rèpliques sota el mateix switch.
  • Confiar en la memòria cau del client NFS per coordinar processos en màquines diferents. La semàntica close-to-open no garanteix que la màquina B vegi el que A acaba d'escriure fins que A tanqui i B obri. Si necessites coordinació, fes servir un bloqueig distribuït (etcd, 03-03) o una cua, no el sistema de fitxers.
  • Fer servir WebHDFS sense planificar la resolució de noms dels DataNodes. La redirecció retorna el nom de host del DataNode; el client ha de poder resoldre'l i arribar al seu port. És l'error número u en fer servir WebHDFS des de fora del clúster.
  • Triar Ceph "perquè escala més" per a un problema de 10 TB. La seva complexitat operativa és real. Per a un llac de dades per lots, HDFS (o un magatzem d'objectes, 04-03) és més simple; per a volums compartits moderats, NFS continua sent la resposta correcta.

Exercicis

Exercici 1. Quilòmetre Zero genera al dia 250 000 esdeveniments de comandes (uns 150 MB en JSON Lines) i 12 milions de línies de clics (uns 4 GB, en fitxers d'una hora per a cadascun de 3 servidors web). Un enginyer proposa pujar a HDFS cada fitxer de clics tal qual (3 servidors × 24 hores = 72 fitxers/dia de ~55 MB) i, a més, un fitxer per comanda per als esdeveniments (40 000 fitxers/dia). Calcula per a un any: (a) el nombre de fitxers i de blocs de cada opció, (b) la memòria aproximada del NameNode (150 bytes per fitxer i per bloc), i (c) proposa una organització millor.

Exercici 2. Amb el docker-compose de la pràctica, explica què passa pas a pas, en termes del NameNode, els DataNodes i el client, quan executes hdfs dfs -cat /km0/esdeveniments/2026-09-14/comandes.jsonl mentre datanode-1 està aturat (sense esperar els 10 minuts de detecció). Quina capçalera de WebHDFS faria fallar l'script de Python en la mateixa situació, i com el faries robust?

Exercici 3. Escriu una funció verificar_replicacio(ruta) que faci servir l'operació GETFILEBLOCKLOCATIONS de WebHDFS (op=GETFILEBLOCKLOCATIONS) per llistar, per bloc, els hosts que tenen rèplica, i que retorni la llista de blocs amb menys rèpliques de les configurades. Explica en què s'assembla aquesta comprovació al que fa el mateix NameNode i per què, en un clúster gran, no convé executar-la sobre tot el llac cada minut.

Solucions

Solució 1:

(a) Proposta de l'enginyer, un any (365 dies): clics: 72 × 365 = 26 280 fitxers, cadascun de 55 MB cap en 1 bloc → 26 280 blocs. Esdeveniments: 40 000 × 365 = 14,6 milions de fitxers d'uns 4 KB, 1 bloc cadascun → 14,6 milions de blocs. Total ≈ 14,63 milions de fitxers i altres tants blocs. Organització millor: un fitxer de clics per dia (4 GB → 32 blocs de 128 MB): 365 fitxers i 11 680 blocs; un fitxer d'esdeveniments per dia (150 MB → 2 blocs): 365 fitxers i 730 blocs. Total: 730 fitxers i 12 410 blocs.

(b) Memòria del NameNode: proposta de l'enginyer ≈ (14,63 M fitxers + 14,63 M blocs) × 150 B ≈ 4,4 GB de heap l'any només per a aquesta dada (i creixent; cada rèplica afegeix a més una entrada al mapa de blocs). Organització millor ≈ (730 + 12 410) × 150 B ≈ 2 MB. Tres ordres de magnitud de diferència amb el mateix volum de dades (uns 1,5 TB l'any).

(c) Consolidar per dia i font (/km0/clics/2026-09-14/clics.jsonl, /km0/esdeveniments/2026-09-14/comandes.jsonl), acumulant al consumidor de Kafka i pujant per lots amb append cada hora; opcionalment comprimir amb un format divisible (Parquet o Avro, que el Mòdul 5 llegirà millor que JSON) i aplicar un cicle de vida que arxivi els anys antics amb factor de replicació 2 o amb codificació d'esborrament (Hadoop 3 la suporta).

Solució 2:

El client demana al NameNode getBlockLocations; com que encara no han passat els 10 minuts, el NameNode continua creient que datanode-1 és viu i retorna el bloc amb dues ubicacions, [datanode-1, datanode-2] (o en ordre invers, segons la proximitat calculada). El client intenta connectar amb el primer; si és datanode-1, la connexió falla (rebutjada o amb timeout de dfs.client.socket-timeout, 60 s per defecte, cosa que es pot notar com una pausa), marca aquest DataNode com a "mort per a aquest client" i passa al següent de la llista, datanode-2, que serveix el bloc. El -cat funciona, potser amb un retard inicial. Res no es re-replica fins que el NameNode detecti la caiguda. A WebHDFS, el problema és la capçalera Location de la redirecció: el NameNode pot redirigir un OPEN o CREATE a datanode-1, i el segon pas fallarà amb error de connexió. Per fer-ho robust: capturar requests.ConnectionError al segon pas i reintentar l'operació des del primer pas (el NameNode triarà un altre DataNode en reintentar, sobretot si se li passa el paràmetre excludedatanodes=datanode-1 en lectures), amb un petit nombre de reintents i un timeout explícit a requests (timeout=(5, 60)), que és exactament el patró de reintents amb timeout que formalitzarà 07-04.

Solució 3:

def verificar_replicacio(ruta: str) -> list[dict]:
    estat = requests.get(_url(ruta, "GETFILESTATUS")).json()["FileStatus"]
    esperades = estat["replication"]
    r = requests.get(_url(ruta, "GETFILEBLOCKLOCATIONS"))
    r.raise_for_status()
    blocs = r.json()["BlockLocations"]["BlockLocation"]
    infra = []
    for i, b in enumerate(blocs):
        hosts = b["hosts"]
        print(f"bloc {i}: offset={b['offset']} len={b['length']} hosts={hosts}")
        if len(hosts) < esperades:
            infra.append({"bloc": i, "hosts": hosts, "falten": esperades - len(hosts)})
    return infra

GETFILESTATUS dona el factor de replicació del fitxer i GETFILEBLOCKLOCATIONS (disponible des de Hadoop 2.8/3.x) retorna, per bloc, els hosts amb rèplica. El NameNode fa contínuament una cosa equivalent però des de l'altra banda: creua els informes de blocs que cada DataNode li envia amb el factor esperat i manté una cua de blocs infrareplicats que va corregint amb un límit d'amplada de banda. Executar la comprovació externa sobre tot el llac cada minut és mala idea perquè cada crida de metadades competeix amb les operacions reals pel mateix NameNode (un sol fil d'escriptura de l'espai de noms, amb bloqueig global) i perquè amb desenes de milers de fitxers són desenes de milers de peticions HTTP; el raonable és consultar les mètriques agregades que el NameNode ja calcula (UnderReplicatedBlocks al seu JMX o a hdfs dfsadmin -report) i reservar fsck o aquesta funció per investigar fitxers concrets, un tema que reprendrem a 07-01.

Conclusió

Un sistema de fitxers distribuït ofereix la interfície més familiar de l'emmagatzematge (directoris i fitxers) sobre molts nodes, i el seu caràcter el decideix la semàntica que tria per a la compartició. NFS manté un servidor únic i una consistència close-to-open recolzada en memòries cau de client; continua sent la resposta correcta per a directoris compartits moderats. GFS i HDFS renuncien a POSIX per aconseguir el que NFS no pot: petabytes en màquines barates, amb un NameNode que desa les metadades en memòria, DataNodes que serveixen blocs de 128 MB replicats tres vegades amb consciència de racks, un pipeline d'escriptura síncron, un únic escriptor per fitxer i alta disponibilitat basada en un quòrum de JournalNodes i una elecció a ZooKeeper, que són les eines de 03-03 i 03-04 posades a treballar. Aquest disseny falla, per construcció, amb molts fitxers petits i amb l'accés aleatori. GlusterFS i Ceph eliminen el servidor de metadades i localitzen les dades per càlcul, amb CRUSH com a parent del hashing consistent de 04-01 que a més respecta els dominis de fallada. A Quilòmetre Zero, HDFS és el llac de dades: un fitxer per dia per als esdeveniments de comandes.esdeveniments i els logs de clics, que hem muntat amb apache/hadoop a docker-compose.yml, inspeccionat amb hdfs dfs i hdfs fsck, i alimentat amb pujar_esdeveniments_hdfs.py a través de WebHDFS i la seva redirecció en dos passos.

Les fotos dels productes s'han quedat fora a propòsit: milers de fitxers petits que la web ha de servir directament per HTTP, amb metadades, versions i una durabilitat que no depengui d'un NameNode. Per a això existeix un model diferent, sense directoris ni append, amb una API que s'ha convertit en l'estàndard de facto del núvol: l'emmagatzematge d'objectes, que veurem a continuació amb MinIO i boto3.

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