A distributed shuffle redistributes records by key so each reducer can see all values for its assigned keys.
Shuffle boundaries and local combiners
Find the network boundary
A row filter or projection can run inside each input partition. Grouping receipts by account cannot: one account may appear on 47 workers, so partial results must meet at a shared reducer. Mark that boundary in the execution plan and count bytes crossing it. Table partition design helps reads; it does not guarantee the compute engine has already colocated every grouping key.
Reduce before transfer
For associative operations such as integer sum, each mapper can collapse many receipt rows into one account subtotal before the shuffle. If a worker reads 83 lines for one account, it can send one subtotal. The reducer combines those subtotals. An average needs both sum and count; averaging mapper averages gives the wrong answer when mapper sizes differ.
Keep exactness in the reduction law
A combiner must be safe under arbitrary partitioning and repeated merge order. Integer addition works within its range. Distinct counts do not combine by summing counts because the same ID can occur on several workers. Use an exact set when bounded, or an explicitly approximate sketch when the product accepts error. Document overflow and decimal rounding for financial totals.
Observe the plan, not a slogan
Record input rows, mapper output records, shuffle bytes, reducer spill and longest task duration. If local combining barely reduces bytes because most keys appear once per partition, the added CPU may not be worthwhile. Large reducers can spill to disk even when total cluster memory looks free; hot-key analysis explains why.
Make reruns comparable
Pin the input snapshot, code revision and partitioning configuration when comparing a tuned job with its baseline. A later source correction can change both byte volume and result. Validate the output count and checksum through a replay manifest, then change one execution variable at a time.
Implementation
from collections import defaultdict
receipt_partitions = [
[("acct-47", 4300), ("acct-47", 2700), ("acct-83", 1900)],
[("acct-47", 1200), ("acct-83", 900)],
]
def local_subtotals(partitions):
partials = []
for partition in partitions:
subtotal = defaultdict(int)
for account_id, amount_cents in partition:
subtotal[account_id] += amount_cents
partials.extend(subtotal.items())
return partials
shuffled = local_subtotals(receipt_partitions)
assert len(shuffled) == 4
assert sum(amount for account, amount in shuffled if account == "acct-47") == 8200Performance and operating cost
Reading N rows is O(N) time; local maps hold O(Kp) keys per partition. If there are P partitions and K keys, the network sends at most P × K partials rather than N rows, although it can send nearly N when keys rarely repeat. The reducer still needs memory or spill space for its assigned key groups.
Common Mistakes
- Do not average partial averages without their counts.
- Do not assume a table partition removes a compute shuffle.
- Do not tune from wall time alone while spill and correctness are unmeasured.
