The NameNode is the metadata server of HDFS. It knows every directory, file, permission and block ID in the filesystem, and it knows where each block's replicas live right now. It never handles file data. Because every open, list, create and rename goes through it, the NameNode's memory, lock and RPC queue set the ceiling on how many files a cluster can hold and how many operations per second it can serve.

This article opens up the process: what it keeps in memory and what it persists, how a write moves through it, how the edit log reaches disk without serializing every client, what DataNode reports do, how to size the heap, and how to read its metrics when it is slow. High availability and checkpointing have their own articles, HDFS high availability and checkpoints and JournalNodes; here they appear only where they touch the core process.

Inside the NameNode processClientsClientProtocol RPCDataNodesheartbeats, reportsRPC servercall queue, handlersservice portFSDirectory: namespace treeinodes, permissions, quotasBlockManager: blocks mapblock to locations, from reportsboth guarded by the FSNamesystem lockEdit logappend, then logSync outside the lockLocal edits dirsdfs.namenode.edits.dirJournalNodesquorum, HAReplication monitorunder/over-replicatedFsImagecheckpoint on diskPersisted:namespace + block IDsNOT persisted:block locations
RPC handlers act on the namespace tree and the block map under one lock; mutations are logged and synced outside it. Block locations are rebuilt from DataNode reports, never persisted.

What the NameNode owns

The NameNode's state splits into two kinds, and the split explains most of its behaviour. The namespace is the directory tree: inodes for directories and files, names, owners, permissions, ACLs, quotas, modification times, and for each file the ordered list of block IDs and their generation stamps. It is persisted, as an FsImage checkpoint plus an edit log of every change since.

The block map records, for each block, which DataNodes hold a replica. It is never written to disk. After a restart the NameNode knows that /logs/day1 consists of blocks 1073741825 and 1073741826, but not where they are, until DataNodes send block reports. That is why a restarting NameNode sits in safe mode: it refuses writes until a configured fraction of blocks (dfs.namenode.safemode.threshold-pct, default 0.999) has at least the minimum number of reported replicas.

Blocks are not inodes. A file inode points to block objects; the blocks map holds those objects and their replica locations. Both live on the JVM heap, which is why the commonly used estimate is about 150 bytes per namespace object: per file, per directory and per block.

What the NameNode owns

The NameNode's state splits into two kinds, and the split explains most of its behaviour. The namespace is the directory tree: inodes for directories and files, names, owners, permissions, ACLs, quotas, modification times, and for each file the ordered list of block IDs and their generation stamps. It is persisted, as an FsImage checkpoint plus an edit log of every change since.

The block map records, for each block, which DataNodes hold a replica. It is never written to disk. After a restart the NameNode knows that /logs/day1 consists of blocks 1073741825 and 1073741826, but not where they are, until DataNodes send block reports. That is why a restarting NameNode sits in safe mode: it refuses writes until a configured fraction of blocks (dfs.namenode.safemode.threshold-pct, default 0.999) has at least the minimum number of reported replicas.

Blocks are not inodes. A file inode points to block objects; the blocks map holds those objects and their replica locations. Both live on the JVM heap, which is why the commonly used estimate is about 150 bytes per namespace object: per file, per directory and per block.

Inside the process: lock, block manager and RPC server

Inside the process, FSNamesystem is the coordinator. It owns FSDirectory (the tree) and the BlockManager (the blocks map, replication queues and placement policy). A single read-write lock guards the namespace: reads such as getFileInfo, listStatus and getBlockLocations share it; mutations such as create, rename and delete take it exclusively. One long write-locked operation, for example deleting a directory with ten million files, stalls every other caller. Deletes limit the damage by unlinking the subtree in one write-lock hold and then removing its blocks from the block map in small batches, releasing the lock between them.

In front sits the Hadoop IPC server. A listener accepts connections, readers decode calls into a call queue, and dfs.namenode.handler.count handler threads (default 10) take calls off the queue and execute them. DataNode heartbeats and block reports can be moved to a separate service RPC port with its own handlers by setting dfs.namenode.servicerpc-address, so a storm of client calls cannot starve the reports that keep the block map accurate.

Clients retry RPCs after timeouts, and a retried create or rename must not apply twice. The NameNode keeps a retry cache keyed by client ID and call ID, so a duplicate of a mutation that already succeeded receives the original answer rather than an error. The cache survives failover in an HA pair, which is what makes retries across a failover safe.

A write as the NameNode sees it

A file write is several RPCs, not one, and each touches the NameNode differently. create checks permissions and quotas, adds a file inode under construction, grants the client a lease on the path, logs an OP_ADD transaction and returns. No DataNodes are chosen yet. For each block, the client calls addBlock: the NameNode allocates a block ID and generation stamp, asks the placement policy for targets, logs the new block and returns the pipeline. The client streams data directly to those DataNodes. Finally complete succeeds once the last block has the minimum replicas, and the lease is released.

Leases are how the NameNode enforces a single writer without trusting clients to behave. The writing client renews its leases periodically. If it stops renewing past the soft limit of about a minute, another client may take the file over; past the hard limit the NameNode itself starts lease recovery: it picks a primary DataNode, has the replicas of the last block agree on a common length and a new generation stamp, and closes the file. This is why a crashed writer's file is briefly unreadable at its tail and then appears with a shorter length than the writer believed it had written.

# NameNode side of addBlock(), simplified
def add_block(path, client, previous):
    with fsn.read_lock():                      # phase 1: validate
        inode = fsdir.get_under_construction(path, lease_holder=client)
        replication = inode.replication
    targets = placement.choose(replication=replication,     # phase 2: no namesystem lock
                               writer=client.host, excluded=client.bad_nodes)
    with fsn.write_lock():                     # phase 3: re-validate and commit
        inode = fsdir.get_under_construction(path, lease_holder=client)
        commit_or_complete(previous)           # previous block now has its final length
        blk = Block(next_block_id(), gen_stamp=next_gen_stamp())
        inode.append(blk); block_manager.add_under_construction(blk, targets)
        txid = edit_log.log_add_block(path, blk)   # buffered, not yet durable
    edit_log.log_sync(txid)                     # outside the lock: group commit
    return LocatedBlock(blk, targets)

Validating under the read lock, choosing targets with no namesystem lock held, and only then committing under the write lock keeps the comparatively expensive placement work from blocking anyone; the re-check in the last phase catches a file deleted or a lease lost in between. Default placement puts the first replica on the writer's node when the writer is a DataNode, the second on a different rack and the third on the second's rack, as described in blocks and replication.

The edit log and group commit

Every mutation becomes a numbered transaction in the edit log, and the client is told success only after the transaction is durable. Syncing to disk, and with HA to a quorum of JournalNodes, takes milliseconds; doing that while holding the write lock would cap the cluster at a few hundred mutations per second. The NameNode instead appends the record to an in-memory buffer under the lock, releases the lock, and calls logSync. One thread flushes the buffer; any other thread whose transaction ID is already covered returns immediately. Hundreds of transactions share one fsync. This group commit is why mutation throughput depends far more on journal latency than on CPU.

Edits go to every directory in dfs.namenode.edits.dir and, in an HA pair, to a qjournal:// URI. A write to the journal quorum must succeed on a majority; if it cannot, the NameNode aborts rather than continue with an edit that might be lost. That abort looks alarming but is the correct choice.

The edit log cannot grow forever, because replaying it is what makes restarts slow. A checkpoint merges the FsImage with recent edits into a new FsImage, triggered every dfs.namenode.checkpoint.period (default 3600 seconds) or after dfs.namenode.checkpoint.txns (default 1,000,000) transactions, whichever comes first. In HA the standby performs it and uploads the image to the active.

Heartbeats, block reports and re-replication

DataNodes keep the block map true. Every dfs.heartbeat.interval (3 seconds) each DataNode reports capacity, usage and active transfers, and the reply carries commands: replicate this block, delete that one. Incremental block reports announce replicas as they are received or deleted, so a new block appears in the map within seconds. A full block report of every replica on the node is sent at startup and then every six hours by default (dfs.blockreport.intervalMsec).

A DataNode that stops heartbeating is marked dead after about 10.5 minutes with defaults: twice the five-minute recheck interval plus ten heartbeats. Only then does the replication monitor queue its blocks for re-replication, prioritising blocks with one live replica left. The delay is deliberate: a rolling restart of a node should not trigger a terabyte of copying. Full reports from large nodes are costly, since processing them takes the write lock, so very dense DataNodes split reports per storage directory.

More on the DataNode side is in the DataNode deep dive.

Sizing the heap: a worked example

Heap size is the NameNode's capacity limit. Work an example. A cluster holds 100 million files averaging 1.5 blocks each, so 150 million blocks, plus 10 million directories: 260 million objects. At 150 bytes each that is about 39 GB of live metadata. Real heaps also carry replica location entries, the RPC and edit buffers, lease and snapshot state, and the free headroom a garbage collector needs, which is why operators plan conservatively with the widely quoted rule of roughly 1 GB of heap per million blocks. For this cluster that points to a heap in the low hundreds of gigabytes.

Practical settings: set -Xms equal to -Xmx in HDFS_NAMENODE_OPTS so the heap never resizes, use G1 or another low-pause collector on large heaps, leave memory for the OS, and size the standby identically, because it holds the same namespace. Watch the Detected pause in JVM or host machine warnings from the JVM pause monitor; a multi-second pause means clients time out and, in HA, the failover controller may decide the active is dead. The cheapest heap savings come from fewer objects: the small files problem is a heap problem first.

Observing a slow NameNode

The NameNode exposes everything through JMX at http://<namenode>:9870/jmx and a few commands answer most questions:

# Namespace size and block health
curl -s 'http://nn1:9870/jmx?qry=Hadoop:service=NameNode,name=FSNamesystem' \
  | jq '.beans[0] | {FilesTotal, BlocksTotal, MissingBlocks, UnderReplicatedBlocks}'

# RPC queueing vs processing time on the client port
curl -s 'http://nn1:9870/jmx?qry=Hadoop:service=NameNode,name=RpcActivityForPort8020' \
  | jq '.beans[0] | {CallQueueLength, RpcQueueTimeAvgTime, RpcProcessingTimeAvgTime}'

hdfs dfsadmin -report            # live/dead DataNodes, capacity
hdfs dfsadmin -safemode get
hdfs haadmin -getServiceState nn1
hdfs dfsadmin -fetchImage /tmp/img && hdfs oiv -p Delimited -i /tmp/img/fsimage_* -o ns.tsv

Read the RPC numbers together. High queue time with low processing time means too few handlers or one abusive client; high processing time means lock contention or a slow journal. The offline image viewer turns a fetched FsImage into a table you can query for file counts and sizes per directory without touching the live NameNode, which is the safe way to find who created ten million tiny files.

Failure modes

FailureWhat you seeResponse
Long GC pauseClient timeouts, JVM pause warnings, possible failoverBigger heap or fewer objects; tune GC; check host swap
Handler saturationRising CallQueueLength and queue timeRaise handler count; separate service port; FairCallQueue
One heavy userEveryone slow, one user dominates audit logFairCallQueue with backoff; fix the job's listing pattern
Slow journalHigh processing time on mutations onlyCheck JournalNode disks and network; dedicated disks
Edits directory full or failedNameNode abortsMultiple edits dirs, disk alerts, retain policy
Missing blocks after restartSafe mode does not exitFind dead DataNodes first; never force-leave blindly

The FairCallQueue (ipc.8020.callqueue.impl=org.apache.hadoop.ipc.FairCallQueue) replaces the single FIFO with priority levels scored by each user's recent call volume, so one runaway job is slowed instead of the whole cluster. It is the first thing to enable on a shared cluster.

Trade-offs and scaling out

Keeping the whole namespace in one JVM's memory is the source of both the NameNode's speed and its limit. Metadata operations are memory lookups under a lock, so they take microseconds, but the namespace cannot exceed one heap and mutations serialise on one lock. There are three ways out, each with a cost. Federation splits the namespace across independent NameNodes sharing DataNodes, with router-based federation presenting one mount table to clients; you gain capacity and throughput and accept cross-namespace renames being impossible. Observer NameNodes serve consistent reads from a standby that tails the journal, offloading read-heavy workloads at the price of another node and a client configuration change. Fewer, bigger files through compaction or container formats cost pipeline work and help everything else. Many teams also move cold data to object storage, which has no NameNode at all and weaker rename semantics.

What to do next

  1. Pull FilesTotal and BlocksTotal from JMX and compute objects per GB of heap today.
  2. Graph RpcQueueTimeAvgTime against RpcProcessingTimeAvgTime for a week and note peaks.
  3. Configure a service RPC port and confirm DataNodes use it.
  4. Enable FairCallQueue on the client port of a shared cluster.
  5. Alert on JVM pause warnings over one second and on MissingBlocks above zero.
  6. Fetch an FsImage, run hdfs oiv, and rank directories by file count to target small-file cleanup.
Key takeaway: The NameNode is an in-memory database with a write-ahead log: the namespace is persisted, block locations are rebuilt from reports, mutations group-commit outside one global lock, and heap size bounds the object count. Measure objects, RPC queueing and GC pauses, and scale by reducing files or splitting the namespace.