Identify skew and wasted work before increasing compute.
Mechanism and reasoning
Performance tuning starts with where time and resources go. Separate source read, parsing, shuffle, join, aggregation and output write. A long total duration does not reveal which stage limits the job. More workers can increase cost without improving a stage constrained by one hot key or a slow source.
Partitioning controls parallel work and data movement. Too few partitions can leave workers idle; too many tiny partitions can create scheduling and file overhead. A hot key can concentrate most of the data in one task despite a large partition count. Inspect task duration and input-size distributions, not only their average.
Join strategy depends on actual data size and engine behavior. Broadcasting a genuinely small dimension can avoid a large shuffle, but a mistaken size estimate can exhaust executor memory. Adaptive execution can revise some choices at runtime. It does not remove the need to understand skew, statistics and output correctness.
Apply filters and projection when semantics permit. Reading only needed columns and partitions can reduce I/O substantially. Be careful with functions that prevent partition pruning or with filters applied after an outer join, where moving the filter can change meaning. An optimization that changes results is a bug.
Record before and after using the same input snapshot and output checks. Compare elapsed time, resource-hours, shuffle bytes and cost where available. A faster run using many more machines may be a poor choice for a daily deadline that already had margin. In interviews, name the bottleneck, predict the effect of one change and preserve a correctness test.
Read the long tail of a stage profile
A distributed stage completes when its required tasks finish. The average task duration can be small while one partition controls the wall-clock result. Compare maximum and percentile durations with input bytes and records. A slow task with much more input suggests skew; a slow task with similar input may point to a host, spill, retry or external dependency problem.
Task group
Tasks
Input per task
Duration
Spill
Ordinary keys
99
1 GiB
1 minute
None
Hot tenant key
1
70 GiB
20 minutes
40 GiB
Same hot key after extra executors only
1
70 GiB
20 minutes
40 GiB
The third row is a hypothetical repeated measurement under an unchanged partitioning plan. It shows why additional idle executors need not split one key. The right next change targets how the key's work is represented. For associative operations such as sum, partial aggregation can distribute work and combine partials later. For exact median or order-sensitive logic, the decomposition requires a different argument.
sql
1-- Teaching two-stage sum; salt assignment must be stable per input row.2WITH partials AS (3 SELECT tenant_id, stable_salt, SUM(amount) AS subtotal4 FROM events5 GROUP BY tenant_id, stable_salt6)7SELECT tenant_id, SUM(subtotal) AS total8FROM partials9GROUP BY tenant_id;
This query illustrates the algebra; it does not guarantee that every engine chooses a faster physical plan. The input must avoid duplicate events under the dataset contract, and stable_salt must distribute the hot tenant's rows without changing their inclusion. Floating-point sums can vary slightly with evaluation order; use appropriate decimal semantics or documented tolerances when that matters. Null handling and overflow also belong to the correctness check.
A join requires separate analysis. Broadcasting a small dimension can avoid shuffling the large fact table, but “small” must refer to the actual serialized and in-memory representation under the engine's limits. Stale statistics can underestimate the dimension. A dimension that grows unexpectedly can trigger memory pressure or a different plan. Inspect the executed plan and runtime metrics instead of assuming a configuration hint guarantees a safe result.
Predicate placement is another source of accidental semantic change. Filtering a right-side column in the WHERE clause after a LEFT JOIN can remove rows with no match. Moving that predicate into the ON clause can preserve those left rows while changing the joined values. Either can be the intended query, but they are not automatically equivalent. A performance rewrite must preserve the original business result or explicitly change the contract.
Partition pruning should be demonstrated in scan metrics. If a date function wraps the partition column, the engine may or may not push the filter into the scan depending on its optimizer. Read the actual plan and bytes scanned. An estimated reduction from eight equal partitions to one is a useful expectation, but unequal partition sizes or metadata behavior can change the observed value. State the assumptions behind the arithmetic.
Compare resource use as well as elapsed time. Suppose a job falls from twenty minutes on ten workers to twelve minutes on twenty workers. Worker-minutes rise from 200 to 240, a 20 percent increase, while elapsed time falls 40 percent. That can be worthwhile for a strict deadline and wasteful for a flexible overnight job. Use the relevant cost and service objective rather than declaring every faster run better.
A disciplined tuning loop changes one main factor, reuses the same input snapshot, records the executed plan and checks output equivalence. If several changes are required together, state that their individual contributions remain uncertain. Keep the original profile and the successful result so the next engineer can see why the change was chosen and which growth condition may require another review.
Worked example
Invented job has 100 tasks. Ninety-nine finish in one minute; one hot-key task takes twenty minutes. Doubling the worker count does not split that key automatically under this plan. Investigate the key's share and whether a two-stage partial aggregation is valid. For a sum, aggregate salted subkeys first, then combine by the original key. For an order-sensitive operation, salting needs a different correctness argument. The optimization follows the operation's algebra, not a general rule that salting is always safe.
Exercise
A scan reads 800 GiB, but only one of eight equal date partitions is required. A predicate can prune the other seven without changing semantics. Estimate bytes read after pruning and one check before accepting the speedup.
Model solution and rubric
The simplified read volume becomes 100 GiB, an 87.5% reduction. Check the physical plan or scan metrics to confirm pruning happened, and compare result keys and aggregates against the original query on the same snapshot. The runtime improvement need not be 87.5% because parsing, shuffle and writes still contribute. Avoid claiming proportional speedup from bytes alone.
Score out of four: one point for the correct result, one for showing the intermediate reasoning, one for identifying the stated failure case, and one for a verification that could disprove the answer. Do not award the reasoning point for a tool name alone.
Failure modes and misconceptions
“More workers split every hot key.” The partitioning and operator determine that. One large key can remain one task.
“Fewer scanned bytes imply the same percentage reduction in runtime.” Shuffle, computation, scheduling and writes remain. Measure the actual stage and total costs.
Interview probe
Evidence class: recommended. Original practice.
A job has one task twenty times slower than the rest. Would you add workers?
Strong answer: First inspect skew and the task's input. If a single key dominates one partition, more idle workers may not help. I would consider a semantics-preserving split or pre-aggregation and verify output equivalence.
Follow-up: When can a broadcast join make the job less reliable?
Weak answer indicators: Adding workers without a stage profile; assuming bytes saved equal time saved; changing an outer join filter without checking semantics.
Sources
Technical references: Spark SQL performance; Apache Iceberg evolution. Sources support the documented mechanisms. The numbers, decisions, rubrics and interview prompts in this lesson are original teaching examples, not measurements or employer question claims.
One hot-key task remains 20 minutes after adding idle executors. What is the targeted next step?
AIncrease every timeout so the task can run longerBInspect key distribution and a valid partial-aggregation or repartitioning planCIncrease partitions without checking whether the key remains togetherDEnable broadcast for every join regardless of dimension size
A filter on the right table moves from WHERE into a LEFT JOIN condition. What must be checked?
AThe two placements always preserve all row valuesBThe optimizer guarantees faster executionCUnmatched left rows may now survive, changing output semanticsDThe move automatically deduplicates the dimension
Explain how you would identify skew and wasted work before increasing compute without reading the solution. State one assumption that could change your answer, and one observation that would make you revise it.
Not yetGetting thereConfident
Wrap-up
Tune the slow stage and keep the output contract fixed. Report both elapsed time and resource cost.