Big Data and AnalyticsUnit 47 min read
MapReduce: Model, Phases, Use Cases & Optimization
Unit 4 of Big Data and Analytics explores MapReduce—a scalable programming model for processing vast datasets across distributed clusters. This note covers its architecture, phases (map, shuffle, reduce), real-world applications (e.g., Google search indexing), and optimizations like combiners and partitioning, with vis
Core Concepts
1. What is MapReduce?
MapReduce is a programming model and execution framework designed by Google (2004) to process large-scale datasets (terabytes/petabytes) across clusters of commodity hardware. It abstracts:
- Distributed storage (HDFS in Hadoop).
- Parallel processing (splitting work into tasks).
- Fault tolerance (retrying failed tasks).
Why MapReduce?
- Handles data too big for a single machine.
- Tolerates hardware failures (no single point of failure).
- Simplifies writing distributed programs (hide complexity of parallelism).
How MapReduce Works: The 4-Phase Pipeline
MapReduce divides work into four key phases, visualized below. Each phase runs on a cluster of nodes (master + workers).
flowchart TD
A["Input Data (Split into Chunks)"] --> B["Map Phase\n(Process Key-Value Pairs)"]
B --> C["Shuffle & Sort\n(Group by Key)"]
C --> D["Reduce Phase\n(Aggregate Results)"]
D --> E["Output\n(Stored in HDFS)"]Phase 1: Input Splitting
- Input is split into fixed-size blocks (e.g., 64MB–128MB per block in Hadoop).
- Each block is assigned to a Map task on a worker node.
- Example: A 1TB log file → 16,000 blocks (if 64MB each).
Phase 2: Map Phase
- Map task processes each key-value pair independently.
- Output: Intermediate
<key, value>pairs (e.g.,(word, 1)for word count). - Key Design Principle:
- Embarrassingly parallel: No communication between mappers.
- Fault tolerance: If a mapper fails, the task is reassigned.
Worked Example: Word Count
| Input (Line) | Map Output (Key-Value) |
|---|---|
"hello world" |
(hello, 1), (world, 1) |
"hello mapreduce" |
(hello, 1), (mapreduce, 1) |
Phase 3: Shuffle & Sort
- Shuffle: Intermediate data is grouped by key and sent to reducers.
- Sort: Data is partitioned (e.g., via hash partitioning) to ensure all values for a key go to the same reducer.
- Optimization: Combiners (local reducers) merge intermediate data to reduce network traffic.
Phase 4: Reduce Phase
- Reduce task aggregates values for each key (e.g., summing counts).
- Output: Final result (e.g.,
(hello, 3),(world, 1)). - Fault tolerance: If a reducer fails, its work is reassigned.
MapReduce Architecture: Master-Slave Model
MapReduce runs on a cluster with:
- 1 Master (JobTracker/ResourceManager): Assigns tasks, monitors progress.
- N Workers (TaskTrackers/NodeManagers): Execute map/reduce tasks.
classDiagram
class JobTracker {
+Assigns tasks
+Monitors workers
}
class TaskTracker {
+Executes map/reduce tasks
+Reports progress
}
class HDFS {
+Stores input/output
}
JobTracker --> TaskTracker : "Assigns"
TaskTracker --> HDFS : "Reads/Writes"Optimizations in MapReduce
| Technique | Purpose | Example Use Case |
|---|---|---|
| Combiners | Reduce map output size | Word count: sum counts locally. |
| Partitioning | Distribute reduce load evenly | Hash partitioning for key distribution. |
| Speculative Execution | Mitigate slow tasks | Run duplicate tasks if one lags. |
| Data Locality | Minimize network I/O | Assign tasks to nodes storing data. |
Advantages and Disadvantages
| Advantages | Disadvantages |
|---|---|
| Scales to petabytes of data. | High latency (not real-time). |
| Fault-tolerant (handles node failures). | Not ideal for iterative algorithms (e.g., ML). |
| Simple programming model (functional). | Inefficient for complex joins. |
| Works on commodity hardware. | Shuffle phase can bottleneck performance. |
In the Real World
Google Search Indexing
- How: MapReduce processes web pages to build inverted indexes (word → URLs).
- Why: Scales to billions of pages; tolerates hardware failures.
Nepal Electricity Authority (NEA) Load Prediction
- How: Uses MapReduce to analyze historical consumption data (from smart meters) to predict peak hours.
- Impact: Optimizes power distribution and reduces outages.
Daraz/Nepal’s E-Commerce Logs
- How: MapReduce aggregates user behavior (clicks, purchases) to personalize recommendations.
- Example: If 80% of users who buy
Xalso buyY, suggestYto buyers ofX.
WhatsApp Message Processing
- How: MapReduce (or Spark) analyzes global message traffic to detect spam/bots at scale.
- Example: Count
(user, message_count)pairs to flag suspicious activity.
Worked Example: Traffic Route Optimization (Kathmandu)
Problem: Find the most congested roads in Kathmandu using GPS data. Solution: MapReduce pipeline:
- Map: Each GPS point →
(road_segment, 1). - Shuffle: Group by
road_segment. - Reduce: Sum counts →
(road_segment, total_vehicles). - Output: Rank roads by congestion.
Visualization:
flowchart TD
A["GPS Data\n(1M points/day)"] --> B["Map\n(road_segment, 1)"]
B --> C["Shuffle\nGroup by road_segment"]
C --> D["Reduce\nSum counts"]
D --> E["Output\nTop 10 congested roads"]MapReduce vs. Spark: Key Differences
| Feature | MapReduce | Apache Spark |
|---|---|---|
| Processing Model | Batch-only | Batch + streaming + ML |
| Latency | High (minutes/hours) | Low (seconds) |
| Fault Tolerance | Checkpointing | In-memory caching + RDDs |
| Use Case | ETL, log analysis | Real-time analytics, ML |
| Example | Google search indexing | Fraud detection in banks (e.g., NMB) |
Exam Tip
- Diagrams are Key: Draw the 4-phase pipeline (map → shuffle → reduce) in exams. Label:
- Input splits, map outputs, shuffle/sort, reducers.
- Word Count Example: Always ready to explain how
(word, 1)→(word, count)works. - Optimizations: Mention combiners and partitioning when asked about efficiency.
- Real-World Tie: Link MapReduce to Google, NEA, or Daraz in application questions.
- Limitations: Know why MapReduce is not used for real-time systems (e.g., stock trading).
Practice Questions
- How would you modify the MapReduce word count program to ignore stop words (e.g., "the", "and")?
- Why is the shuffle phase a bottleneck in MapReduce? Suggest two optimizations.
- Compare MapReduce’s data locality advantage with Spark’s in-memory processing. Which would you choose for:
- Analyzing 10TB of call logs (Ncell)?
- Detecting fraud in real-time (NMB Bank)?
Based on the TU BIM syllabus for Big Data and Analytics (IT278), unit 4.
Discussion
Loading…