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

  1. Qué es un sistema de archivos distribuido y qué transparencias ofrece
  2. NFS: el modelo cliente-servidor
  3. GFS y HDFS: ficheros enormes en máquinas baratas
  4. Alta disponibilidad del NameNode
  5. Lo que HDFS no sabe hacer
  6. GlusterFS y Ceph: sin metadatos centralizados
  7. Tabla comparativa y casos de uso
  8. HDFS en Kilómetro Cero: el lago de datos
  9. Práctica: HDFS en docker-compose.yml y la API WebHDFS
  10. Errores Comunes y Consejos
  11. Ejercicios
  12. Conclusión

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

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

  1. 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.replication DataNodes, 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.

  1. 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/2 de 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 ZKFailoverController que 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.

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

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

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

  1. 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 envoltura id_evento/tipo/version/fecha_ms/origen/datos de 02-05). Kafka retiene los eventos unos días; un consumidor de analitica los 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.

  1. Práctica: HDFS en docker-compose.yml y la API WebHDFS

Añ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: 2 porque solo hay dos DataNodes; con el valor por defecto (3) cada bloque quedaría permanentemente "infrarreplicado" y fsck lo mostraría como aviso.
  • ENSURE_NAMENODE_DIR hace que el contenedor ejecute hdfs namenode -format si 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   # 3

hdfs 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 -locations

Salida 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:

  • _url construye la URL WebHDFS: la ruta HDFS va en el path y la operación en op. user.name identifica 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).
  • crear hace el baile en dos pasos con allow_redirects=False: queremos ver el 307 y leer la cabecera Location, que apunta a http://datanode-1:9864/webhdfs/v1/.... Desde fuera de Docker ese nombre no resuelve; añade 127.0.0.1 datanode-1 datanode-2 a /etc/hosts (y el puerto 9864 solo funciona para datanode-1, por eso datanode-2 publica en 9865: en un entorno real el cliente estaría en la misma red que los DataNodes, y esta molestia desaparece).
  • append es el mismo patrón con POST. Solo un cliente puede tener el fichero abierto para escritura; un segundo APPEND concurrente recibe un error AlreadyBeingCreatedException, que es la semántica de un solo escritor del apartado 3 asomando por HTTP.
  • listar devuelve, entre otros, replication y blockSize, que confirman lo que fsck mostraba.

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.dir en 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.replication al 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 infra

GETFILESTATUS 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

Módulo 2: Comunicación en Sistemas Distribuidos

Módulo 3: Consistencia y Replicación

Módulo 4: Almacenamiento Distribuido

Módulo 5: Computación Distribuida

Módulo 6: Seguridad en Sistemas Distribuidos

Módulo 7: Monitoreo y Mantenimiento

Módulo 8: Casos de Estudio y Aplicaciones

© Copyright 2026. Todos los derechos reservados