skip to content

MPP Execution & Shuffle

A shared-nothing cluster splits the table across nodes, runs the same plan fragment everywhere, and redistributes rows over the network whenever a join or group-by needs matching keys co-located. Shuffle cost and data skew are the inevitable follow-ups.

on this pageshow

questions

6

In an MPP query plan, what is an exchange operator and when must the planner insert one?

level: middleimportance: must knowfreq 74%

answer

  1. rows that must meet must be on one machine
  2. the planner adds a step that moves rows
  3. it cuts the plan into stages
  4. hash the key, send to the owning node
  5. aggregate locally before sending

basics

~20 s

An exchange moves rows between compute nodes over the network. The planner inserts one wherever an operator needs matching key values on the same node and the current data placement does not already guarantee that — joins, grouping, global sorts, and the final gather.

solid answer

~50 s

Every operator that combines rows — a join, a `GROUP BY`, a `DISTINCT`, a global sort — requires that all rows sharing a key value be processed together on one machine. In a shared-nothing cluster that is only true if the data is already placed on that key, so the planner inserts an **exchange** (shuffle) operator to re-place the rows: hash the key, send each row to the node that owns that hash bucket. Exchanges cut the plan into stages — everything below one exchange is purely local, everything above depends on rows arriving. Common variants are hash repartition, broadcast (send one side to every node), and gather/merge (collect to the coordinator). A good planner pushes a **partial aggregation** below the exchange so only per-node partial results cross the network instead of raw rows, which is usually the single biggest reduction in shuffle bytes.

code

text · 7 lines
text
Stage 2  (runs on all nodes)
  FinalAggregate  group by user_id  sum(partial_amount)
    Exchange  HASH(user_id)   <-- rows cross the network here
      Stage 1  (runs on all nodes)
        PartialAggregate  group by user_id  sum(amount)
          Filter  event_date >= DATE '2026-01-01'
            Scan  events

go deeper

for a junior

Recall that an MPP engine sometimes has to send rows between machines mid-query, and that this step appears in the query plan as its own operator.

for a middle

Explain the co-location rule that forces the shuffle, name the hash-repartition, broadcast and gather variants, and describe how partial aggregation below the exchange shrinks what crosses the network.

for a senior

Read a real plan stage by stage: identify which exchange moves the most bytes, whether estimates were wrong, and which rewrite or layout change removes or shrinks it.

for a principal

Frame shuffle as the scarce cluster-wide resource. Argue for physical layouts and modelling choices that make the organization's hottest joins and groupings co-located by default rather than tuned query by query.

## The co-location invariant An MPP engine can only combine rows that are physically on the same machine. Two rows that must be compared, joined, or aggregated together have to arrive in the same worker's memory. That gives every relational operator a placement precondition: - A **join** on `a.k = b.k` needs all `a` rows and all `b` rows with the same `k` on one node. - A **`GROUP BY g`** needs all rows with the same `g` on one node to produce a final group total. - A **`DISTINCT`** or `COUNT(DISTINCT ...)` needs all occurrences of a value on one node to deduplicate it. - A **global `ORDER BY`** needs a single ordered stream, which ultimately means one merge point. When the current placement already satisfies the precondition — the table was distributed on exactly that key — nothing moves. When it does not, the engine has to re-place the rows at runtime. That re-placement is the exchange. ## What the exchange operator does An exchange has a sender side and a receiver side. On each node the sender applies a partitioning function to every outgoing row — typically `hash(key) mod N` — and writes it into a buffer destined for the node that owns that bucket. The receiver on each node accepts rows from all senders and feeds them to the operator above. Because every sender may talk to every receiver, a hash exchange is an all-to-all pattern; its cost grows with the number of rows moved and with how many bytes each row carries. Three variants show up in almost every engine's plan output: - **Hash repartition** — the general case above, used to co-locate join or grouping keys. - **Broadcast (replicate)** — every node sends its whole input to every other node, so each node ends up with a full copy of one side. Used when one input is small. - **Gather / merge** — all nodes send to a single point, usually the coordinator, for the final result or a global sorted `LIMIT`. Some engines also expose a *round-robin* or *rebalance* exchange, used to even out parallelism rather than to co-locate a key. ## Stages and why they behave like barriers Exchanges cut the plan into stages (fragments). Everything below an exchange runs entirely locally on the rows a node already holds; everything above it consumes rows that arrived from the network. Reading a plan therefore means reading it stage by stage: how many exchanges are there, how many rows and bytes does each one move, and which stage dominates the runtime. The exchange itself normally *streams* — rows flow as they are produced. What creates a barrier is usually the operator sitting above it: a hash join must fully build its hash table before probing, a final aggregate must see all input before emitting a group, a sort must see everything before emitting the first ordered row. So the stage boundary behaves like a synchronization point, and the slowest sender holds up the whole stage. ## Partial aggregation: the standard shuffle reduction The most important optimization around an exchange is doing work *before* it. For `SELECT g, SUM(x) FROM t GROUP BY g`, a naive plan shuffles every row of `t` on `g`. A good plan instead aggregates locally first: ```sql -- conceptually, per node: -- partial: SELECT g, SUM(x) AS s FROM local_rows GROUP BY g -- exchange on g -- final: SELECT g, SUM(s) FROM partials GROUP BY g ``` Now the network carries at most one row per distinct group per node instead of one row per input row. When the group count is small relative to the row count, this turns a terabyte shuffle into a megabyte shuffle. It works because `SUM`, `COUNT`, `MIN`, `MAX` and `AVG` (as sum plus count) are decomposable. It does not work as cleanly for `COUNT(DISTINCT ...)`, which is why exact distinct counts are notoriously the most shuffle-heavy aggregates. The same principle applies to projection: only the columns an upper stage actually needs should cross the exchange. Selecting fewer columns, filtering earlier, and pushing predicates below the exchange all reduce shuffle bytes directly. ## When no exchange is needed If both join inputs are already distributed on the join key, matching rows are already co-located and the join runs entirely locally — the plan shows no exchange between the scans and the join. The same holds for a `GROUP BY` on the distribution key: each node can produce final group totals without talking to anyone, and only the small result set is gathered. Designing physical layout so that the hottest joins and groupings are already co-located is the main lever an engineer has over shuffle cost. ## What to look for in a plan Count the exchanges. For each, note the partitioning (hash on which column, or broadcast), the estimated versus actual row count, and the bytes moved. A stage whose actual row count is orders of magnitude above the estimate is the classic cause of a bad exchange choice, and an exchange placed below rather than above a filter is a classic missed optimization.

  • Why does the planner place a partial aggregation below the exchange for a GROUP BY?
    Because SUM, COUNT, MIN and MAX are decomposable: each node can pre-aggregate its own rows and send one partial row per distinct group instead of every raw row. The final aggregate above the exchange combines the partials. When groups are few relative to rows, this cuts shuffled bytes by orders of magnitude, which is usually the dominant cost of the query.
  • Does a GROUP BY on the table's distribution key still require an exchange?
    No. If the table is hash-distributed on the grouping column, every row of a group is already on one node, so each node can compute final totals locally in a single aggregation step. Only the small result set is gathered. This is why aligning the distribution key with the most common grouping or join key is such a powerful design choice.
  • Is an exchange a blocking operator?
    The exchange itself streams rows as they are produced. What blocks is usually the operator above it — a hash build, a final aggregate, or a sort — which must consume all its input before emitting output. That is what makes a stage boundary behave like a barrier and why the slowest sending node determines when the next stage can finish.

saying these in an interview costs you the question

  • Believing the whole table is always sent to the coordinator
  • Thinking an exchange is inserted for every query
  • Assuming shuffle cost depends only on row count, not row width
  • Not knowing local pre-aggregation happens before the shuffle
  • Confusing shuffle with reading data from storage

context

open as a page

One MPP node runs at 100% while the rest idle during a large GROUP BY — what is happening and how do you fix it?

level: seniorimportance: must knowfreq 66%

basics

~20 s

The shuffle key is skewed: a few values, often NULL or a placeholder, own most of the rows, so one node receives that whole bucket. Runtime equals the slowest node. Fix by isolating or salting the hot values, or by pre-aggregating before the shuffle.

open as a page

In a shared-nothing MPP warehouse, how is one table's data divided across the cluster?

level: juniorimportance: should knowfreq 55%

basics

~20 s

Each compute node owns an exclusive subset of the table's rows and can read only that subset. Rows are placed by hashing a column, by round-robin, or by replicating the whole table everywhere. Every node then runs the same plan over its own share.

open as a page

How does an MPP planner choose between broadcasting one join input and hash-redistributing both?

level: middleimportance: should knowfreq 63%

basics

~20 s

It compares network cost. Broadcasting sends the smaller input's bytes to every node, so cost scales with node count; redistributing sends each row of both inputs once. Broadcast wins when one side is small enough to fit in every node's memory.

open as a page

Why does an MPP query spill to disk after a shuffle even when total cluster memory looks sufficient?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Cluster memory is not one pool. Each worker gets a fraction of one node's RAM, divided again among concurrent queries, and a shuffle can hand one worker far more than its even share. That worker spills alone while the cluster looks half empty.

open as a page

Doubling an MPP cluster's nodes barely speeds up a shuffle-heavy query — how do you diagnose it and what do you change?

level: principalimportance: should knowfreq 40%

basics

~20 s

Scale-out only shrinks the local work. Redistribution volume is roughly constant in node count, broadcast volume grows with it, skewed keys stay on one node, and the final gather stays serial. Attribute the runtime per stage before buying capacity.

open as a page