skip to content

Hadoop

The original big-data stack: HDFS for storage, YARN for resources, MapReduce for compute. Even where Spark has replaced the compute layer, interviewers ask about Hadoop because HDFS and YARN still sit under many clusters and explain why later engines look the way they do.

on this pageshow

explore

questions

30

In a Hadoop 3 cluster, which daemons run on the master nodes and which run on every worker node?

level: juniorimportance: must knowfreq 52%

answer

  1. small brain tier, wide muscle tier
  2. two systems, same machines
  3. one stores, one schedules
  4. the pair that sits on every worker
  5. NameNode and ResourceManager coordinate, they do not store

basics

~20 s

Master nodes run the coordinators: the HDFS NameNode (plus JournalNodes and ZKFC when HA is on) and the YARN ResourceManager. Every worker runs a DataNode for storage and a NodeManager for compute, co-located so containers read local blocks.

solid answer

~40 s

Hadoop splits coordination from capacity. The **master** tier runs the HDFS `NameNode`, which holds the namespace and block map in memory, and the YARN `ResourceManager`, which hands out containers. In an HA cluster the master tier also runs a second NameNode, three `JournalNode`s, one `DFSZKFailoverController` (ZKFC) next to each NameNode, and an external ZooKeeper ensemble. Every **worker** node runs a `DataNode` — it stores blocks on the disks listed in `dfs.datanode.data.dir` and heartbeats and block-reports to the NameNode — and a `NodeManager`, which advertises the node's resources and launches containers. DataNode and NodeManager are deliberately co-located so a task can usually read its input off local disk. Auxiliary services (MapReduce JobHistoryServer, HttpFS, Timeline Server) live on master or edge nodes. `jps` on a host tells you which tier it is.

code

text · 9 lines
text
# master node
14231 NameNode
14550 DFSZKFailoverController
14802 JournalNode
15104 ResourceManager

# worker node
9012 DataNode
9310 NodeManager

go deeper

for a junior

Be able to name the four core daemons and where each runs: NameNode and ResourceManager on masters, DataNode and NodeManager on every worker. Knowing that jps shows them is worth a sentence.

for a middle

Explain what each daemon actually holds — the NameNode's in-memory namespace and block map, the NodeManager's advertised memory and vcores — and why DataNode and NodeManager are co-located for locality.

for a senior

Show you have operated this: which hosts carry JournalNodes and ZKFC, why edge nodes exist, what a missing DataNode on a NodeManager host does to job performance, and how you verify daemon health from the UIs.

for a principal

Own the placement policy — master-node count and isolation, whether to keep storage and compute co-located at all versus a disaggregated object-store design, and what that choice costs in network capacity and locality.

## The shape of a Hadoop cluster A Hadoop cluster is two tiers of long-running Java daemons. The **master** tier is small (typically three to five machines) and holds all the coordination state; the **worker** tier is everything else and holds the data and does the work. Both HDFS and YARN follow the same master/worker split, and the two systems are deployed onto the same machines so that computation can be scheduled next to its data. ## Master-side daemons **NameNode (HDFS).** Keeps the entire filesystem namespace — the directory tree, file names, permissions, and the list of blocks each file is made of — in memory, and persists changes to an on-disk `fsimage` plus an edit log under `dfs.namenode.name.dir`. It also holds the *block map*: which DataNode currently has which block replica. That map is never persisted; it is rebuilt from DataNode block reports at startup. **ResourceManager (YARN).** The cluster-wide scheduler. It tracks the resources every NodeManager advertises, accepts application submissions, and allocates containers to applications. It does not run user code itself. **JournalNode (HDFS HA only).** A small daemon holding a copy of the shared edit log. You run an odd number, usually three, spread across separate machines. **ZKFC / DFSZKFailoverController (HDFS HA only).** One process per NameNode host. It health-checks its local NameNode and participates in a ZooKeeper election to decide which NameNode is active. **ZooKeeper.** Not a Hadoop daemon, but an external ensemble (3 or 5 nodes) that both HDFS automatic failover and ResourceManager HA depend on. **Secondary NameNode.** Present only in a *non*-HA cluster, where it periodically merges the edit log into a new `fsimage`. In an HA cluster it does not exist — the standby NameNode does that work. ## Worker-side daemons **DataNode.** Stores blocks as ordinary files on each of the local disks listed in `dfs.datanode.data.dir` (JBOD, one directory per disk — you do not RAID DataNode disks). It sends a heartbeat every few seconds (`dfs.heartbeat.interval`, 3 seconds by default) and a full block report periodically, and it serves and receives block data on its own data-transfer port. **NodeManager.** Advertises the node's usable memory and vcores (`yarn.nodemanager.resource.memory-mb`, `yarn.nodemanager.resource.cpu-vcores`), launches and monitors containers on request from the ResourceManager, kills containers that exceed their limits, and runs auxiliary services such as `mapreduce_shuffle`. Every worker runs exactly one of each. One of the containers a NodeManager launches will be an **ApplicationMaster** — a per-application process, not a cluster daemon, that asks the ResourceManager for the rest of its containers. ## Auxiliary and edge services The MapReduce **JobHistoryServer**, the YARN **Timeline Server**, **HttpFS**, and the HDFS **Router** (for federation) are optional daemons usually placed on master or edge hosts. *Edge* or *gateway* nodes run no cluster daemon at all: they just carry `$HADOOP_CONF_DIR` and the client jars so users can submit jobs. ## How the daemons find each other Nothing is auto-discovered. Workers find the masters through the site XML files pushed to every node — `fs.defaultFS` in `core-site.xml` points DataNodes and clients at HDFS, `yarn.resourcemanager.hostname` (or the HA `rm-ids` set) in `yarn-site.xml` points NodeManagers at the ResourceManager. The `etc/hadoop/workers` file (called `slaves` before Hadoop 3) is only used by the `start-dfs.sh`/`start-yarn.sh` helper scripts to SSH out and start daemons; a config-managed cluster often does not use it at all. ## Checking a node `jps` lists the JVMs on a host, and their class names are the daemon names above. The web UIs are the other quick check — and their default ports moved in Hadoop 3: the NameNode UI is on 9870 (it was 50070 in Hadoop 2), the DataNode UI on 9864, the Secondary NameNode on 9868; the ResourceManager UI stayed on 8088 and the NodeManager on 8042. A worker that shows a NodeManager but no DataNode is the classic cause of "all my tasks read over the network": YARN will happily schedule there, but there is no local data to read.

  • Why are the DataNode and NodeManager deliberately placed on the same machine?
    For data locality. The ResourceManager is told where each block lives, so it can place a container on a node that already holds the input block, and the task reads from local disk instead of pulling the block across the rack uplink. Separating storage and compute nodes is a legitimate design today, but it trades locality for elasticity and leans much harder on the network.
  • Is the ApplicationMaster a cluster daemon like the NodeManager?
    No. The ApplicationMaster is created per application, inside a YARN container on some worker node, and it dies when the application finishes. It negotiates the application's remaining containers with the ResourceManager and monitors its own tasks. The NodeManager, by contrast, is a permanent daemon on every worker regardless of what is running.
  • What does an edge or gateway node run?
    No cluster daemons — only the Hadoop client jars and a copy of `$HADOOP_CONF_DIR`. Users log in there to run `hdfs dfs`, submit jobs, or run Hive/Spark clients. Keeping submission off the master nodes stops user JVMs from competing with the NameNode and ResourceManager for memory.

saying these in an interview costs you the question

  • Says the NameNode stores the actual file data
  • Thinks the Secondary NameNode is a hot standby that takes over
  • Claims DataNodes discover the NameNode automatically on the network
  • Says the ApplicationMaster is a permanent daemon on each worker
  • Puts a NameNode on every node so metadata is replicated

context

open as a page

In Hive, what happens to the data when you DROP a managed table versus an external table?

level: juniorimportance: must knowfreq 50%

basics

~20 s

Dropping a managed table removes its metastore entry and deletes its data files. Dropping an external table removes only the metastore entry and leaves the files untouched. Declare EXTERNAL whenever another system owns the data.

open as a page

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

level: juniorimportance: must knowfreq 60%

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.

open as a page

What phases does a Hadoop MapReduce job move through from input split to final output?

level: juniorimportance: must knowfreq 55%

basics

~20 s

A Hadoop MapReduce job runs map, then shuffle and sort, then reduce. One map task handles each InputSplit; its output is partitioned and sorted on local disk, fetched by reducers over the network, merged by key, and written out by the OutputFormat.

open as a page

In YARN, what do the ResourceManager, NodeManager and ApplicationMaster each do?

level: juniorimportance: must knowfreq 58%

basics

~20 s

YARN splits cluster management three ways: the ResourceManager schedules containers across the whole cluster, a NodeManager runs and monitors containers on each machine, and every application gets its own ApplicationMaster that requests containers and drives that job.

open as a page

What does the Hive Metastore provide that query engines outside Hive still need?

level: middleimportance: must knowfreq 55%

basics

~20 s

The Hive Metastore is a relational catalog mapping table and partition names to storage paths, columns, types, file formats and statistics. Engines such as Trino, Presto and Flink read it so that directories of files behave as SQL tables.

open as a page

Why can adding a combiner to a Hadoop MapReduce job make an average come out wrong?

level: middleimportance: must knowfreq 48%

basics

~20 s

A combiner is a reduce-style function the framework may run zero, one or many times over a mapper's own output. An arithmetic mean is not decomposable that way, so averaging partial groups and then averaging those averages gives a wrong result.

open as a page

How does HDFS NameNode high availability keep a standby ready and stop two NameNodes from both writing?

level: seniorimportance: must knowfreq 48%

basics

~20 s

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

open as a page

Why does the HDFS NameNode struggle with ten million tiny files?

level: seniorimportance: must knowfreq 50%

basics

~20 s

Every file, directory and block is an object held in the NameNode's heap, and the cost is per object, not per byte. Ten million tiny files consume as much master memory as ten million huge ones while storing almost no data.

open as a page

Why does a YARN NodeManager kill a container for exceeding physical memory limits?

level: seniorimportance: must knowfreq 50%

basics

~20 s

Each YARN container is granted a fixed memory amount, and the NodeManager monitors the container's whole process tree against it. When total resident memory crosses the grant the container is killed, because YARN protects the node from oversubscription rather than letting the machine swap or OOM.

open as a page

Why does an HA HDFS cluster have no Secondary NameNode, and what performs checkpointing instead?

level: middleimportance: should knowfreq 44%

basics

~20 s

The Secondary NameNode only merges the edit log into a new fsimage so restarts stay fast — it is not a standby. In an HA cluster the standby NameNode already tails every edit, so it performs the checkpoint and the Secondary NameNode role disappears.

open as a page

In Hadoop, what do core-site.xml, hdfs-site.xml and yarn-site.xml each configure, and what overrides them?

level: middleimportance: should knowfreq 46%

basics

~20 s

core-site.xml holds cluster-wide settings like fs.defaultFS, hdfs-site.xml configures NameNode and DataNode behaviour, yarn-site.xml the ResourceManager and NodeManagers. Each overrides the matching *-default.xml shipped in the jars, and a job can override again unless the property is marked final.

open as a page

When would you store a dataset in HBase rather than as a Hive table on HDFS?

level: middleimportance: should knowfreq 45%

basics

~20 s

Choose HBase when you need millisecond random reads and writes of individual rows by key, and updates in place. Choose a Hive table on HDFS when the workload is large sequential scans and aggregations over immutable files.

open as a page

Why is the default HDFS block size 128 MB rather than a few kilobytes?

level: middleimportance: should knowfreq 52%

basics

~20 s

A large HDFS block amortises disk seek time over a long sequential read and keeps the NameNode's per-byte metadata cost tiny. Blocks are logical, so a 5 MB file consumes 5 MB of disk, not 128 MB.

open as a page

Where does HDFS put the three replicas of a block under its default rack-aware policy?

level: middleimportance: should knowfreq 38%

basics

~20 s

With the default policy, the first replica goes on the writing client's own node (or a random DataNode if the client is off-cluster), the second on a node in a different rack, and the third on another node in that same second rack.

open as a page

How does an HDFS client write one block through a DataNode replication pipeline?

level: middleimportance: should knowfreq 45%

basics

~20 s

The client asks the NameNode for a block and a list of DataNodes, then streams packets to the first DataNode, which forwards to the second, which forwards to the third. Acknowledgements travel back up the chain before packets are released.

open as a page

In Hadoop MapReduce, what is an InputSplit and what determines its size?

level: middleimportance: should knowfreq 42%

basics

~20 s

An InputSplit is the logical byte range one map task will process, plus host hints for locality — not a physical copy of data. With FileInputFormat its size is max(minSize, min(maxSize, blockSize)), so by default it equals the HDFS block size of 128 MB.

open as a page

How do YARN's Capacity Scheduler and Fair Scheduler differ when sharing one cluster?

level: middleimportance: should knowfreq 45%

basics

~20 s

Both divide a YARN cluster into hierarchical queues, but the Capacity Scheduler starts from guaranteed percentages per queue with an elastic ceiling, while the Fair Scheduler starts from equal sharing among running applications and pulls resources back toward each queue's fair share over time.

open as a page

A Hive dynamic-partition INSERT produced thousands of tiny files — what caused it and what do you tune?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Each writing task creates its own file in every partition it touches, so files multiply as tasks times partitions. Funnel each partition to one writer with DISTRIBUTE BY the partition columns, enable Hive's merge settings, and use coarser partition granularity.

open as a page

One reducer in a Hadoop MapReduce job runs for hours after every other reducer finishes — why?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Almost always key skew. The default HashPartitioner sends every value for a given key to one reduce task, and a single key's group cannot be split across reducers, so one hot key pins one task. Speculative execution does not help, because the retry processes the same data.

open as a page

A YARN application sits in ACCEPTED state for 20 minutes and never runs — how do you diagnose it?

level: seniorimportance: should knowfreq 45%

basics

~20 s

ACCEPTED means YARN took the application but has not launched its ApplicationMaster container yet. The cause is always missing capacity: the queue is full, its ApplicationMaster budget is exhausted, a per-user limit binds, or no healthy node can host the requested container size.

open as a page

How would you size a new Hadoop cluster's DataNode storage, NameNode heap and worker node shape?

level: principalimportance: should knowfreq 34%

basics

~20 s

Size three budgets separately: raw storage as logical data times the replication factor plus headroom for non-DFS use and node rebuilds; NameNode heap from the count of files, directories and blocks rather than bytes; and worker memory, cores and disks from the workload's container profile.

open as a page

What would make you replace Oozie with Airflow to orchestrate a Hadoop platform's workflows?

level: principalimportance: should knowfreq 35%

basics

~20 s

Airflow wins when workflows must be generated and tested as code, span systems beyond the cluster, and be operated by people who need a usable UI. Oozie's XML workflows are static, YARN-bound and awkward to test.

open as a page

How do you decide which legacy Hadoop MapReduce pipelines to rewrite on Spark and which to leave alone?

level: principalimportance: should knowfreq 35%

basics

~20 s

Rank jobs by how much the MapReduce model actually costs them. Multi-job chains and iterative algorithms pay an HDFS round trip between every stage and repay a rewrite; stable single-pass map-only jobs are already IO-bound and rarely justify the risk.

open as a page

How would you design YARN queues for a cluster shared by teams with different SLAs?

level: principalimportance: should knowfreq 32%

basics

~20 s

Shape queues around workload classes and SLAs rather than around org charts: give SLA-bound pipelines guaranteed capacity with limited borrowing, let ad-hoc and batch work share a large elastic queue, cap concurrency and per-user limits, and enable preemption only where redoing work is cheap.

open as a page

How does a Sqoop import split an RDBMS table across parallel mappers?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

Sqoop takes a split column, queries its minimum and maximum, divides that range into equal slices, and gives each mapper a SELECT with its own WHERE range. Each mapper opens a separate connection to the source database.

open as a page

What does a Hadoop MapReduce job do differently when mapreduce.job.reduces is set to 0?

level: middleimportance: nice to knowfreq 30%

basics

~10 s

It becomes a map-only job: no partitioning, no sort, no shuffle. Each map task writes its records straight through the OutputFormat to HDFS as part-m-NNNNN files, and any configured combiner never runs.

open as a page

In YARN, what happens when an ApplicationMaster container fails, and what caps the retries?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

The ResourceManager notices the ApplicationMaster is gone and starts a fresh attempt in a new container elsewhere, up to yarn.resourcemanager.am-max-attempts (default 2). When attempts run out the whole application is marked FAILED, and by default its earlier containers are killed too.

open as a page

What happens to running YARN applications when the active ResourceManager fails over in an HA cluster?

level: seniorimportance: nice to knowfreq 32%

basics

~20 s

With work-preserving restart, containers already running on NodeManagers keep running. The standby ResourceManager becomes active, rebuilds application state from its ZooKeeper state store, and NodeManagers and ApplicationMasters resync with it. Clients retry transparently; new scheduling stalls for the failover window.

open as a page

When would you choose HDFS erasure coding over three-way replication in Hadoop 3?

level: seniorimportance: nice to knowfreq 22%

basics

~20 s

Erasure coding suits large, cold, rarely-read files: a Reed-Solomon scheme such as RS-6-3 tolerates three failures at about 1.5x storage overhead instead of replication's 3x. It costs read locality, CPU on recovery, and does not support append.

open as a page