skip to content

In HDFS, what does the NameNode store and what do the DataNodes store?

level: juniorimportance: must knowfreq 60%

answer

  1. one master, many workers
  2. one side holds names, the other bytes
  3. the master is never in the data path
  4. one map is rebuilt, never saved
  5. block reports supply what fsimage omits

basics

~10 s

The NameNode keeps the filesystem namespace in memory: directories, files, permissions, and which blocks make up each file. DataNodes store the actual block bytes on local disks. File data never flows through the NameNode.

solid answer

~50 s

HDFS separates **metadata** from **data**. The NameNode is a single master process that holds the entire namespace in RAM — the directory tree, file names, permissions, replication factor, and the ordered list of block IDs that make up each file. It persists that namespace to an `fsimage` file plus an append-only edit log, and it does **not** persist block locations: DataNodes tell it where blocks live through periodic block reports, so the location map is rebuilt at startup. DataNodes are the workers: each stores block replicas as ordinary files on its local disks, sends heartbeats so the NameNode knows it is alive, and serves bytes directly to clients. A client asks the NameNode *where* a block is, then talks straight to the DataNode — so the NameNode is never in the data path, which is why one master can front petabytes.

code

bash · 3 lines
bash
hdfs dfs -put events.log /data/events.log
# which blocks make up the file, and which DataNodes hold each replica
hdfs fsck /data/events.log -files -blocks -locations

go deeper

for a junior

Be ready to state the split in one breath: NameNode holds names, structure and the block list in memory; DataNodes hold the actual bytes. Knowing that clients read straight from DataNodes already puts you ahead.

for a middle

Explain the persistence mechanics — fsimage plus edit log, checkpointing to bound restart time, and why block locations are rebuilt from DataNode block reports rather than saved. Mention safe mode as the visible consequence.

for a senior

Show you have operated this: heartbeat timeouts triggering re-replication storms, safe-mode duration scaling with block count on a cold start, JBOD over RAID because replication is the redundancy, and NameNode heap as the real capacity ceiling.

for a principal

Own the architectural tradeoff. A single in-memory metadata master buys strong, cheap namespace consistency and pays with a heap-bounded object count and a failure domain that needs HA or federation. Be able to argue when that trade stops being worth it versus object storage.

## The core split HDFS is built on one structural decision: **metadata lives in one place, data lives everywhere else**. The NameNode is the single authority on what files exist and what blocks they are made of. The DataNodes are dumb-by-design byte stores that hold block replicas and hand them to whoever asks. A client that wants to read a file gets *directions* from the NameNode and then fetches the bytes from DataNodes itself. No file content ever passes through the NameNode. This is what lets a single-master design scale: metadata is tiny relative to data, and the throughput-hungry part of the workload (moving terabytes) is spread across every machine in the cluster. ## What the NameNode holds in memory The NameNode keeps the whole **namespace** resident in the JVM heap: - the directory tree, with file and directory names - per-file attributes: owner, group, permissions, modification time, replication factor, block size - for each file, the ordered list of **block IDs** that compose it - the **block map**: for each block ID, which DataNodes currently report a replica Because it is all in RAM, namespace operations (`ls`, `mkdir`, `rename`, resolving a path) are fast — and because it is all in RAM, the number of files, directories and blocks a cluster can hold is bounded by the NameNode's heap. That single fact drives the small-files problem, NameNode heap sizing, and HDFS Federation. ## What the NameNode persists Two artefacts on the NameNode's local disk: - **`fsimage`** — a checkpoint: a serialized snapshot of the namespace at some point in time. - **the edit log (`edits`)** — an append-only journal of every namespace mutation since that checkpoint. On startup the NameNode loads `fsimage`, replays `edits`, and has the namespace back. Because replaying a long edit log is slow, the log is periodically merged into a new `fsimage` — historically by the Secondary NameNode, and in a highly-available setup by the Standby NameNode, which also tails a shared journal so it can take over. (The HA machinery itself — journal quorum, ZooKeeper-based failover — is a cluster-operations concern rather than part of the storage model.) ## The one thing the NameNode does *not* persist **Block locations.** The mapping from block ID to "which machines have it" is never written to `fsimage`. It is reconstructed from **block reports**: each DataNode scans its local block directories and tells the NameNode the full list of block replicas it holds, at startup and periodically thereafter (every six hours by default). This is deliberate — DataNodes are the ground truth for what is actually on disk, and a persisted map would immediately be stale after a disk failure. A consequence: after a cold NameNode start the cluster sits in **safe mode**, read-only, until enough blocks have been reported at minimum replication. On a large cluster that wait is minutes, and it is dominated by the number of blocks, not the number of bytes. ## What a DataNode does A DataNode stores each block replica as a plain file (plus a checksum metadata file) in a configured set of local directories, typically one per physical disk — HDFS deliberately uses JBOD rather than RAID, because replication already provides redundancy and RAID would cap a disk group at its slowest spindle. A DataNode's duties: - **heartbeat** to the NameNode (every three seconds by default) so the NameNode knows it is alive; sustained silence marks it dead and triggers re-replication of everything it held - **block reports** describing its replicas - **serve and receive block data** directly from and to clients, and forward blocks along a write pipeline to the next DataNode - **verify checksums** on read and via a background block scanner, reporting corrupt replicas so the NameNode can replace them DataNodes hold no knowledge of files or directories. A DataNode does not know that block `blk_1073742001` is the third block of `/data/events.log`; only the NameNode knows that. ## How a read flows The client calls `open()`, the NameNode returns the block list along with the DataNode locations for each block, sorted by network distance from the client. The client connects to the nearest replica and streams bytes; if that DataNode fails or a checksum mismatches, the client transparently retries the next replica and reports the bad one. Data locality — scheduling computation on a node that already holds the block — falls straight out of this design. ## Traps to avoid The NameNode is not a proxy, a cache, or a metadata *database on disk* consulted per request — it is an in-memory index. The Secondary NameNode is not a hot standby; it only checkpoints. And DataNodes are not aware of replication factor: the NameNode decides when a block is under-replicated and instructs a DataNode to copy it elsewhere.

  • Why does the NameNode rebuild block locations from DataNode reports instead of persisting them in fsimage?
    Because DataNodes are the only reliable source of truth about what is physically on disk. Disks fail, machines are re-imaged, and replicas are moved by the balancer, so a persisted location map would be wrong the moment it was written. Rebuilding from block reports at startup guarantees the NameNode's view matches reality — at the cost of a safe-mode wait while reports arrive.
  • Is the Secondary NameNode a hot standby that takes over when the NameNode dies?
    No. Despite the name, it performs checkpointing only: it periodically fetches the fsimage and edit log, merges them into a fresh fsimage, and ships it back, so restart time stays bounded. It cannot serve clients. Real failover requires a proper HA deployment with an active and a standby NameNode over a shared journal.
  • What happens to the blocks a DataNode held when it stops sending heartbeats?
    After the NameNode declares it dead, every block that node held drops below its target replication factor. The NameNode queues those blocks for re-replication and instructs surviving DataNodes to copy them to new nodes, throttled so recovery does not saturate the cluster. Nothing is lost as long as another replica survives.

The NameNode is the library's card catalogue: it knows every title and which shelf each volume sits on, but holds no books. The DataNodes are the shelves — and readers walk to the shelf themselves rather than asking the catalogue to fetch the book.

saying these in an interview costs you the question

  • Claiming client file data is streamed through the NameNode
  • Calling the Secondary NameNode a hot standby or failover node
  • Saying DataNodes know which file a block belongs to
  • Believing block locations are stored in fsimage on disk
  • Describing the NameNode as a disk-based metadata database

context