How unit 4 is examined
This unit covers the MapReduce model, how jobs are built, run and monitored, the Hadoop daemons, HDFS and the execution modes; the marks sit in MapReduce, HDFS (storage process, architecture, goals) and the Spark versus MapReduce justification.
Employing Hadoop Map Reduce
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Medium weight</span>
Definition. <mark>MapReduce is a programming model in which a job is split into a map phase that turns input records into intermediate key-value pairs and a reduce phase that aggregates all values of each key, so that huge data sets are processed in parallel on a cluster of commodity machines.</mark>
Diagram.
<figure class="ds-fig" style="margin:1.4rem 0;overflow-x:auto"><svg xmlns="http://www.w3.org/2000/svg" id="dsfig-u4-01" viewBox="0 0 596 80" width="596" height="80" role="img" aria-label="MapReduce data flow. In = input file in HDFS, Spl = input splits, Map = map tasks, Shf = shuffle and sort, Red = reduce tasks, Out = output file in HDFS"><style>#dsfig-u4-01 .e{stroke:#454C5A;stroke-width:1.4;fill:none}#dsfig-u4-01 .e.hi{stroke:#2340B8;stroke-width:2.6}#dsfig-u4-01 .n{fill:#FFFFFF;stroke:#16181D;stroke-width:1.4}#dsfig-u4-01 .n.hi{fill:#E3E9FC;stroke:#2340B8;stroke-width:2.2}#dsfig-u4-01 .n.rb-b{fill:#16181D;stroke:#16181D}#dsfig-u4-01 .n.rb-r{fill:#BD3227;stroke:#BD3227}#dsfig-u4-01 text{font-family:"JetBrains Mono",ui-monospace,Menlo,Consolas,monospace;font-size:13px}#dsfig-u4-01 .t{fill:#16181D;font-weight:500}#dsfig-u4-01 .t.inv{fill:#FFFFFF;font-weight:700}#dsfig-u4-01 .kd{stroke:#16181D;stroke-width:1.2}#dsfig-u4-01 .dot{fill:#16181D}#dsfig-u4-01 .ann{fill:#2340B8;font-size:11px;font-weight:700}#dsfig-u4-01 .lbl{fill:#6F7787;font-family:system-ui,-apple-system,sans-serif;font-size:12px;font-weight:700}#dsfig-u4-01 .ptr{fill:#2340B8;font-size:12px;font-weight:700}#dsfig-u4-01 .ah{fill:#454C5A}#dsfig-u4-01 .ah.hi{fill:#2340B8}#dsfig-u4-01 .wl rect{fill:#FFFFFF;stroke:#DCE0E7}#dsfig-u4-01 .wl .t{font-size:12px;font-weight:700}#dsfig-u4-01 .wl.hi rect{fill:#2340B8;stroke:#2340B8}#dsfig-u4-01 .wl.hi .t{fill:#FFFFFF}html.dark #dsfig-u4-01 .e{stroke:#B1B7C3}html.dark #dsfig-u4-01 .e.hi{stroke:#8FA3FF}html.dark #dsfig-u4-01 .n{fill:#161920;stroke:#E6E8ED}html.dark #dsfig-u4-01 .n.hi{fill:#1E2748;stroke:#8FA3FF}html.dark #dsfig-u4-01 .n.rb-b{fill:#E6E8ED;stroke:#E6E8ED}html.dark #dsfig-u4-01 .n.rb-r{fill:#FF7E71;stroke:#FF7E71}html.dark #dsfig-u4-01 .t{fill:#E6E8ED}html.dark #dsfig-u4-01 .t.inv{fill:#0F1115}html.dark #dsfig-u4-01 .kd{stroke:#E6E8ED}html.dark #dsfig-u4-01 .dot{fill:#E6E8ED}html.dark #dsfig-u4-01 .ann{fill:#8FA3FF}html.dark #dsfig-u4-01 .lbl{fill:#858D9C}html.dark #dsfig-u4-01 .ptr{fill:#8FA3FF}html.dark #dsfig-u4-01 .ah{fill:#B1B7C3}html.dark #dsfig-u4-01 .ah.hi{fill:#8FA3FF}html.dark #dsfig-u4-01 .wl rect{fill:#161920;stroke:#2A2E37}html.dark #dsfig-u4-01 .wl.hi rect{fill:#8FA3FF;stroke:#8FA3FF}html.dark #dsfig-u4-01 .wl.hi .t{fill:#0F1115}</style><defs><marker id="ah4" viewBox="0 0 10 10" refX="9" refY="5" markerWidth="7" markerHeight="7" orient="auto-start-reverse"><path class="ah" d="M0,1 L9,5 L0,9 z"/></marker><marker id="ahh4" viewBox="0 0 10 10" refX="9" refY="5" markerWidth="7" markerHeight="7" orient="auto-start-reverse"><path class="ah hi" d="M0,1 L9,5 L0,9 z"/></marker></defs><path class="e" d="M59,40 L122.2,40" marker-end="url(#ah4)"/><path class="e" d="M162.2,40 L225.4,40" marker-end="url(#ah4)"/><path class="e" d="M265.4,40 L328.6,40" marker-end="url(#ah4)"/><path class="e" d="M368.6,40 L431.8,40" marker-end="url(#ah4)"/><path class="e" d="M471.8,40 L535,40" marker-end="url(#ah4)"/><circle class="n" cx="40" cy="40" r="18"/><text class="t" x="40" y="40" dy=".35em" text-anchor="middle">In</text><circle class="n" cx="143.2" cy="40" r="18"/><text class="t" x="143.2" y="40" dy=".35em" text-anchor="middle">Spl</text><circle class="n" cx="246.4" cy="40" r="18"/><text class="t" x="246.4" y="40" dy=".35em" text-anchor="middle">Map</text><circle class="n" cx="349.6" cy="40" r="18"/><text class="t" x="349.6" y="40" dy=".35em" text-anchor="middle">Shf</text><circle class="n" cx="452.8" cy="40" r="18"/><text class="t" x="452.8" y="40" dy=".35em" text-anchor="middle">Red</text><circle class="n" cx="556" cy="40" r="18"/><text class="t" x="556" y="40" dy=".35em" text-anchor="middle">Out</text></svg><figcaption style="font-size:.82em;opacity:.72;margin-top:.45rem">MapReduce data flow. In = input file in HDFS, Spl = input splits, Map = map tasks, Shf = shuffle and sort, Red = reduce tasks, Out = output file in HDFS</figcaption></figure>
Key points.
- The input file is divided into input splits (normally one per HDFS block) and each split is given to one map task, which runs on the node that stores the block.
- In the Map phase the mapper reads one record at a time and emits intermediate pairs $\langle k_2, v_2 \rangle$; for word count it emits $\langle word, 1 \rangle$ for every word.
- Shuffle and sort collect all pairs with the same key from every mapper, send them over the network to one reducer and sort them by key.
- In the Reduce phase the reducer receives $\langle k_2, list(v_2) \rangle$, aggregates the list (sum, count, max) and writes the final pairs to HDFS.
- Formally, $map: (k_1,v_1) \to list(k_2,v_2)$ and $reduce: (k_2, list(v_2)) \to list(k_3,v_3)$.
- The main components are the client that submits the job, the JobTracker (master, schedules jobs) with a TaskTracker on each slave (runs tasks); in Hadoop 2 these are the ResourceManager, NodeManager and ApplicationMaster.
- Failed tasks are simply re-run on another node, so the framework tolerates machine failure without losing work.
Example. Input lines "cat dog" and "dog cat dog": map gives (cat,1) (dog,1) (dog,1) (cat,1) (dog,1); after shuffle cat -> [1,1], dog -> [1,1,1]; reduce outputs (cat,2) (dog,3).
Spark versus MapReduce. Spark is faster because it keeps intermediate data in memory while MapReduce writes it to disk after every job.
| Point | MapReduce | Spark |
|---|---|---|
| Processing | Disk-based, results of map and reduce written to HDFS | In-memory, data cached as RDDs |
| Execution plan | Fixed map then reduce, one job per step | DAG of transformations optimised as a whole |
| Iterative jobs (ML, graphs) | Reads and writes disk every iteration, slow | Reuses cached RDD, often 10-100 times faster |
| Latency | High, minutes, batch only | Low, seconds, also streaming |
| Fault tolerance | Re-run task, data replicated in HDFS | Recompute lost partition from RDD lineage |
| Best for | Huge one-pass batch jobs | Iterative, interactive and real-time analytics |
Answer frame. Open with the definition; draw the data flow figure; develop the split, map, shuffle, reduce points, then the word-count example and the components; close with "MapReduce gives scalable, fault-tolerant parallel processing". For Spark, open with disk versus memory, give the table, close with an iterative workload such as logistic regression.
Pitfall: Do not say the reducer runs before shuffle; the reducer starts only after its map outputs are shuffled and sorted.
Asked: [7 marks] (Jun 2020, Nov 2023) What is Map Reduce programming model? Explain. Brief about the main component of Map Reduce. Asked: [7 marks] (Dec 2020) Justify: SPARK is faster than Map reduce.
Creating the components of Hadoop Map Reduce jobs
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. A MapReduce job is built from three user-written components, the Mapper, the Reducer and the Driver, plus optional Combiner and Partitioner.
Key points.
- The Mapper class overrides
map()to convert each input record into intermediate key-value pairs. - The Reducer class overrides
reduce()to aggregate the list of values of each key. - The Driver sets the job name, mapper, reducer, key and value types, and input and output paths, then submits the job.
- A Combiner is a local mini-reducer that cuts network traffic; a Partitioner decides which reducer gets each key.
Distributing data processing across server farms
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. A server farm is a cluster of many commodity machines on which Hadoop spreads both the data and the computation.
Key points.
- Data is split into blocks stored on many nodes, and processing is moved to the data, which is called data locality.
- Moving code is cheaper than moving terabytes, so network traffic falls sharply.
- Adding nodes increases capacity and speed almost linearly, which is horizontal scaling.
- Failed nodes are tolerated because blocks are replicated and tasks are re-run elsewhere.
Executing Hadoop Map Reduce jobs
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. Executing a job means submitting it to the cluster, where the scheduler splits it into tasks and runs them on the nodes.
Key points.
- The job is submitted with
hadoop jar job.jar Driver input output. - The client gets a job ID, copies the jar and splits to HDFS, and the ResourceManager (YARN) allocates containers.
- Map tasks run first, then shuffle, then reduce tasks; the output directory must not already exist.
- A failed task is retried on another node, and slow tasks may get speculative duplicates.
monitoring the progress of job flows
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. Monitoring means tracking the state, progress and counters of running jobs through web pages and the command line.
Key points.
- The JobTracker web UI (port 50030) in Hadoop 1, and the ResourceManager UI (port 8088) in Hadoop 2, list jobs with map and reduce percentage complete.
hadoop job -listandmapred job -status <id>show job state from the terminal.- Counters record bytes read, records processed and failed tasks.
- Task logs help to debug failures.
The Building Blocks of Hadoop Map Reduce Distinguishing Hadoop daemons
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. Daemons are the background processes that make up a Hadoop cluster.
Key points.
- NameNode is the HDFS master that stores metadata; DataNode stores the blocks and reports by heartbeat.
- Secondary NameNode periodically merges the edit log into the fsimage; it is not a hot standby.
- JobTracker is the Hadoop 1 master that schedules MapReduce jobs; TaskTracker runs tasks on each slave.
- In Hadoop 2 the ResourceManager and NodeManager (YARN) replace JobTracker and TaskTracker.
Investigating the Hadoop Distributed File System
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Medium weight</span>
Definition. <mark>HDFS is a distributed, fault-tolerant file system that stores very large files as fixed-size blocks replicated across the DataNodes of a commodity cluster, under the control of a single NameNode.</mark>
Diagram.
<figure class="ds-fig" style="margin:1.4rem 0;overflow-x:auto"><svg xmlns="http://www.w3.org/2000/svg" id="dsfig-u4-02" viewBox="0 0 467 338" width="467" height="338" role="img" aria-label="HDFS architecture. Cl = client, NN = NameNode (metadata), D1-D3 = DataNodes holding replicated blocks"><style>#dsfig-u4-02 .e{stroke:#454C5A;stroke-width:1.4;fill:none}#dsfig-u4-02 .e.hi{stroke:#2340B8;stroke-width:2.6}#dsfig-u4-02 .n{fill:#FFFFFF;stroke:#16181D;stroke-width:1.4}#dsfig-u4-02 .n.hi{fill:#E3E9FC;stroke:#2340B8;stroke-width:2.2}#dsfig-u4-02 .n.rb-b{fill:#16181D;stroke:#16181D}#dsfig-u4-02 .n.rb-r{fill:#BD3227;stroke:#BD3227}#dsfig-u4-02 text{font-family:"JetBrains Mono",ui-monospace,Menlo,Consolas,monospace;font-size:13px}#dsfig-u4-02 .t{fill:#16181D;font-weight:500}#dsfig-u4-02 .t.inv{fill:#FFFFFF;font-weight:700}#dsfig-u4-02 .kd{stroke:#16181D;stroke-width:1.2}#dsfig-u4-02 .dot{fill:#16181D}#dsfig-u4-02 .ann{fill:#2340B8;font-size:11px;font-weight:700}#dsfig-u4-02 .lbl{fill:#6F7787;font-family:system-ui,-apple-system,sans-serif;font-size:12px;font-weight:700}#dsfig-u4-02 .ptr{fill:#2340B8;font-size:12px;font-weight:700}#dsfig-u4-02 .ah{fill:#454C5A}#dsfig-u4-02 .ah.hi{fill:#2340B8}#dsfig-u4-02 .wl rect{fill:#FFFFFF;stroke:#DCE0E7}#dsfig-u4-02 .wl .t{font-size:12px;font-weight:700}#dsfig-u4-02 .wl.hi rect{fill:#2340B8;stroke:#2340B8}#dsfig-u4-02 .wl.hi .t{fill:#FFFFFF}html.dark #dsfig-u4-02 .e{stroke:#B1B7C3}html.dark #dsfig-u4-02 .e.hi{stroke:#8FA3FF}html.dark #dsfig-u4-02 .n{fill:#161920;stroke:#E6E8ED}html.dark #dsfig-u4-02 .n.hi{fill:#1E2748;stroke:#8FA3FF}html.dark #dsfig-u4-02 .n.rb-b{fill:#E6E8ED;stroke:#E6E8ED}html.dark #dsfig-u4-02 .n.rb-r{fill:#FF7E71;stroke:#FF7E71}html.dark #dsfig-u4-02 .t{fill:#E6E8ED}html.dark #dsfig-u4-02 .t.inv{fill:#0F1115}html.dark #dsfig-u4-02 .kd{stroke:#E6E8ED}html.dark #dsfig-u4-02 .dot{fill:#E6E8ED}html.dark #dsfig-u4-02 .ann{fill:#8FA3FF}html.dark #dsfig-u4-02 .lbl{fill:#858D9C}html.dark #dsfig-u4-02 .ptr{fill:#8FA3FF}html.dark #dsfig-u4-02 .ah{fill:#B1B7C3}html.dark #dsfig-u4-02 .ah.hi{fill:#8FA3FF}html.dark #dsfig-u4-02 .wl rect{fill:#161920;stroke:#2A2E37}html.dark #dsfig-u4-02 .wl.hi rect{fill:#8FA3FF;stroke:#8FA3FF}html.dark #dsfig-u4-02 .wl.hi .t{fill:#0F1115}</style><defs><marker id="ah5" viewBox="0 0 10 10" refX="9" refY="5" markerWidth="7" markerHeight="7" orient="auto-start-reverse"><path class="ah" d="M0,1 L9,5 L0,9 z"/></marker><marker id="ahh5" viewBox="0 0 10 10" refX="9" refY="5" markerWidth="7" markerHeight="7" orient="auto-start-reverse"><path class="ah hi" d="M0,1 L9,5 L0,9 z"/></marker></defs><path class="e" d="M56.3,159.2 L238.7,49.8"/><path class="e" d="M53.4,182.4 L155.6,284.6"/><path class="e" d="M57,177.5 L281,289.5"/><path class="e" d="M188,298 L277,298" marker-end="url(#ah5)"/><path class="e" d="M317,298 L406,298" marker-end="url(#ah5)"/><path class="e" d="M249,58 L175,280"/><path class="e" d="M258.1,58.7 L294.9,279.3"/><path class="e" d="M265.5,55.8 L416.5,282.2"/><g class="wl"><rect x="127.1" y="95.5" width="40.8" height="18" rx="9"/><text class="t" x="147.5" y="104.5" dy=".35em" text-anchor="middle">meta</text></g><g class="wl"><rect x="84.1" y="224.5" width="40.8" height="18" rx="9"/><text class="t" x="104.5" y="233.5" dy=".35em" text-anchor="middle">data</text></g><g class="wl"><rect x="148.6" y="224.5" width="40.8" height="18" rx="9"/><text class="t" x="169" y="233.5" dy=".35em" text-anchor="middle">data</text></g><g class="wl"><rect x="213.1" y="289" width="40.8" height="18" rx="9"/><text class="t" x="233.5" y="298" dy=".35em" text-anchor="middle">repl</text></g><g class="wl"><rect x="342.1" y="289" width="40.8" height="18" rx="9"/><text class="t" x="362.5" y="298" dy=".35em" text-anchor="middle">repl</text></g><circle class="n" cx="40" cy="169" r="18"/><text class="t" x="40" y="169" dy=".35em" text-anchor="middle">Cl</text><circle class="n" cx="255" cy="40" r="18"/><text class="t" x="255" y="40" dy=".35em" text-anchor="middle">NN</text><circle class="n" cx="169" cy="298" r="18"/><text class="t" x="169" y="298" dy=".35em" text-anchor="middle">D1</text><circle class="n" cx="298" cy="298" r="18"/><text class="t" x="298" y="298" dy=".35em" text-anchor="middle">D2</text><circle class="n" cx="427" cy="298" r="18"/><text class="t" x="427" y="298" dy=".35em" text-anchor="middle">D3</text></svg><figcaption style="font-size:.82em;opacity:.72;margin-top:.45rem">HDFS architecture. Cl = client, NN = NameNode (metadata), D1-D3 = DataNodes holding replicated blocks</figcaption></figure>
Key points.
- The NameNode is the master that keeps the namespace and the block-to-DataNode map in memory (fsimage and edit log) but stores no file data.
- DataNodes are the slaves that store the blocks on local disks and send heartbeats and block reports to the NameNode.
- A file is split into blocks of 128 MB (64 MB in Hadoop 1); a 500 MB file becomes 4 blocks, three of 128 MB and one of 116 MB.
- Each block is replicated, by default 3 times, so 500 MB occupies 1500 MB of raw disk.
- Write flow: the client asks the NameNode, receives a list of DataNodes, sends the block to the first DataNode, which forwards it down a pipeline to the second and third; acknowledgements return along the pipeline.
- Read flow: the client asks the NameNode for block locations and reads each block directly from the nearest DataNode.
- Rack awareness places the first replica on the writer's node, the second on a different rack and the third on another node of that second rack, which survives a rack failure while limiting cross-rack traffic.
- Fault tolerance: if a DataNode stops sending heartbeats, the NameNode re-replicates its blocks on other nodes.
Goals of HDFS.
- Fault tolerance: hardware failure is normal, so the system detects it and recovers automatically through replication.
- Very large data sets: files of gigabytes to terabytes, with a cluster scaling to thousands of nodes.
- Streaming data access: high throughput for reading whole files is preferred over low latency.
- Simple coherency: files follow write-once-read-many, so they are not edited in place, only appended.
- Moving computation to the data: tasks run where blocks live, since moving code is cheaper than moving data.
- Portability and commodity hardware: it runs on cheap machines and is written in Java, so it works across platforms.
Answer frame. For storage, open with the HDFS definition; draw the architecture figure; develop NameNode, DataNode, blocks, the 500 MB example, write pipeline, replication and rack awareness; close with fault tolerance. For goals, open with the definition and design assumptions, list goals 1-6 in order and close with write-once-read-many.
Asked: [7 marks] (Dec 2020, Nov 2023) Explain the process of data storage in Hadoop Distributed File System (HDFS) with the help of a suitable example. Describe the structure of HDFS in a Hadoop ecosystem using a diagram. Asked: [7 marks] (Jun 2020) Write down the goals of HDFS.
Selecting appropriate execution modes: local, pseudo-distributed, fully distributed
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. Hadoop can run in three modes that differ in how many JVMs and machines the daemons use.
Key points.
- Local (standalone) mode runs everything in one JVM on the local file system with no daemons, and is used to test and debug code.
- Pseudo-distributed mode runs all daemons as separate processes on one machine using HDFS, and simulates a cluster for development.
- Fully distributed mode runs the daemons across many machines, and is used in production.
- Configuration files such as
core-site.xml,hdfs-site.xmlandmapred-site.xmldiffer per mode.
Last-minute revision
- MapReduce: map emits $\langle k, v \rangle$, shuffle groups by key, reduce aggregates.
- $map: (k_1,v_1) \to list(k_2,v_2)$; $reduce: (k_2, list(v_2)) \to list(k_3,v_3)$.
- Word count output for "cat dog / dog cat dog" is (cat,2) (dog,3).
- Three job components: Mapper, Reducer, Driver; optional Combiner and Partitioner.
- Move computation to data (data locality).
- HDFS: NameNode holds metadata, DataNodes hold blocks; block 128 MB, replication 3.
- 500 MB file: 4 blocks, 1500 MB raw storage.
- Rack awareness: 1 replica local, 2 on another rack.
- HDFS is write-once-read-many, built for streaming access on commodity hardware.
- Spark is faster because of in-memory RDDs and DAG scheduling, mainly for iterative jobs.
- Modes: local, pseudo-distributed, fully distributed.
Memory hooks
- Map-Shuffle-Reduce: "split, sort, sum".
- NameNode = librarian with the catalogue; DataNodes = shelves with the books.
- HDFS replication "3-2-1": 3 copies, 2 racks, 1 local.
- Spark = RAM, MapReduce = disk.
- Modes go small to large: one JVM, one machine, many machines.
Coverage checklist
- Employing Hadoop Map Reduce: Q1 (MapReduce model and components), Q4 (Spark faster than MapReduce)
- Creating the components of Hadoop Map Reduce jobs: covered (no past question)
- Distributing data processing across server farms: covered (no past question)
- Executing Hadoop Map Reduce jobs: covered (no past question)
- monitoring the progress of job flows: covered (no past question)
- The Building Blocks of Hadoop Map Reduce Distinguishing Hadoop daemons: covered (no past question)
- Investigating the Hadoop Distributed File System: Q2 (data storage and structure), Q3 (goals of HDFS)
- Selecting appropriate execution modes: local, pseudo-distributed, fully distributed: covered (no past question)