Big Data and AnalyticsUnit 49 min read
MapReduce: Model, Phases, Use Cases & Optimization
Unit 4 of Big Data and Analytics explores MapReduce—its core architecture, the map and reduce phases, data partitioning, combiners, and real-world deployments (e.g., Hadoop). Learn how it processes petabytes of data across clusters, with comparisons to Spark and worked examples like log analysis or Ncell call-detail re
What is MapReduce?
MapReduce is a programming model and execution framework designed to process large-scale datasets (terabytes to petabytes) in parallel across clusters of commodity hardware. Developed by Google (2004) and later open-sourced as Apache Hadoop’s core engine, it simplifies distributed computing by breaking tasks into two key phases: map and reduce.
Why MapReduce?
- Scalability: Handles data too large for a single machine.
- Fault Tolerance: Automatically re-runs failed tasks.
- Simplicity: Abstracts low-level distributed computing details.
The MapReduce Model: Two-Phase Processing
MapReduce divides work into two phases, each with a distinct role:
1. Map Phase
- Input: Splits data into key-value pairs (e.g.,
(word, 1)for word counts). - Processing: Applies a user-defined map function to each pair, producing intermediate results.
- Output: Emits intermediate key-value pairs (e.g.,
(hello, 1),(world, 1)).
2. Shuffle and Sort Phase
- Shuffle: Groups intermediate values by key (e.g., all
(hello, x)pairs go to one node). - Sort: Orders values by key to prepare for reduction.
3. Reduce Phase
- Input: Receives grouped key-value pairs (e.g.,
(hello, [1, 1, 1])). - Processing: Applies a user-defined reduce function (e.g.,
sum) to aggregate values. - Output: Produces final results (e.g.,
(hello, 3)).
flowchart TD
A["Input Data\n(e.g., logs, CSV)"] --> B["Split into\nKey-Value Pairs"]
B --> C["Map Phase\n(user-defined function)"]
C --> D["Shuffle & Sort\n(by key)"]
D --> E["Reduce Phase\n(user-defined function)"]
E --> F["Final Output\n(e.g., word counts)"]How MapReduce Works: A Step-by-Step Trace
Example: Counting Words in a Corpus
Input: A file with text:
hello world
hello hadoop
world bigdata
Step 1: Map Phase
Each line is split into words, and each word emits a (word, 1) pair.
Map Task 1:
("hello", 1)
("world", 1)
Map Task 2:
("hello", 1)
("hadoop", 1)
Map Task 3:
("world", 1)
("bigdata", 1)
Step 2: Shuffle & Sort
Intermediate data is grouped by key:
hello: [1, 1]
world: [1, 1]
hadoop: [1]
bigdata: [1]
Step 3: Reduce Phase
The reducer sums values for each key:
("hello", 2)
("world", 2)
("hadoop", 1)
("bigdata", 1)
Key Components of MapReduce
| Component | Role | Example |
|---|---|---|
| InputSplit | Divides input into chunks for parallel processing. | 64MB or 128MB blocks of a log file. |
| Mapper | Processes input splits into intermediate key-value pairs. | map(word, 1) → (word, 1). |
| Combiner | Local reducer (optional) to minimize data transfer. | Sums (hello, 1) pairs before shuffle. |
| Partitioner | Distributes reduce tasks across nodes (default: hash partitioning). | hello → Node 1, world → Node 2. |
| Reducer | Aggregates intermediate values. | sum(1, 1, 1) → 3. |
| OutputFormat | Writes final results to storage (e.g., HDFS). | CSV, JSON, or a database. |
Optimizations in MapReduce
MapReduce’s performance depends on partitioning, combiners, and speculative execution:
1. Partitioning Strategies
- Hash Partitioning: Default method (e.g.,
hash(key) % num_reducers). - Range Partitioning: Useful for ordered data (e.g., timestamps).
- Custom Partitioners: Override
getPartition()for skewed data.
2. Combiners
- Local Aggregation: Reduces data transfer by pre-aggregating on mappers.
Example: In word count, a combiner sums
(hello, 1)pairs before sending to reducers.
3. Speculative Execution
- Problem: Slow tasks delay the job.
- Solution: MapReduce launches duplicate tasks on faster nodes.
4. Data Skew Handling
- Issue: Uneven key distribution (e.g., one word appears 1M times).
- Fix:
- Salting: Append random prefixes to keys (e.g.,
hello_1,hello_2). - Dynamic Partitioning: Adjust reducers at runtime.
- Salting: Append random prefixes to keys (e.g.,
MapReduce vs. Other Frameworks
| Feature | MapReduce (Hadoop) | Apache Spark | Google FlumeJava |
|---|---|---|---|
| Processing Model | Batch-only | Batch + Streaming | Batch + Incremental |
| Latency | High (minutes/hours) | Low (seconds) | Low |
| Language | Java, Python (via Streaming) | Scala, Python, Java, R | Java |
| Fault Tolerance | Checkpointing + replication | RDD lineage + checkpointing | Custom recovery |
| Use Case | ETL, log analysis | ML, real-time analytics | Interactive queries |
In the Real World
1. Google Search (PageRank)
- Idea Used: MapReduce processes trillions of web pages to compute PageRank scores.
- How:
- Map: Emits
(page_url, outbound_links)and(linked_page, vote)pairs. - Reduce: Aggregates votes to rank pages.
- Map: Emits
- Impact: Powers Google’s search engine since 2004.
2. Ncell Call Detail Records (CDR) Analysis
- Idea Used: MapReduce processes millions of daily CDRs to detect fraud or optimize networks.
- How:
- Map: Extracts
(customer_id, call_duration)from logs. - Reduce: Sums durations to flag anomalies (e.g., sudden spikes).
- Map: Extracts
- Impact: Reduces revenue loss from fake calls.
3. Daraz Order Fulfillment
- Idea Used: MapReduce routes orders to nearest warehouses based on inventory.
- How:
- Map: Emits
(product_id, location)pairs from order data. - Reduce: Assigns orders to the closest warehouse with stock.
- Map: Emits
- Impact: Faster deliveries, lower shipping costs.
Worked Example: Analyzing NTC Traffic Data
Problem: Nepal Telecom (NTC) wants to analyze call logs to predict peak hours.
Data: CSV with call_id, customer_id, start_time, duration.
MapReduce Solution
Mapper (Python-like Pseudocode)
for log in input_logs:
key = log.start_time.hour # Group by hour (e.g., 14 for 2–3 PM)
value = log.duration
emit_intermediate(key, value)
Reducer
total_duration = 0
for duration in grouped_values:
total_duration += duration
emit_final(key, total_duration) # e.g., (14, 4500) → 4500 sec of calls at 2–3 PM
Output
| Hour | Total Call Duration (seconds) |
|---|---|
| 14 | 4500 |
| 18 | 6200 |
| 22 | 3800 |
Insight: NTC can schedule maintenance during low-traffic hours (e.g., 3 AM).
Advantages and Disadvantages
✅ Advantages
- Scalability: Runs on thousands of nodes.
- Fault Tolerance: Automatically retries failed tasks.
- Simplicity: Hides distributed computing complexity.
❌ Disadvantages
- Batch-Only: Not suitable for real-time analytics (use Spark instead).
- High Latency: Jobs take minutes/hours.
- Resource Overhead: Requires HDFS storage.
Exam Tip
What Examiners Look For
- Phases: Clearly distinguish map, shuffle/sort, and reduce phases.
- Key-Value Pairs: Show how input/output is structured (e.g.,
(word, count)). - Optimizations: Mention combiners, partitioning, and speculative execution.
- Real-World Tie: Relate to Google Search, Ncell CDRs, or Daraz logistics.
- Code Snippets: Use pseudocode to explain mappers/reducers (no full implementation).
Common Pitfalls
- Ignoring Shuffle/Sort: Many students skip this critical phase.
- Confusing MapReduce with Spark: Spark is faster but not a replacement for all MapReduce use cases.
- Overlooking Skew: Always discuss how to handle uneven data distribution.
Visual Summary
classDiagram
class Mapper {
+map(key, value) (key, value)[]
}
class Combiner {
+reduce(key, values) value
}
class Partitioner {
+getPartition(key) int
}
class Reducer {
+reduce(key, values) value
}
Mapper --> Combiner : "Optional Local Aggregation"
Combiner --> Partitioner : "Groups by Key"
Partitioner --> Reducer : "Distributes Tasks"
Reducer --> Output : "Final Result"Based on the TU BITM syllabus for Big Data and Analytics (IT278), unit 4.
Discussion
Loading…