skip to content

In Spark, how does a reduce task learn where to fetch its shuffle blocks?

level: middleimportance: should knowfreq 50%

answer

  1. someone has to keep the address book
  2. the driver learns it when a map task finishes
  3. block ids grouped by the executor holding them
  4. MapStatus, then a direct pull

basics

~20 s

Every finished Spark map task reports a MapStatus to the driver's MapOutputTracker, recording which executor holds its output and how large each partition is. A reduce task asks the tracker for its partition's locations, then fetches those byte ranges directly.

solid answer

~40 s

When a map task finishes it sends the driver a `MapStatus` naming its `BlockManagerId` and the size of each reduce partition it produced. The driver's `MapOutputTrackerMaster` collects these for the whole map stage. When the reduce stage starts, each task's executor asks its `MapOutputTrackerWorker` for the block locations of its partition (`getMapSizesByExecutorId`); the answer is fetched once from the driver and cached per executor, not per task. A `ShuffleBlockFetcherIterator` then reads local blocks straight off disk and pulls remote ones over Netty from the owning executor's block manager, or from the external shuffle service if it is enabled. In-flight data is capped by `spark.reducer.maxSizeInFlight` (default 48m), split across several concurrent requests so one slow peer does not stall the reader.

code

text · 3 lines
text
Shuffle Read Size / Records: 41.2 GB / 812,004,113
Shuffle Read Blocks Fetched: local 1,204  remote 398,412
Shuffle Read Fetch Wait Time: 6.1 min

go deeper

for a junior

Know that finished tasks report where their output lives and that the next stage pulls that data over the network. You are not expected to name the tracker components yet.

for a middle

Explain the sequence: MapStatus to the driver, location lookup from the executor's tracker, then direct block fetches with local reads bypassing the network. Mention that in-flight bytes are bounded.

for a senior

Use the fetch path to read shuffle-read metrics in anger: fetch wait time, remote versus local bytes, block counts. Tie a stalled reader to a saturated peer, tiny blocks or a lost executor rather than to slow user code.

for a principal

Recognise the driver as a scaling bottleneck for very wide shuffles and reason about it when setting platform-wide partition-count conventions or deciding to push shuffle serving off the executors.

## The problem the fetch path solves After the map stage, one reduce partition's rows are scattered: a slice of it sits inside every map task's data file, on whichever node that task ran. Reduce task 57 must locate all those slices and stream them in. Spark solves this with a driver-side directory plus a client-side fetcher. ## Registration: MapStatus As each map task completes, it returns a `MapStatus` to the driver. It carries two things: the `BlockManagerId` (host, port and executor id) where the output lives, and the size of each reduce partition inside that task's data file. Sizes are stored compactly, and when a shuffle has a large number of partitions Spark switches to a `HighlyCompressedMapStatus` that records the exact sizes only of unusually large blocks and an average for the rest — otherwise the status objects themselves would bloat driver memory on a shuffle with tens of thousands of partitions. The driver's `MapOutputTrackerMaster` holds these statuses per shuffle id, which is why it is also the component that knows when map output has been *lost* and must be regenerated. ## Lookup: MapOutputTrackerWorker Every executor runs a `MapOutputTrackerWorker`. The first task on that executor needing shuffle X calls `getMapSizesByExecutorId`, which asks the driver for the serialized statuses, caches them locally, and returns the block ids grouped by the executor address that holds them, each with its expected size. Subsequent tasks on the same executor reuse the cache; the driver broadcasts large status maps rather than answering thousands of identical RPCs. This is one of the driver-side costs of very wide shuffles. ## Transfer: ShuffleBlockFetcherIterator The reduce task then builds a `ShuffleBlockFetcherIterator`, which splits requests into two classes: - **Local blocks** — output written by a map task that ran on this same executor. Read straight from disk with a seek into the data file, using the index file for the offset. No network involved. - **Remote blocks** — fetched over Netty from the peer executor's block manager, or from the node's external shuffle service when one is running. Requests are grouped by target address so many small blocks from one host travel together. Flow control matters here because a reduce task can be pulling from hundreds of peers at once. `spark.reducer.maxSizeInFlight` (default 48m) caps how much data may be outstanding, and the iterator issues several concurrent requests within that budget so a single slow node does not idle the reader. Companion limits bound requests and blocks in flight per remote address, which protects a hot node from being hammered by every reduce task at once. Blocks arrive out of order; the iterator hands them to the operator as they land, which is why shuffle reads are streamed and only aggregations or sorts that need all input buffer it. ## Where the fetch can fail Because the fetcher talks to a specific host and port that the tracker supplied, anything that invalidates that address surfaces here: the executor exited, the node was reclaimed, the disk holding the file is gone, or the shuffle service is saturated and times out. After a bounded number of retries the task raises a fetch failure, and the scheduler — not the task retry logic — takes over, because the fix is to regenerate the missing map output rather than to rerun the reader. ## Reading the metrics The stage's shuffle-read metrics decompose exactly along this path: remote bytes read, local bytes read, remote blocks fetched, and fetch wait time. A large fetch-wait time with modest bytes read points at many tiny blocks or a struggling peer, not at a slow computation. Local-versus-remote split tells you how much locality the scheduler actually achieved. ## What interviewers are checking Whether you can say who knows the locations (the driver), who does the reading (the reduce task's executor), and what limits the traffic. Candidates who believe reduce tasks broadcast a request to the cluster, or that the driver relays shuffle data, are describing a system that would not scale — and the misconception usually goes with an inability to explain fetch failures.

  • Why does Spark cap in-flight shuffle data with spark.reducer.maxSizeInFlight?
    A reduce task may be pulling from hundreds of peers simultaneously. Without a cap, arriving blocks would accumulate in the executor's memory faster than the operator consumes them and push the JVM into GC pressure or OOM. The 48m default bounds outstanding bytes while still issuing several parallel requests, so throughput stays high without letting the reader's buffers grow unbounded.
  • Why is a shuffle with very many partitions expensive for the driver?
    The driver stores a MapStatus per map task, each carrying per-partition sizes, so status memory grows with map tasks times reduce partitions. Spark mitigates this by switching to a highly compressed status above a configured partition count and by broadcasting the status map instead of answering per-task RPCs, but very wide shuffles still show up as driver memory and RPC pressure.
  • What do the local and remote shuffle-read metrics tell you?
    They split the read by where the block came from. A high remote share means map output rarely sat on the same executor as the reader, so the stage is network-bound. Large fetch wait time with small byte counts usually means too many tiny blocks — a partition count set far too high — or one degraded peer serving slowly.

saying these in an interview costs you the question

  • Says the driver relays the actual shuffle data
  • Thinks reduce tasks broadcast a request to all executors
  • Claims map output locations are stored in ZooKeeper or HDFS
  • Believes every shuffle read crosses the network
  • Assumes blocks must arrive in map-task order

context