For a nightly Spark job with a recurring hot key, how do you choose between AQE, salting, and isolating the key?
answer
- measure before choosing a remedy
- free configuration first, code last
- is the hot key even a real key?
- every hard-coded hot list rots
- watch the max-versus-median ratio over time
basics
~20 sStart with the free options: leave adaptive skew handling on and broadcast the small side if you can. Salt when a known hot key recurs and the operator is an aggregation or an unbroadcastable join. Isolate the key into its own path when one or two values dominate everything.
solid answer
~50 sOrder the options by cost of ownership. **AQE's skew-join splitting** costs nothing to leave on and handles shuffle joins whose skewed partition clears its thresholds — but it acts after the giant shuffle is written and ignores aggregations. **Broadcasting** the small side removes the shuffle entirely and is the best answer whenever it fits in memory. **Fixing the data** — dropping or separating null and sentinel keys, pre-aggregating, re-partitioning the source — is often the real fix and the cheapest to maintain. **Salting** is powerful and covers aggregations, but it hard-codes an assumption about the distribution that decays. **Isolating** the hot key into a separate branch (broadcast join or dedicated job) is best when one or two values dominate and the rest of the data is healthy. Then instrument it: track the stage's max-versus-median task duration so the day the distribution shifts is visible.
go deeper
Know the menu of fixes exists — broadcast, adaptive skew handling, salting — and that the right one depends on whether the operator is a join or an aggregation and on how big the small side is.
Compare the options concretely: what adaptive skew splitting does and does not cover, when a broadcast is possible, and what salting costs in replicated rows.
Show a worked decision on a real pipeline: the measurement, the option chosen, the option rejected and why, and how the fix was verified in the Spark UI afterwards.
Own the strategy beyond one job: whether the hot key should be designed out of the published dataset, whether a dominant tenant deserves isolation, and how skew is monitored so the fix's decay is caught before an SLA is.
## Frame the decision before picking a technique Three things determine the answer, and none of them are Spark configuration: 1. **Is the skew stable or shifting?** One dominant tenant every night is a design problem you can solve once. A different hot key every day needs a mechanism that adapts at runtime. 2. **Which operator concentrates the data?** A shuffle join, an aggregation and a window have different remedies, and the adaptive rule covers only the first. 3. **Is the hot key meaningful?** Nulls, `-1`, `"unknown"`, a default account, a test tenant, a bot — these carry no business value in the join at all, and the cheapest fix is upstream of Spark. ## The options, in increasing cost of ownership **Leave AQE on.** Adaptive skew-join splitting is enabled by default and costs nothing when it does not fire. It should be the baseline everywhere. Its limits are real, though: it acts only after the oversized shuffle partition has been written, only for shuffle joins, and only when the partition clears both a factor above the median and an absolute byte threshold. A stage that stalls at a smaller scale, or an aggregation, gets no help. **Broadcast the small side.** If one input fits comfortably in driver and executor memory, the join needs no shuffle, and the skew disappears with it. This is the highest-value fix per line of code, and it also argues for keeping dimension tables small and pruned enough to stay broadcastable — a modelling decision, not a tuning one. **Fix the data.** Filter or separately handle null foreign keys, which never match in an inner join anyway. Pre-aggregate the fact side before the join so the hot key contributes fewer rows. Re-partition the source so the hot value is not concentrated in one input file. Split a dominant tenant into its own table or its own daily run. These are usually the changes that survive, because they remove the imbalance rather than compensating for it. **Salt.** Random salt on the skewed side, exhaustive replication on the other, join or group on the composite key. It is the only technique that also fixes aggregation skew, and the only one that balances the map side and the shuffle write rather than just the read. The cost is permanent complexity in the pipeline, an N that must be sized, and often a hard-coded hot-key list — an assumption about the data that rots silently. Salting is worth it when the skew is large, recurring and cannot be broadcast away. **Isolate the hot key.** Split the job in two: the handful of dominant values take one path (typically a broadcast join or a specialised aggregation), everything else takes the normal path, and the results are unioned. This is often cleaner than salting when the hot set is tiny and well known, because each branch is a simple, readable query and neither carries salt machinery. It shares salting's weakness — the hot list must be maintained. ## What a lead is actually being asked The interviewer is testing whether you optimise for the incident or for the next two years. Signals of the latter: - **Measure first.** Task-duration and shuffle-read quantiles for the stage, then a count of the key. Do not choose a remedy from a guess. - **Cheapest reversible fix first.** Configuration, then a hint, then a data change, then code. - **Make the assumption visible.** If the fix encodes a hot-key list or a salt count, put it in one named place with a comment about how it was derived, and derive it at runtime where the cost allows. - **Instrument the decay.** Emit the stage's max/median task-duration ratio from the event log into whatever dashboard the team already watches. Skew fixes fail quietly. - **Consider isolation at the cluster level too.** A single tenant that skews every pipeline may deserve its own scheduled run and its own resource pool, rather than a skew remedy in each of a dozen jobs. - **Know when the answer is architectural.** If one customer is a third of the fact table, every downstream consumer will hit this. Fixing the shape of the published dataset once beats fixing it in ten places. ## The trap answers Adding executors, raising `spark.sql.shuffle.partitions`, or enabling speculation do not help a skewed task: the critical path is one partition, and none of those change how many partitions the hot key hashes to. Being able to say *why* those three fail is a large part of what separates a considered answer from a list of remembered knobs.
- How would you monitor for skew returning after the fix ships?Derive a per-stage max-versus-median task duration (or shuffle-read bytes) ratio from the event logs and publish it alongside the job's runtime. A ratio that creeps from 2x toward 20x is skew re-forming, and it is visible weeks before the SLA breaks. Alerting on runtime alone tells you only after the pipeline is already late.
- When is the right fix to change the published dataset rather than the job?When more than one consumer hits the same hot key. If one tenant is a third of the fact table, every downstream join and aggregation inherits the problem, and ten teams each write their own salt. Pre-aggregating, bucketing, or partitioning the published table — or splitting that tenant out — fixes it once at the source.
- A stakeholder asks to just double the cluster to fix the two-hour stage. What do you say?That a skewed stage's critical path is a single task holding a single key's rows, so extra executors sit idle behind it and the wall clock barely moves while the bill doubles. Show the max-versus-median task metrics as evidence, then propose the actual remedy and what it costs in engineering time.
saying these in an interview costs you the question
- Reaches for salting before checking whether a broadcast fits
- Proposes more executors or partitions for a single hot key
- Leaves a hard-coded hot-key list with no review or monitoring
- Treats null and sentinel keys as real keys needing distribution
- Optimises the incident without instrumenting for its return