A skewed key sends a disproportionate share of records to one task, making that task set the runtime and spill risk of the stage.
Skewed keys and salted aggregation
Measure the heavy tail
A cluster can report 70% idle capacity while one reducer handles 63% of invoice rows for an umbrella account. Inspect per-task input bytes, spill and duration, then compare key frequency from the actual interval. The mean partition size hides the long tail. Storage skew and execution skew are related but measured at different boundaries.
Salt a decomposable calculation
For an associative sum, spread a heavy account across several deterministic subkeys derived from immutable invoice ID, compute partial sums, then merge them by the original account. The salt must be stable on retry. Random salt changes task assignment and makes incident replay harder, even if the final integer sum happens to match.
Keep the identity intact
Salting is an execution detail, not a new business key. The published table still has one row per account and day. After merging salts, assert uniqueness at that grain and reconcile total cents with the unsalted input. A distinct-customer count cannot be repaired by adding salted distinct counts because the same customer may occur in several salts.
Compare alternatives
Adaptive skew handling, a broadcast lookup, preaggregation, or separate treatment of a small set of heavy keys may be simpler. Salting every key adds an extra grouping stage and can worsen small jobs. Begin with the executed plan and tune only the measured bottleneck; local combiners may cut traffic before any salting.
Bound the split
Eight salt buckets can reduce one reducer input, but they also create eight partial states and a second merge. Choose a bucket count from the heavy key size and target task size. Monitor the largest salted bucket, not just the average, because poor hash distribution or a low-cardinality event ID can preserve skew.
Implementation
from collections import defaultdict
from hashlib import sha256
invoice_lines = [
{"invoice_id": "inv-47", "account": "parent-9", "cents": 4200},
{"invoice_id": "inv-48", "account": "parent-9", "cents": 3100},
{"invoice_id": "inv-49", "account": "parent-9", "cents": 1800},
]
def salted_totals(lines, buckets=4):
partial = defaultdict(int)
for line in lines:
salt = int.from_bytes(sha256(line["invoice_id"].encode()).digest()[:4], "big") % buckets
partial[(line["account"], salt)] += line["cents"]
merged = defaultdict(int)
for (account, _salt), cents in partial.items():
merged[account] += cents
return dict(merged)
assert salted_totals(invoice_lines) == {"parent-9": 9100}Performance and operating cost
Hashing N invoice IDs and summing values is O(N) expected time, with O(K × S) partial state for K heavy keys and S salts. The second grouping adds shuffle and task overhead. This only helps when it lowers the largest task enough to outweigh that overhead.
Common Mistakes
- Do not salt a non-decomposable metric and sum its partial answers.
- Do not publish salted keys as customer identities.
- Do not select bucket count from averages while the largest task still spills.
Read next
- Shuffle boundaries and local combiners
- Partition planning, pruning and skew
- Broadcast joins and build-side limits
- Project: release a distributed invoice rollup
- Fact grain and measure additivity
Continue the workflow: Spatial cell resolution and hotspot layout.
