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.

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.
  • 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).
  • 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.
  • 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

  1. Phases: Clearly distinguish map, shuffle/sort, and reduce phases.
  2. Key-Value Pairs: Show how input/output is structured (e.g., (word, count)).
  3. Optimizations: Mention combiners, partitioning, and speculative execution.
  4. Real-World Tie: Relate to Google Search, Ncell CDRs, or Daraz logistics.
  5. 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…