Doubling an MPP cluster's nodes barely speeds up a shuffle-heavy query — how do you diagnose it and what do you change?
answer
- more nodes shrink only the local part
- redistribution moves the same bytes either way
- copying a side to everyone costs more per node added
- a hot key never splits
- measure per stage before buying capacity
basics
~20 sScale-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.
solid answer
~50 sAdding nodes halves the per-node scan, and nothing else. Hash redistribution still moves about the same total bytes, now over more connections with smaller messages; a broadcast moves the small side to every node, so its cost rises *linearly* with node count; a hot key still lands on exactly one worker; and the coordinator's final merge for a large `ORDER BY ... LIMIT` is serial no matter what. So the fix starts with attribution: from the profile, split runtime into scan, exchange, and the operator above each exchange, and check per-worker balance. If scan dominates, scale-out works. If exchange or one worker dominates, the answer is structural — co-locate the hottest join by aligning the distribution key, denormalize or pre-join so the shuffle disappears, materialize the recurring aggregate, or narrow what crosses the wire. Scaling up (fewer, larger nodes) can beat scaling out for shuffle-bound work because more of the join stays inside one machine.
code
text · 10 lines-- profile attribution before and after doubling the cluster
stage 16 nodes 32 nodes
scan 210 s 98 s <- scaled well
exchange 430 s 405 s <- barely moved
final agg 120 s 112 s
gather 40 s 41 s <- serial tail
total 800 s 656 s
exchange stage, rows received per worker (32 nodes):
median 41,000,000 max 1,900,000,000 <- skew, not capacitygo deeper
Recall that adding machines speeds up the part of the work each machine does alone, and does not help the part where machines must send data to each other.
Explain why redistribution volume stays flat, broadcast volume grows with node count, and a skewed key stays on one worker regardless of cluster size.
Do the attribution: time per stage, per-worker balance, estimate versus actual on each exchange, and the serial tail. Then pick the structural fix that removes or narrows the shuffle.
Own the trade between recurring capacity spend and one-off layout change. Defend spending when attribution shows scan-bound growth, and refuse it when it would only hide skew or an estimate-driven bad plan.
## What scale-out actually buys Adding compute nodes to a shared-nothing cluster divides the *local* portion of the work: each node scans fewer bytes, filters fewer rows, and builds smaller local hash tables. If a query's time is dominated by scanning and local computation, doubling the nodes genuinely approaches halving the runtime. Everything that is not local behaves differently, and a shuffle-heavy query is mostly not local. ## The four things that do not improve **Hash redistribution volume is roughly invariant.** Redistributing an input moves each row exactly once regardless of how many nodes exist. Doubling `N` does not halve the bytes on the wire; it splits the same bytes into twice as many streams. Per-node received volume does halve, which helps the receiving operator, but aggregate network work does not fall, and the all-to-all connection count grows quadratically, so per-message efficiency drops as messages get smaller. Beyond some point the coordination overhead eats the parallelism gain. **Broadcast volume gets worse.** Replicating a side to every node costs `size × N`. Double the cluster and that traffic doubles, while the same hash table is built redundantly on twice as many machines. A plan that was comfortably broadcast-based on a small cluster can become network- and memory-bound on a large one — the query genuinely gets *slower* per unit of hardware. **Skew is invariant.** All rows sharing a key value are indivisible under hashing. If one value owns 30% of the rows, one worker owns 30% of the work at any cluster size. The critical path is unchanged and the extra nodes idle. **Serial tails stay serial.** The coordinator's final merge for a global sort with a large result, single-threaded result serialization to the client, planning and metadata operations, and any per-query fixed startup cost are all outside the parallel region. This is Amdahl's law with a cloud invoice attached: as the parallel part shrinks, the serial remainder becomes the whole runtime. ## The diagnosis, in order 1. **Attribute the runtime by stage.** From the query profile, get wall time (or CPU time) per stage, and within each stage separate scan/local work from exchange transfer from the accumulating operator above it. If scan is 80% of the time, scale-out is the right lever and you are done. 2. **Check per-worker balance in the dominant stage.** Rows and bytes received per worker. A max-to-median ratio far above one is skew, and no amount of hardware fixes it. 3. **Check the exchange types and their estimated-versus-actual rows.** A broadcast whose actual cardinality is orders of magnitude above the estimate is a statistics problem masquerading as a capacity problem. 4. **Check the tail.** How long after the last worker finishes does the query take to return? A large gap points at the gather/merge or at result serialization, not at the cluster. 5. **Check what else was running.** On a shared cluster, added nodes also add concurrency capacity; if admission changed at the same time, per-query memory and network share changed too. Only after this does a capacity conclusion mean anything. ## The structural changes worth making **Co-locate the hot join.** If one join key drives most of the workload, distributing both tables on it removes the exchange entirely for those queries. This is the single largest available win and also the most expensive to change later, which is why it belongs in schema design rather than in tuning. **Remove the join.** Pre-joining or denormalizing the recurring join into a wide table trades storage and load complexity for zero shuffle at query time. In columnar storage the extra columns are cheap to skip when unused, which makes wide tables far more viable than transactional instincts suggest. **Materialize the recurring aggregate.** If the same expensive shuffle runs on a schedule to serve dashboards, computing it once into a summary table converts a repeated large shuffle into a cheap read. **Narrow the wire.** Project only the needed columns, push filters below the exchange, ensure partial aggregation is happening. Often worth a larger factor than any hardware change and it applies to every future run. **Fix skew at the source.** Sentinel defaults from upstream loaders and grain choices that concentrate rows are data-model defects; salting is a workaround, correcting the model is the fix. ## Scale up versus scale out For shuffle-bound work, fewer larger nodes can beat more smaller ones. More cores per machine means more of the join partners are already in the same address space, less traffic crosses the interconnect, and broadcast replication is paid fewer times. For scan-bound work the opposite holds — many small nodes maximize aggregate I/O bandwidth. Knowing which regime a workload is in is the point of the attribution exercise, and a mixed fleet may reasonably run different workloads on differently shaped compute. ## The judgment call Capacity is the fastest lever and the one with a recurring bill; layout and modelling changes are slow, risky and permanent. The defensible position is to spend on capacity when attribution shows scan-bound work, sustained demand growth or a deadline that cannot wait for a model change — and to refuse it when attribution shows skew, an estimate-driven bad plan choice, or a shuffle that a layout change would delete outright. Buying nodes to hide those simply raises the floor of every future bill.
- When would you deliberately scale up to fewer, larger nodes instead of out?When attribution shows the query is shuffle- or broadcast-bound rather than scan-bound. Larger nodes keep more join partners inside one machine, cut interconnect traffic, and pay broadcast replication fewer times. Scan-bound workloads want the opposite — many nodes maximizing aggregate I/O bandwidth. A platform serving both may justify differently shaped compute per workload.
- How do you decide between fixing the physical layout and just accepting a larger cluster?Compare the recurring cost of the capacity against the one-off cost and risk of the change, weighted by how many queries the layout serves. A distribution key that removes shuffle from a dozen daily pipelines is worth a migration; a single quarterly report is not. Refuse capacity when attribution shows skew or a bad plan, because hardware raises the floor of every future bill without fixing either.
- Why can the all-to-all pattern of a hash exchange itself limit scale-out?Every sender may talk to every receiver, so connection count grows roughly with the square of the node count while each message carries less data. Fixed per-message and per-connection overheads then consume a larger share of the transfer, and the stage cannot finish until its slowest sender does. Past some cluster size the coordination overhead outgrows the parallelism the extra nodes provide.
saying these in an interview costs you the question
- Assuming runtime scales inversely with node count
- Sizing the cluster before attributing time per stage
- Forgetting broadcast traffic grows with node count
- Believing more nodes dilute a hot key
- Ignoring the serial final merge and result return