En la lección anterior aprendimos a decidir qué nodo guarda cada clave. Ahora bajamos un nivel y nos ocupamos del sistema que almacena físicamente los bytes de forma que muchos clientes, desde muchas máquinas, los vean como ficheros normales. Un sistema de archivos distribuido (DFS, Distributed File System) es la forma más antigua de almacenamiento compartido en red y sigue siendo la que más se parece a lo que cualquier programador conoce: directorios, ficheros, open, read, write. Pero bajo esa interfaz familiar hay decisiones muy distintas: NFS se diseñó para que una oficina compartiera el disco de un servidor; GFS y HDFS para que miles de máquinas baratas guardaran petabytes de ficheros enormes que se escriben una vez y se leen muchas; GlusterFS y Ceph para no depender de ningún nodo central. Para Kilómetro Cero la pregunta es concreta: dónde acumular los millones de eventos de pedidos.eventos y los logs de clics de la web para que el Módulo 5 los pueda procesar en masa. La respuesta será HDFS, y lo montaremos en el docker-compose.yml de km0/.
Contenido
- Qué es un sistema de archivos distribuido y qué transparencias ofrece
- NFS: el modelo cliente-servidor
- GFS y HDFS: ficheros enormes en máquinas baratas
- Alta disponibilidad del NameNode
- Lo que HDFS no sabe hacer
- GlusterFS y Ceph: sin metadatos centralizados
- Tabla comparativa y casos de uso
- HDFS en Kilómetro Cero: el lago de datos
- Práctica: HDFS en
docker-compose.ymly la API WebHDFS - Errores Comunes y Consejos
- Ejercicios
- Conclusión
- Qué es un sistema de archivos distribuido y qué transparencias ofrece
Un DFS presenta a sus clientes un espacio de nombres jerárquico (directorios y ficheros) cuyo contenido está almacenado en uno o varios servidores remotos. Lo que lo distingue de "copiar ficheros por la red" es que el acceso se integra en el sistema operativo o en una API equivalente, de modo que los programas no notan (o notan lo menos posible) dónde están los datos. Ese "no notar" es exactamente el concepto de transparencia que definimos en 01-01, y un DFS es el mejor ejemplo para repasarlo:
| Transparencia | Qué significa en un DFS | Quién la ofrece |
|---|---|---|
| De acceso | Se usan las mismas llamadas (open, read) para ficheros locales y remotos |
NFS (montaje), Ceph (CephFS), HDFS solo parcialmente (API propia) |
| De ubicación | El nombre del fichero no revela en qué servidor está (/km0/eventos/…, no //datanode2/…) |
Todos |
| De replicación | El cliente no sabe cuántas copias hay ni cuál lee | HDFS, Ceph, GlusterFS |
| De fallo | Si cae un servidor que guarda una copia, el cliente sigue leyendo | HDFS, Ceph, GlusterFS; NFS no (un servidor) |
| De concurrencia | Varios clientes acceden sin corromperse | Todos, con semánticas muy distintas (apartado 2) |
| De migración / escalado | Se pueden añadir discos o nodos sin cambiar rutas | HDFS, Ceph, GlusterFS |
El punto delicado de cualquier DFS es la semántica de compartición: qué ve un cliente cuando otro cliente está escribiendo el mismo fichero. En un sistema local (semántica UNIX) cada write es visible inmediatamente a cualquier read posterior. Por la red, garantizar eso exige que cada operación viaje al servidor, lo que mata el rendimiento. Cada DFS elige un punto distinto entre rendimiento y semántica, y ese punto explica casi todas sus diferencias.
- NFS: el modelo cliente-servidor
NFS (Network File System, Sun, 1984; hoy en su versión 4.2) es el DFS clásico: un servidor exporta un directorio y los clientes lo montan en su árbol local. Cada operación del cliente se convierte en una llamada RPC (el ONC RPC que vimos en 02-02, con XDR como serialización) al servidor.
sequenceDiagram
participant App as Aplicación
participant K as Kernel cliente<br/>(caché de páginas y atributos)
participant S as Servidor NFS
App->>K: open("/mnt/km0/facturas/P-2026-000123.pdf")
K->>S: LOOKUP + GETATTR (RPC)
S-->>K: file handle + atributos (mtime, tamaño)
App->>K: read(4096 bytes)
K->>S: READ (si no está en caché)
S-->>K: datos
K-->>App: datos (y se guardan en la caché de páginas)
App->>K: close()
K->>S: escritura de páginas sucias (close-to-open)
Cachés de cliente y consistencia. Sin caché, cada read sería un viaje de red. Por eso el cliente NFS cachea datos y atributos, y aquí aparece el compromiso: si el cliente A tiene cacheado un bloque y el cliente B escribe ese mismo bloque en el servidor, A leerá datos viejos hasta que revalide. NFS usa la semántica close-to-open: los cambios de un cliente se garantizan visibles a los demás solo cuando el escritor hace close y el lector hace open después. Entre medias, el cliente revalida los atributos del fichero cada pocos segundos (actimeo, por defecto 3–60 s) y descarta su caché si el mtime cambió. Es consistencia eventual con una ventana de segundos, y en el vocabulario de 03-01 no ofrece ni lectura de las propias escrituras entre máquinas distintas. NFSv4 añade delegaciones: el servidor puede ceder a un cliente el derecho exclusivo sobre un fichero mientras nadie más lo pida, lo que permite cachear con seguridad y revocar cuando aparece otro cliente.
Límites. NFS tiene un solo servidor por exportación: es el límite de capacidad, de rendimiento y el punto único de fallo (se puede montar en alta disponibilidad con un par activo/pasivo y almacenamiento compartido, pero no escala horizontalmente). La versión 4.1 introdujo pNFS (parallel NFS), que separa metadatos de datos y permite leer de varios servidores, pero su adopción es limitada.
Cuándo sigue siendo válido. Muchísimas veces. Compartir el directorio /home de un equipo de desarrollo, dar a varias instancias de una aplicación acceso a los mismos ficheros de configuración o a un directorio de intercambio, montar volúmenes ReadWriteMany en Kubernetes para cargas moderadas: NFS es simple, está en todos los kernels, y hasta unos pocos terabytes y cientos de clientes funciona sin más. Kilómetro Cero lo descartó como almacén de fotos de productos por dos razones que veremos en 04-03: quiere replicación real entre nodos y una API HTTP directa para la web, no un montaje.
- GFS y HDFS: ficheros enormes en máquinas baratas
Google publicó en 2003 el diseño de su Google File System (GFS) y Yahoo lo reimplementó en código abierto como HDFS (Hadoop Distributed File System). Sus premisas rompen con NFS:
- Los ficheros son enormes (gigabytes o terabytes), y son pocos millones, no miles de millones.
- Se escriben una vez (o solo se les añade al final) y se leen muchas veces, casi siempre de forma secuencial y completa.
- El hardware es barato y falla constantemente: con 1000 discos, alguno muere cada día. La tolerancia a fallos es la norma, no la excepción.
- Importa el ancho de banda agregado (leer un terabyte en un minuto desde cien máquinas), no la latencia de una lectura pequeña.
Arquitectura. HDFS tiene dos tipos de nodos:
| Rol | Cuántos | Qué guarda | Qué hace |
|---|---|---|---|
| NameNode | 1 activo (+ 1 en espera, apartado 4) | Los metadatos: árbol de directorios, permisos, y para cada fichero la lista de bloques y en qué DataNodes está cada réplica. Todo en memoria | Atiende las operaciones de espacio de nombres (mkdir, ls, open), decide dónde van los bloques nuevos, ordena re-replicar los que pierden copias |
| DataNode | Decenas a miles | Los bloques de datos, como ficheros normales en su disco local | Sirve lecturas y escrituras de bloques directamente a los clientes; envía latidos y el informe de sus bloques al NameNode |
Los conceptos clave:
- Bloques grandes: por defecto 128 MB (frente a los 4 KB de un sistema de ficheros local). Un fichero de 1 GB son 8 bloques. Los bloques grandes reducen el número de metadatos que el NameNode guarda en memoria y hacen que cada lectura sea una transferencia secuencial larga, que es lo que los discos hacen bien. Un fichero más pequeño que un bloque ocupa solo lo que mide (no desperdicia 128 MB), pero sí consume una entrada de metadatos completa.
- Factor de replicación: cada bloque se copia en
dfs.replicationDataNodes, por defecto 3. La replicación es por bloque, no por fichero, y el NameNode la vigila: si un DataNode deja de enviar latidos (10 minutos por defecto), todos sus bloques quedan "infrarreplicados" y el NameNode ordena copiarlos desde las réplicas supervivientes a otros nodos. - Rack awareness: el NameNode conoce en qué armario (rack) está cada DataNode. Con 3 réplicas, coloca la primera en el nodo del cliente (si es un DataNode), la segunda en un nodo de otro rack y la tercera en otro nodo del mismo rack que la segunda. Así un fallo de un rack entero (un switch) no pierde ningún bloque, y solo una de las tres copias cruza entre racks (ancho de banda entre racks, que es el escaso).
- Write-once / append: un fichero se crea, se escribe (por un solo escritor) y se cierra; después es inmutable, salvo que se añada al final (
append). No hay escrituras en medio del fichero. Esto simplifica enormemente la consistencia: no hay dos escritores concurrentes sobre los mismos bytes y las réplicas de un bloque cerrado son idénticas para siempre.
Flujo de escritura y lectura. Lo esencial es que los datos nunca pasan por el NameNode; solo los metadatos:
sequenceDiagram
participant C as Cliente HDFS
participant NN as NameNode
participant D1 as DataNode 1
participant D2 as DataNode 2
participant D3 as DataNode 3
Note over C,D3: ESCRITURA de /km0/eventos/2026-09-14/pedidos.jsonl
C->>NN: create(ruta)
NN-->>C: ok (fichero en construcción)
C->>NN: addBlock()
NN-->>C: bloque B1 → [D1, D2, D3] (rack awareness)
C->>D1: paquetes de B1
D1->>D2: reenvía (pipeline)
D2->>D3: reenvía (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] … (ordenados por cercanía)
C->>D1: read B1
C->>D2: read B2
En la escritura, el cliente pide al NameNode dónde poner el bloque, y envía los datos al primer DataNode, que los reenvía al segundo, que los reenvía al tercero (pipeline de replicación): el cliente solo emite los datos una vez y la replicación consume ancho de banda entre DataNodes, no del cliente. El bloque se considera escrito cuando los tres han confirmado: es replicación síncrona en el sentido de 03-04, y por eso HDFS es CP (una escritura no se acepta si no puede alcanzar el número mínimo de réplicas, dfs.namenode.replication.min, que por defecto es 1). En la lectura, el cliente obtiene la lista de bloques con sus ubicaciones y lee cada bloque directamente del DataNode más cercano; si uno falla, pasa al siguiente de la lista. Cada bloque lleva sumas de comprobación (CRC32 cada 512 bytes) que el cliente verifica: un bloque corrupto se descarta, se lee de otra réplica y se notifica al NameNode para que lo re-replique.
- Alta disponibilidad del NameNode
En el diseño original, el NameNode era un punto único de fallo: si caía, el clúster entero quedaba inaccesible, aunque los datos siguieran intactos en los DataNodes. Desde Hadoop 2 existe la configuración de alta disponibilidad (HA), que aplica lo que ya sabemos de 03-03 y 03-04:
- Dos (o más) NameNodes: uno activo y otro en espera (standby). Ambos tienen el espacio de nombres en memoria.
- El registro de cambios (edit log) del activo se escribe en un quórum de JournalNodes (normalmente 3): la operación de metadatos se confirma cuando la mayoría lo ha persistido. El NameNode en espera lee ese registro continuamente y aplica los cambios, de modo que su copia va solo unos milisegundos por detrás. Este quórum es la aplicación directa del
W > N/2de 03-04. - Los DataNodes envían latidos e informes de bloques a ambos NameNodes, para que el standby conozca la ubicación de cada bloque y pueda tomar el mando sin reconstruirla.
- ZooKeeper (03-03) decide quién es el activo: cada NameNode tiene un proceso
ZKFailoverControllerque mantiene un nodo efímero en ZooKeeper; si el activo deja de renovarlo, el standby adquiere el bloqueo y se promociona. - Fencing: antes de promocionarse, el nuevo activo se asegura de que el antiguo no pueda seguir escribiendo en los JournalNodes (que solo aceptan escrituras del NameNode con la época más alta, igual que los términos de Raft), y opcionalmente lo mata por SSH. Es la defensa contra el "cerebro dividido" que ya discutimos en 03-04.
flowchart LR
ZK[(ZooKeeper<br/>elección de activo)]
NN1[NameNode activo] --- ZK
NN2[NameNode standby] --- ZK
NN1 -- escribe edits --> J1[(JournalNode 1)]
NN1 -- escribe edits --> J2[(JournalNode 2)]
NN1 -- escribe edits --> J3[(JournalNode 3)]
NN2 -. lee edits .-> J1
NN2 -. lee edits .-> J2
NN2 -. lee edits .-> J3
D1[DataNode] -- latidos e informes --> NN1
D1 -- latidos e informes --> NN2
Con HA, el fallo del NameNode activo se resuelve en decenas de segundos sin intervención. Sin HA (como en el docker-compose de práctica), un SecondaryNameNode solo compacta el edit log; no es un respaldo y no puede tomar el mando, un malentendido tan frecuente que merece subrayarse.
- Lo que HDFS no sabe hacer
Sus premisas son sus límites, y conviene tenerlos claros antes de meter en HDFS algo que no encaja:
- Muchos ficheros pequeños. Cada fichero, directorio y bloque ocupa unos 150 bytes en la memoria del NameNode. Cien millones de ficheros de 10 KB son 15 GB de heap y 1 TB de datos que en bloques de 128 MB habrían sido 8 000 entradas. Además, cada lectura de un fichero pequeño paga una llamada al NameNode y una conexión a un DataNode para leer unos pocos KB. Las fotos de los productos (miles de ficheros de 200 KB) son un mal caso para HDFS; por eso van a un almacén de objetos (04-03). Si hay que meter ficheros pequeños en HDFS, se agrupan (ficheros
SequenceFile, HAR o, mejor, particiones diarias de un solo fichero grande, como haremos con los eventos). - Acceso aleatorio de baja latencia. Leer un registro concreto exige localizar el bloque, abrir una conexión y leer al menos un paquete; decenas de milisegundos. HDFS no es una base de datos: HBase se construyó encima precisamente para dar acceso aleatorio, gestionando sus propios índices y ficheros grandes.
- Escritores concurrentes y modificaciones en medio del fichero. Un fichero tiene un único escritor y solo admite
append. Un log que muchos servicios escriben a la vez debe pasar antes por Kafka (02-04) y volcarse a HDFS por lotes. - Semántica POSIX completa. No hay
mmap, ni bloqueos, ni escrituras parciales; el "montaje" con NFS Gateway o FUSE es una capa de compatibilidad con limitaciones.
- GlusterFS y Ceph: sin metadatos centralizados
El NameNode, aun con HA, es un límite: todo el espacio de nombres cabe en la memoria de una máquina y toda operación de metadatos pasa por ella. Una segunda familia de DFS elimina el servidor de metadatos y localiza los datos por cálculo, con las mismas ideas del hashing consistente de 04-01.
GlusterFS agrupa directorios exportados por varios servidores (bricks) en un volumen. No hay servidor de metadatos: la ubicación de cada fichero se calcula con un hash de su nombre sobre un rango asignado a cada brick (elastic hashing), y el cliente, que conoce la configuración del volumen, habla directamente con el brick correcto. Los volúmenes pueden ser distribuidos (cada fichero en un brick), replicados (cada fichero en N bricks) o dispersos (con codificación de borrado, que veremos en 04-03), y se combinan. Es sencillo, se monta como un sistema de ficheros normal (FUSE o NFS) y escala bien para ficheros medianos y grandes; su punto débil son las operaciones de directorio (un ls debe preguntar a todos los bricks) y la autorreparación tras fallos.
Ceph es más ambicioso: un almacén de objetos distribuido (RADOS) sobre el que se construyen tres interfaces: bloques (RBD, discos virtuales para máquinas virtuales y Kubernetes), objetos (RGW, compatible con S3, 04-03) y ficheros (CephFS). Sus piezas:
- OSD (Object Storage Daemon): un proceso por disco, que guarda objetos y se replica con otros OSD entre pares.
- Monitores (MON): un pequeño quórum (Paxos, 03-03) que mantiene el mapa del clúster (qué OSD existen y su estado). No guardan datos ni metadatos de ficheros.
- CRUSH (Controlled Replication Under Scalable Hashing): el algoritmo que, a partir del nombre de un objeto y del mapa del clúster, calcula en qué OSD están sus réplicas. Es un pariente directo del hashing consistente de 04-01, con una diferencia importante: CRUSH entiende la topología (disco → servidor → rack → sala) y unas reglas ("tres réplicas en tres racks distintos"), de manera que las réplicas no solo se reparten uniformemente sino que respetan los dominios de fallo. Cualquier cliente con el mapa calcula la ubicación sin preguntar a nadie: no hay NameNode, y añadir un OSD mueve solo la fracción proporcional de objetos, como en el anillo.
- MDS (Metadata Server): solo para CephFS, gestiona el árbol de directorios; puede haber varios activos, repartiéndose subárboles dinámicamente.
Ceph es el motor de almacenamiento de muchos clouds privados (OpenStack, Proxmox, Kubernetes con Rook). Es también notablemente más complejo de operar que HDFS o GlusterFS.
- Tabla comparativa y casos de uso
| NFS | HDFS | GlusterFS | Ceph | |
|---|---|---|---|---|
| Metadatos | En el único servidor | NameNode centralizado (en memoria) | Sin servidor: hash del nombre | Sin servidor para objetos (CRUSH); MDS para CephFS |
| Datos | Un servidor | Bloques de 128 MB replicados en DataNodes | Ficheros enteros en bricks | Objetos en OSD, con CRUSH |
| Replicación | No (o activo/pasivo externo) | Por bloque, síncrona, rack aware | Por fichero, síncrona | Por objeto, síncrona, topología con reglas |
| Semántica | Close-to-open, POSIX aproximado | Write-once + append, un escritor | POSIX aproximado | POSIX (CephFS), bloque, objeto |
| Escala | Terabytes, cientos de clientes | Petabytes, miles de nodos; millones de ficheros | Petabytes | Petabytes a exabytes |
| Ficheros pequeños | Bien | Mal | Regular | Bien (como objetos) |
| Acceso aleatorio | Bien | Mal | Bien | Bien |
| Interfaz | Montaje del SO | API Java/CLI, WebHDFS (HTTP), FUSE limitado | Montaje (FUSE/NFS) | Montaje, S3, disco de bloques |
| Complejidad operativa | Muy baja | Media | Media | Alta |
| Caso de uso típico | Directorios compartidos, /home, volúmenes ReadWriteMany moderados |
Lago de datos para procesamiento por lotes (Módulo 5) | Almacén de ficheros de tamaño medio, contenido web, backups | Infraestructura de almacenamiento unificada de un cloud privado |
- HDFS en Kilómetro Cero: el lago de datos
En 01-06 fijamos que analitica procesaría por lotes los datos históricos y en streaming los recientes. Los datos históricos necesitan un sitio donde acumularse durante años, barato por terabyte, tolerante a fallos y optimizado para que el Módulo 5 los lea enteros de forma paralela. Ese sitio es el lago de datos (data lake) sobre HDFS, con dos fuentes:
- Los eventos del tópico
pedidos.eventos(pedido.creado,stock.reservado,pago.confirmado, … con la envolturaid_evento/tipo/version/fecha_ms/origen/datosde 02-05). Kafka retiene los eventos unos días; un consumidor deanaliticalos vuelca a HDFS por lotes, en un fichero por día y tipo:/km0/eventos/2026-09-14/pedidos.jsonl. Cada línea es un evento en JSON (formato JSON Lines). Con 40 000 pedidos y unos 6 eventos por pedido, un día de campaña son unos 250 000 eventos y unos 150 MB: uno o dos bloques, un tamaño ideal para HDFS. - Los logs de clics de la web (página vista, producto visto, añadido al carrito), que el servidor web escribe en ficheros rotados cada hora y que un proceso sube a
/km0/clics/2026-09-14/hora=13/web-01.jsonl.
La organización por directorios de fecha (fecha=2026-09-14/) no es casual: es la partición por rango de 04-01 aplicada a ficheros, y permitirá al Módulo 5 procesar "solo la Semana de la Vendimia" sin leer el resto del lago. Los eventos son inmutables (write-once encaja a la perfección), el acceso es secuencial y masivo, y nadie necesita leer "el evento 123" con baja latencia: para eso está la base de datos de pedidos (04-04). Esta es la división que cierra el módulo: HDFS para lo masivo y frío, la base de datos para lo operativo, los objetos para los ficheros que la web sirve.
- Práctica: HDFS en
docker-compose.yml y la API WebHDFS
docker-compose.yml y la API WebHDFSAñadimos al docker-compose.yml de km0/ un NameNode y dos DataNodes con la imagen oficial apache/hadoop. La imagen se configura con variables de entorno cuyo nombre replica el fichero XML de Hadoop (CORE-SITE.XML_fs.defaultFS equivale a la propiedad fs.defaultFS de core-site.xml):
# km0/docker-compose.yml (fragmento)
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" # solo tenemos 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; nunca en producción
services:
namenode:
image: apache/hadoop:3.4.1
hostname: namenode
command: ["hdfs", "namenode"]
environment:
<<: *hadoop-env
ENSURE_NAMENODE_DIR: /tmp/hadoop-root/dfs/name # formatea el NameNode la primera vez
ports:
- "9870:9870" # interfaz web y 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 (necesario para leer/escribir datos desde fuera)
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:Puntos a entender:
dfs.replication: 2porque solo hay dos DataNodes; con el valor por defecto (3) cada bloque quedaría permanentemente "infrarreplicado" yfscklo mostraría como aviso.ENSURE_NAMENODE_DIRhace que el contenedor ejecutehdfs namenode -formatsi el directorio está vacío. Formatear un NameNode con datos borra el espacio de nombres (los bloques quedan huérfanos en los DataNodes), así que va en un volumen persistente.- El puerto 9870 es la consola web del NameNode (
http://localhost:9870), donde se ven los DataNodes vivos, la capacidad y el explorador de ficheros.
Arrancamos y probamos con la CLI, ejecutada dentro del contenedor del NameNode:
docker compose up -d namenode datanode-1 datanode-2
docker compose exec namenode hdfs dfsadmin -report | head -20 # 2 DataNodes vivos, capacidad
# Espacio de nombres del lago de datos
docker compose exec namenode hdfs dfs -mkdir -p /km0/eventos/2026-09-14 /km0/clics
docker compose exec namenode hdfs dfs -ls /km0
# Subir un fichero local (dentro del contenedor) y leerlo
docker compose exec namenode bash -c 'printf "%s\n" \
"{\"id_evento\":\"e-1\",\"tipo\":\"pedido.creado\",\"version\":2,\"fecha_ms\":1789380000000,\"origen\":\"pedidos\",\"datos\":{\"pedido_id\":\"P-2026-000123\",\"cliente\":\"Ana\"}}" \
"{\"id_evento\":\"e-2\",\"tipo\":\"stock.reservado\",\"version\":1,\"fecha_ms\":1789380000450,\"origen\":\"inventario\",\"datos\":{\"pedido_id\":\"P-2026-000123\",\"producto\":\"queso-curado\",\"cantidad\":2}}" \
> /tmp/pedidos.jsonl'
docker compose exec namenode hdfs dfs -put /tmp/pedidos.jsonl /km0/eventos/2026-09-14/pedidos.jsonl
docker compose exec namenode hdfs dfs -ls -h /km0/eventos/2026-09-14
docker compose exec namenode hdfs dfs -cat /km0/eventos/2026-09-14/pedidos.jsonl
# Añadir al final (append) y comprobar
docker compose exec namenode bash -c 'echo "{\"id_evento\":\"e-3\",\"tipo\":\"pago.confirmado\",\"version\":1,\"fecha_ms\":1789380002000,\"origen\":\"pagos\",\"datos\":{\"pedido_id\":\"P-2026-000123\"}}" > /tmp/mas.jsonl'
docker compose exec namenode hdfs dfs -appendToFile /tmp/mas.jsonl /km0/eventos/2026-09-14/pedidos.jsonl
docker compose exec namenode hdfs dfs -cat /km0/eventos/2026-09-14/pedidos.jsonl | wc -l # 3hdfs fsck es la herramienta para ver la anatomía de un fichero: bloques, réplicas y en qué DataNodes están:
docker compose exec namenode hdfs fsck /km0/eventos/2026-09-14/pedidos.jsonl -files -blocks -locationsSalida resumida:
/km0/eventos/2026-09-14/pedidos.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
Vemos un solo bloque (el fichero mide menos de 128 MB), con dos réplicas vivas en dos direcciones distintas, y un único rack (no hemos configurado la topología). Un experimento instructivo: docker compose stop datanode-2, esperar unos 10 minutos (o bajar dfs.namenode.heartbeat.recheck-interval para acelerar) y repetir el fsck: el bloque aparece como under-replicated con Live_repl=1, y el -cat sigue funcionando gracias a la réplica de datanode-1. Al arrancar de nuevo datanode-2, vuelve a Live_repl=2.
WebHDFS desde Python. HDFS expone su API completa por HTTP (WebHDFS), lo que evita instalar un cliente Java en analitica. Las operaciones que tocan datos usan una redirección en dos pasos que reproduce el flujo del apartado 3: se pide al NameNode CREATE (sin enviar datos), el NameNode responde 307 con la URL de un DataNode, y el cliente envía los datos a ese DataNode. El script servicios/analitica/subir_eventos_hdfs.py sube el fichero de eventos del día:
# km0/servicios/analitica/subir_eventos_hdfs.py
"""Sube (o añade a) /km0/eventos/<día>/pedidos.jsonl usando la API WebHDFS."""
import sys
import requests
NAMENODE = "http://localhost:9870/webhdfs/v1"
USUARIO = "km0"
def _url(ruta: str, op: str, **params) -> str:
query = "&".join(f"{k}={v}" for k, v in {"op": op, "user.name": USUARIO, **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 se pudo crear {ruta}"
def existe(ruta: str) -> bool:
return requests.get(_url(ruta, "GETFILESTATUS")).status_code == 200
def crear(ruta: str, datos: bytes, replicacion: int = 2) -> None:
# Paso 1: el NameNode nos redirige al DataNode que escribirá el primer bloque
r1 = requests.put(_url(ruta, "CREATE", overwrite="false", replication=replicacion),
allow_redirects=False)
assert r1.status_code == 307, r1.text
datanode_url = r1.headers["Location"]
# Paso 2: enviamos los bytes al DataNode; él replica por pipeline
r2 = requests.put(datanode_url, data=datos, headers={"Content-Type": "application/octet-stream"})
r2.raise_for_status() # 201 Created
def append(ruta: str, datos: 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=datos,
headers={"Content-Type": "application/octet-stream"})
r2.raise_for_status() # 200 OK
def listar(ruta: str) -> list[dict]:
r = requests.get(_url(ruta, "LISTSTATUS"))
r.raise_for_status()
return r.json()["FileStatuses"]["FileStatus"]
def leer(ruta: str) -> bytes:
r = requests.get(_url(ruta, "OPEN")) # aquí sí seguimos la redirección automáticamente
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"eventos/{dia}/pedidos.jsonl"
remoto = f"/km0/eventos/{dia}/pedidos.jsonl"
with open(local, "rb") as f:
contenido = f.read()
mkdirs(f"/km0/eventos/{dia}")
if existe(remoto):
append(remoto, contenido)
print(f"añadidos {len(contenido)} bytes a {remoto}")
else:
crear(remoto, contenido)
print(f"creado {remoto} con {len(contenido)} bytes")
for e in listar(f"/km0/eventos/{dia}"):
print(f"{e['pathSuffix']:20} {e['length']:>10} bytes repl={e['replication']} bloque={e['blockSize']//2**20} MB")Explicación:
_urlconstruye la URL WebHDFS: la ruta HDFS va en el path y la operación enop.user.nameidentifica al usuario (sin Kerberos, HDFS se fía del nombre: es la razón por la que en producción se activa la autenticación, tema de 06-01).crearhace el baile en dos pasos conallow_redirects=False: queremos ver el307y leer la cabeceraLocation, que apunta ahttp://datanode-1:9864/webhdfs/v1/.... Desde fuera de Docker ese nombre no resuelve; añade127.0.0.1 datanode-1 datanode-2a/etc/hosts(y el puerto 9864 solo funciona paradatanode-1, por esodatanode-2publica en 9865: en un entorno real el cliente estaría en la misma red que los DataNodes, y esta molestia desaparece).appendes el mismo patrón conPOST. Solo un cliente puede tener el fichero abierto para escritura; un segundoAPPENDconcurrente recibe un errorAlreadyBeingCreatedException, que es la semántica de un solo escritor del apartado 3 asomando por HTTP.listardevuelve, entre otros,replicationyblockSize, que confirman lo quefsckmostraba.
Con el consumidor Kafka de analitica acumulando los eventos del día en eventos/2026-09-14/pedidos.jsonl y este script ejecutado cada hora (o al cerrar el día), el lago de datos crece un directorio por día, listo para el Módulo 5.
Errores Comunes y Consejos
- Tratar HDFS como un disco de red. No es POSIX, no admite escrituras aleatorias ni escritores concurrentes, y cada fichero pequeño le cuesta memoria al NameNode. Agrupa siempre: un fichero grande por día y fuente.
- Creer que el SecondaryNameNode es un respaldo. Solo compacta el edit log. Sin JournalNodes y ZooKeeper no hay alta disponibilidad, y sin copias del directorio de metadatos (
dfs.namenode.name.diren varios discos) una avería puede dejar los bloques huérfanos e irrecuperables. - Factor de replicación mayor que el número de DataNodes. Todo queda infrarreplicado para siempre y el NameNode reintenta sin descanso. Ajusta
dfs.replicational clúster. - Ignorar la topología de racks. Sin
net.topology.script.file.name, HDFS cree que todo está en un rack y puede poner las tres réplicas bajo el mismo switch. - Confiar en la caché del cliente NFS para coordinar procesos en máquinas distintas. La semántica close-to-open no garantiza que la máquina B vea lo que A acaba de escribir hasta que A cierre y B abra. Si necesitas coordinación, usa un bloqueo distribuido (etcd, 03-03) o una cola, no el sistema de ficheros.
- Usar WebHDFS sin planificar la resolución de nombres de los DataNodes. La redirección devuelve el nombre de host del DataNode; el cliente debe poder resolverlo y alcanzar su puerto. Es el error número uno al usar WebHDFS desde fuera del clúster.
- Elegir Ceph "porque escala más" para un problema de 10 TB. Su complejidad operativa es real. Para un lago de datos por lotes, HDFS (o un almacén de objetos, 04-03) es más simple; para volúmenes compartidos moderados, NFS sigue siendo la respuesta correcta.
Ejercicios
Ejercicio 1. Kilómetro Cero genera al día 250 000 eventos de pedidos (unos 150 MB en JSON Lines) y 12 millones de líneas de clics (unos 4 GB, en ficheros de una hora por cada uno de 3 servidores web). Un ingeniero propone subir a HDFS cada fichero de clics tal cual (3 servidores × 24 horas = 72 ficheros/día de ~55 MB) y, además, un fichero por pedido para los eventos (40 000 ficheros/día). Calcula para un año: (a) el número de ficheros y de bloques de cada opción, (b) la memoria aproximada del NameNode (150 bytes por fichero y por bloque), y (c) propón una organización mejor.
Ejercicio 2. Con el docker-compose de la práctica, explica qué ocurre paso a paso, en términos del NameNode, los DataNodes y el cliente, cuando ejecutas hdfs dfs -cat /km0/eventos/2026-09-14/pedidos.jsonl mientras datanode-1 está parado (sin esperar los 10 minutos de detección). ¿Qué cabecera de WebHDFS haría fallar el script de Python en la misma situación, y cómo lo harías robusto?
Ejercicio 3. Escribe una función verificar_replicacion(ruta) que use la operación GETFILEBLOCKLOCATIONS de WebHDFS (op=GETFILEBLOCKLOCATIONS) para listar, por bloque, los hosts que tienen réplica, y que devuelva la lista de bloques con menos réplicas de las configuradas. Explica en qué se parece esta comprobación a lo que hace el propio NameNode y por qué, en un clúster grande, no conviene ejecutarla sobre todo el lago cada minuto.
Soluciones
Solución 1:
(a) Propuesta del ingeniero, un año (365 días): clics: 72 × 365 = 26 280 ficheros, cada uno de 55 MB cabe en 1 bloque → 26 280 bloques. Eventos: 40 000 × 365 = 14,6 millones de ficheros de unos 4 KB, 1 bloque cada uno → 14,6 millones de bloques. Total ≈ 14,63 millones de ficheros y otros tantos bloques. Organización mejor: un fichero de clics por día (4 GB → 32 bloques de 128 MB): 365 ficheros y 11 680 bloques; un fichero de eventos por día (150 MB → 2 bloques): 365 ficheros y 730 bloques. Total: 730 ficheros y 12 410 bloques.
(b) Memoria del NameNode: propuesta del ingeniero ≈ (14,63 M ficheros + 14,63 M bloques) × 150 B ≈ 4,4 GB de heap al año solo para este dato (y creciendo; cada réplica añade además una entrada en el mapa de bloques). Organización mejor ≈ (730 + 12 410) × 150 B ≈ 2 MB. Tres órdenes de magnitud de diferencia con el mismo volumen de datos (unos 1,5 TB al año).
(c) Consolidar por día y fuente (/km0/clics/2026-09-14/clics.jsonl, /km0/eventos/2026-09-14/pedidos.jsonl), acumulando en el consumidor de Kafka y subiendo por lotes con append cada hora; opcionalmente comprimir con un formato divisible (Parquet o Avro, que el Módulo 5 leerá mejor que JSON) y aplicar un ciclo de vida que archive los años antiguos con factor de replicación 2 o con codificación de borrado (Hadoop 3 la soporta).
Solución 2:
El cliente pide al NameNode getBlockLocations; como aún no han pasado los 10 minutos, el NameNode sigue creyendo que datanode-1 está vivo y devuelve el bloque con dos ubicaciones, [datanode-1, datanode-2] (o en orden inverso, según la cercanía calculada). El cliente intenta conectar con el primero; si es datanode-1, la conexión falla (rechazada o con timeout de dfs.client.socket-timeout, 60 s por defecto, lo que puede notarse como una pausa), marca ese DataNode como "muerto para este cliente" y pasa al siguiente de la lista, datanode-2, que sirve el bloque. El -cat funciona, quizá con un retraso inicial. Nada se re-replica hasta que el NameNode detecte la caída. En WebHDFS, el problema es la cabecera Location de la redirección: el NameNode puede redirigir un OPEN o CREATE a datanode-1, y el segundo paso fallará con error de conexión. Para hacerlo robusto: capturar requests.ConnectionError en el segundo paso y reintentar la operación desde el primer paso (el NameNode elegirá otro DataNode al reintentar, sobre todo si se le pasa el parámetro excludedatanodes=datanode-1 en lecturas), con un pequeño número de reintentos y un timeout explícito en requests (timeout=(5, 60)), que es exactamente el patrón de reintentos con timeout que formalizará 07-04.
Solución 3:
def verificar_replicacion(ruta: str) -> list[dict]:
estado = requests.get(_url(ruta, "GETFILESTATUS")).json()["FileStatus"]
esperadas = estado["replication"]
r = requests.get(_url(ruta, "GETFILEBLOCKLOCATIONS"))
r.raise_for_status()
bloques = r.json()["BlockLocations"]["BlockLocation"]
infra = []
for i, b in enumerate(bloques):
hosts = b["hosts"]
print(f"bloque {i}: offset={b['offset']} len={b['length']} hosts={hosts}")
if len(hosts) < esperadas:
infra.append({"bloque": i, "hosts": hosts, "faltan": esperadas - len(hosts)})
return infraGETFILESTATUS da el factor de replicación del fichero y GETFILEBLOCKLOCATIONS (disponible desde Hadoop 2.8/3.x) devuelve, por bloque, los hosts con réplica. El NameNode hace continuamente algo equivalente pero desde el otro lado: cruza los informes de bloques que cada DataNode le envía con el factor esperado y mantiene una cola de bloques infrarreplicados que va corrigiendo con un límite de ancho de banda. Ejecutar la comprobación externa sobre todo el lago cada minuto es mala idea porque cada llamada de metadatos compite con las operaciones reales por el mismo NameNode (un solo hilo de escritura del espacio de nombres, con bloqueo global) y porque con decenas de miles de ficheros son decenas de miles de peticiones HTTP; lo razonable es consultar las métricas agregadas que el NameNode ya calcula (UnderReplicatedBlocks en su JMX o en hdfs dfsadmin -report) y reservar fsck o esta función para investigar ficheros concretos, un tema que retomaremos en 07-01.
Conclusión
Un sistema de archivos distribuido ofrece la interfaz más familiar del almacenamiento (directorios y ficheros) sobre muchos nodos, y su carácter lo decide la semántica que elige para la compartición. NFS mantiene un servidor único y una consistencia close-to-open apoyada en cachés de cliente; sigue siendo la respuesta correcta para directorios compartidos moderados. GFS y HDFS renuncian a POSIX para conseguir lo que NFS no puede: petabytes en máquinas baratas, con un NameNode que guarda los metadatos en memoria, DataNodes que sirven bloques de 128 MB replicados tres veces con conciencia de racks, un pipeline de escritura síncrono, un único escritor por fichero y alta disponibilidad basada en un quórum de JournalNodes y una elección en ZooKeeper, que son las herramientas de 03-03 y 03-04 puestas a trabajar. Ese diseño falla, por construcción, con muchos ficheros pequeños y con el acceso aleatorio. GlusterFS y Ceph eliminan el servidor de metadatos y localizan los datos por cálculo, con CRUSH como pariente del hashing consistente de 04-01 que además respeta los dominios de fallo. En Kilómetro Cero, HDFS es el lago de datos: un fichero por día para los eventos de pedidos.eventos y los logs de clics, que hemos montado con apache/hadoop en docker-compose.yml, inspeccionado con hdfs dfs y hdfs fsck, y alimentado con subir_eventos_hdfs.py a través de WebHDFS y su redirección en dos pasos.
Las fotos de los productos se han quedado fuera a propósito: miles de ficheros pequeños que la web debe servir directamente por HTTP, con metadatos, versiones y una durabilidad que no dependa de un NameNode. Para eso existe un modelo distinto, sin directorios ni append, con una API que se ha convertido en el estándar de facto de la nube: el almacenamiento de objetos, que veremos a continuación con MinIO y boto3.
Curso de Arquitecturas Distribuidas
Módulo 1: Introducción a los Sistemas Distribuidos
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
