How does HDFS NameNode high availability keep a standby ready and stop two NameNodes from both writing?
answer
- two brains, one pen
- the log lives on a quorum
- a majority must acknowledge each edit
- ZooKeeper picks, ZKFC acts
- epoch numbers stop the stale writer
basics
~20 sThe active NameNode writes edits to a quorum of JournalNodes; the standby tails them and receives DataNode block reports, so it is hot. ZooKeeper plus a ZKFC per NameNode elects the active, and journal epochs plus fencing block the stale writer.
solid answer
~40 sTwo NameNodes share one edit log stored on an odd set of `JournalNode`s, addressed as `dfs.namenode.shared.edits.dir=qjournal://jn1:8485;jn2:8485;jn3:8485/<nameservice>`. An edit is committed once a **majority** of JournalNodes acknowledge it; if the active cannot reach a majority it aborts rather than diverge. The standby continuously tails those edits, and DataNodes heartbeat and block-report to *both* NameNodes, so the standby's block map is current and failover takes seconds rather than minutes of block reporting. With `dfs.ha.automatic-failover.enabled=true`, a `ZKFC` process beside each NameNode health-checks it and contests a ZooKeeper election; the winner is transitioned to active. Split brain is prevented by the journal itself — it accepts writes only from the highest epoch — with `dfs.ha.fencing.methods` (`sshfence`, or `shell(...)`) additionally stopping the old process from serving stale reads. Clients address `hdfs://<nameservice>/` and fail over transparently via `ConfiguredFailoverProxyProvider`.
code
xml · 24 lines<property>
<name>dfs.nameservices</name>
<value>prod</value>
</property>
<property>
<name>dfs.ha.namenodes.prod</name>
<value>nn1,nn2</value>
</property>
<property>
<name>dfs.namenode.shared.edits.dir</name>
<value>qjournal://jn1:8485;jn2:8485;jn3:8485/prod</value>
</property>
<property>
<name>dfs.client.failover.proxy.provider.prod</name>
<value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value>
</property>
<property>
<name>dfs.ha.fencing.methods</name>
<value>sshfence</value>
</property>
<property>
<name>dfs.ha.automatic-failover.enabled</name>
<value>true</value>
</property>go deeper
Know that production HDFS runs two NameNodes, active and standby, with a shared edit log, and that the Secondary NameNode is a different thing entirely and is not a failover target.
Explain the mechanics: journal quorum with majority acknowledgement, the standby tailing edits, DataNodes reporting to both NameNodes, and ZooKeeper plus ZKFC electing the active.
Demonstrate that you have run this — fencing configuration and its failure modes, what happens on journal quorum loss, graceful failover for maintenance, and the alerts you set on lag and ZKFC health.
Own the availability target end to end: how many standbys, JournalNode and ZooKeeper placement across failure domains, whether Observer reads are worth the operational complexity, and what HDFS unavailability actually costs the platform above it.
## What HA has to solve Without HA, the NameNode is a single point of failure for the whole platform: it holds the only copy of the live namespace in memory, and if it dies every HDFS read and write, and therefore every job, Hive query and HBase region, stops until an operator restarts it and it finishes loading `fsimage`, replaying edits and collecting block reports — which on a large cluster is tens of minutes. HA replaces that with a second NameNode that can take over in seconds. ## The shared edit log: Quorum Journal Manager Every namespace mutation (create, rename, delete, block allocation) is appended to an edit log before it is acknowledged. In an HA cluster that log does not live on one disk; it lives on a set of `JournalNode` daemons — an odd number, in practice three — configured as: ``` dfs.namenode.shared.edits.dir = qjournal://jn1:8485;jn2:8485;jn3:8485/prod ``` The active NameNode writes each edit batch to all of them and considers it durable once a **majority** acknowledge. Three JournalNodes therefore tolerate one failure, five tolerate two. The critical rule is what happens when the majority is *not* reachable: the active NameNode aborts itself rather than accept writes it cannot durably record. Losing JournalNode quorum is an outage, not a silent degradation — which is why JournalNodes go on three separate machines with their own local disk for `dfs.journalnode.edits.dir`. ## Why the standby is hot, not cold The standby does two things continuously. It **tails the journal**, applying every edit to its own in-memory namespace, so its view of the directory tree is seconds behind the active. And it receives **DataNode heartbeats and block reports**: DataNodes are configured with the nameservice, not a single NameNode address, and report to every NameNode in it. That second point is what makes fast failover possible at all, because the block map — which node holds which replica — is never written to the edit log. A NameNode that has not been receiving block reports cannot serve reads until the cluster has reported in, which is the multi-minute startup cost HA is designed to avoid. ## Automatic failover: ZooKeeper and the ZKFC Set `dfs.ha.automatic-failover.enabled=true` and point `ha.zookeeper.quorum` (in `core-site.xml`) at a ZooKeeper ensemble. Each NameNode host then runs a `DFSZKFailoverController` (started as `hdfs --daemon start zkfc`, initialised once with `hdfs zkfc -formatZK`). The ZKFC health-checks its local NameNode over RPC and, if healthy, tries to hold an ephemeral znode representing "active". If the NameNode becomes unhealthy or the host dies, the session expires, the znode disappears, the peer ZKFC wins the election and transitions its NameNode to active. ## Fencing and split brain The scary case is an active NameNode that is *not* dead — just GC-paused, or partitioned — while a second one is promoted. With the Quorum Journal Manager, the journal itself resolves this: each writer takes an increasing **epoch** number when it becomes active, and JournalNodes reject writes from any lower epoch. A stale NameNode therefore cannot corrupt the namespace; its next write fails and it shuts down. Fencing methods configured in `dfs.ha.fencing.methods` — `sshfence` (SSH in and kill the process) or `shell(...)` — add belt and braces, mainly to stop the old process from answering stale *reads* to clients that have not failed over yet. Many QJM deployments configure `shell(/bin/true)` and rely on epoch fencing; if you use `sshfence` you must ensure the key is present and reachable, or a failover can hang waiting to fence a dead host. ## Clients Clients no longer address a host. They use `hdfs://prod/path`, and `dfs.client.failover.proxy.provider.prod` set to `ConfiguredFailoverProxyProvider` makes the client try each NameNode in `dfs.ha.namenodes.prod` and retry on a `StandbyException`. A running job survives a failover as a retry, not a failure. ## Beyond two NameNodes Hadoop 3 supports more than one standby, so you can tolerate two NameNode failures. Later Hadoop 3 releases also add the **Observer NameNode**, a standby that can serve consistent reads (clients issue an `msync` to bound staleness), which offloads read RPC from the active on very large clusters. Both are worth naming; neither is common on a small cluster. ## Operating it `hdfs haadmin -getServiceState nn1` reports active/standby. `hdfs haadmin -failover nn1 nn2` performs a graceful failover for planned maintenance. `-transitionToActive --forcemanual` exists but is dangerous and only valid when automatic failover is disabled. The alerts that matter: JournalNode quorum health, edit-log lag on the standby, ZKFC liveness, and the state where *neither* NameNode is active — that is a ZooKeeper or ZKFC problem, not a NameNode one.
- What happens to the active NameNode if a majority of JournalNodes becomes unreachable?It aborts. An edit is only committed once a majority of JournalNodes acknowledge it, so a NameNode that cannot reach quorum shuts itself down rather than accept writes it cannot durably log. That is deliberate: it converts a potential namespace divergence into a clean outage. It also means JournalNode quorum loss takes HDFS down even though both NameNodes are alive.
- If QJM epochs already prevent two writers, why configure fencing at all?Epoch fencing protects the *edit log*, not the clients. A stale active that is merely partitioned can still answer read RPCs from its now-outdated in-memory namespace until it tries to write and fails. A fencing method kills that process outright. It also covers deployments on shared storage rather than QJM, where the journal offers no epoch protection.
- Why does an HA cluster fail over in seconds when a cold NameNode start takes many minutes?Because the standby is already warm on both axes. It has tailed every edit, so its namespace is current, and DataNodes block-report to both NameNodes, so its block map is populated. A cold start must load fsimage, replay edits, and then wait in safe mode for enough block reports to arrive before it can serve reads.
- How do you fail over deliberately for planned NameNode maintenance?Use `hdfs haadmin -failover nn1 nn2`, which coordinates a graceful transition: it fences and demotes the current active, then promotes the target once it has caught up on the journal. Verify with `hdfs haadmin -getServiceState` on both before and after. Avoid `-transitionToActive --forcemanual`, which bypasses the coordination and is only appropriate when automatic failover is off.
Two pilots share one flight recorder that only accepts the pilot holding the current stamped authority code; the moment a new code is issued, the recorder ignores the old pilot no matter what they do.
saying these in an interview costs you the question
- Says the Secondary NameNode takes over when the NameNode dies
- Thinks both NameNodes can be active to share load
- Proposes two JournalNodes, which cannot form a majority
- Believes failover is fast because the standby reads fsimage from disk
- Thinks ZooKeeper stores the HDFS namespace or the block map