SysDesignPrep.com
Study guide 99 of 183

MapReduce and Spark

How large-scale batch processing works: the MapReduce model of map, shuffle and reduce, combiners, fault tolerance by re-execution, why Spark replaced it with in-memory DAGs, partitions and shuffles, joins and skew, and when batch jobs appear in system designs.

Reading is half of it. See this used in a real interview: walk through Design a Web Crawler →

Building a search index from a web crawl, computing daily billing from billions of click events, rebuilding autocomplete rankings from query logs, training data for recommendations: these are batch jobs over terabytes. MapReduce introduced the model that made such jobs simple to write and tolerant of machine failures; Spark made them much faster. Knowing how they work explains what a batch layer in your design actually does and where it gets slow.

The MapReduce model

You write two functions; the framework does the rest:

  1. Map: each input record produces zero or more (key, value) pairs. For word count: each word in a line emits (word, 1).
  2. Shuffle: the framework groups all pairs by key across the cluster, sorting and sending each key's values to one reducer.
  3. Reduce: for each key, combine its values. For word count: sum the ones.

The framework splits input into chunks (often aligned with distributed file system blocks), runs map tasks near the data, partitions map output by key hash for reducers, and writes final output back to storage.

Combiners

A combiner runs a local reduce on each mapper's output before the shuffle (summing counts per word on each machine), so far less data crosses the network. It works when the reduce operation is associative and commutative (sum, max, count).

Fault tolerance

Machines fail during long jobs. Map and reduce tasks are deterministic functions of their input, so the framework simply re-runs failed tasks elsewhere. Slow machines (stragglers) are handled with speculative execution: run a duplicate of a slow task and take whichever finishes first. See tail latency.

Why Spark replaced it

MapReduce writes intermediate results to disk between every stage, and multi-step pipelines become chains of jobs, each reading and writing the file system. Spark:

  • Represents the whole job as a DAG of transformations, optimised as a unit.
  • Keeps intermediate data in memory where possible, spilling to disk when needed.
  • Recovers lost partitions by lineage: recomputing them from their inputs rather than relying on replicated intermediate files.
  • Offers higher-level APIs: DataFrames and SQL with a query optimiser, plus streaming and machine learning libraries.

Iterative and multi-stage jobs run many times faster. The core ideas (partitioned data, shuffles, re-execution) remain the same.

Partitions and shuffles

  • Data is split into partitions processed in parallel by tasks.
  • Narrow transformations (filter, map) work within a partition; cheap.
  • Wide transformations (group by, join, distinct) need a shuffle: data moves across the network by key. Shuffles dominate cost and are where jobs fail from memory pressure.

Reduce shuffles by filtering and projecting early, pre-aggregating, and partitioning data by the keys you join and group on.

Joins

  • Shuffle (sort-merge) join: both sides are partitioned by the join key and merged. General but expensive.
  • Broadcast join: send a small table to every task and join locally, with no shuffle of the big table. The fastest option when one side fits in memory.
  • Bucketed tables: store data pre-partitioned by the join key to avoid repeated shuffles.

Skew

If one key has a huge share of records (a celebrity user, a "null" key), one task gets most of the work and the whole job waits for it. Fixes: salting the key in a first aggregation stage, isolating heavy keys, broadcast joins, and adaptive execution that splits skewed partitions automatically. See hot keys and skew.

Where batch jobs appear in designs

Batch versus stream

Streaming gives fresh results; batch gives complete, reproducible ones over all data and is easier to fix by rerunning. Most data-heavy designs use both: streaming for dashboards and alerts, batch as the source of truth. See batch and stream processing and windowing and watermarks.

Making jobs reliable

  • Idempotent outputs: write each run's output to a partition (by date) and overwrite it on rerun, so retries and backfills never duplicate.
  • Orchestrate dependencies and schedules with a workflow tool; backfill past dates on logic changes. See workflow orchestration.
  • Data quality checks before publishing results (row counts, null rates, totals versus yesterday).

Checklist

  • Map, shuffle, reduce; combiners for associative aggregations.
  • Re-execution and speculative tasks for fault tolerance.
  • Spark DAGs, in-memory processing and lineage-based recovery.
  • Minimise shuffles; broadcast small tables; bucket by join keys.
  • Handle skew with salting, isolation or adaptive execution.
  • Idempotent, partitioned outputs; orchestrated schedules and quality checks.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.