Designing Data-Intensive Applications
Ch. 10

MapReduce and Distributed Filesystems

MapReduce brought Unix-style batch processing to clusters with automatic parallelization and fault tolerance.

MapReduce wraps the Unix pipeline in a fault-tolerant distributed runtime. Map tasks extract data in parallel; shuffle groups by key; reduce tasks aggregate. The runtime handles scheduling, retries, and data locality.

Diagram
In practice

Early Twitter and LinkedIn analytics ran on Hadoop MapReduce over HDFS. Today AWS EMR and Databricks run Spark over S3 — same map-shuffle-reduce idea, but in-memory stages avoid writing every intermediate to disk. Immutable Parquet files on S3 replaced HDFS for most new data lakes.

LinkedIn at scale

Early member analytics and ad targeting ran on Hadoop MapReduce over HDFS — immutable input files, parallel map tasks, shuffle by key, reduce aggregates. The pattern scaled to petabytes before Spark replaced disk-bound intermediates.

typescript — Spark batch over immutable Parquet
// LinkedIn analytics — Spark batch over immutable Parquet on S3
const dailyActive = spark.read.parquet("s3://events/dt=2025-06-01/")
  .filter(col("event") === "login")
  .groupBy("user_id")
  .agg(count("*").as("logins"));
dailyActive.write.mode("overwrite").parquet("s3://metrics/dau/");
// Lineage graph retries failed stages — no full job restart
Diagram
Key Takeaways
  • Input files are split across cluster nodes; each runs a map function.
  • Shuffle sorts and groups intermediate key-value pairs by key.
  • Reduce functions aggregate all values for each key.
  • Failed map/reduce tasks are automatically retried on other nodes.
  • HDFS stores large immutable files replicated across the cluster.
  • Hadoop MapReduce on HDFS is legacy; S3 + Spark replaced it in most cloud stacks.
MapReduceSparkHDFSS3shufflebatch job