Spark
Mid-levelHow do I pick the right shuffle partitions, executor memory, and core count?
There's no universal right answer, but there is a defensible starting point — and, more usefully, a way to tell when your current numbers are wrong.
Every guide gives you a formula. Almost none of them tell you how to verify the formula actually fit your job. Both matter.
Shuffle partitions: start from target partition size, not a fixed number
spark.sql.shuffle.partitions defaults to 200 regardless of data size, which is wrong for
almost every real workload — too many partitions for a small job (task-scheduling overhead dominates), too
few for a large one (each partition spills to disk). A better starting point: divide your shuffle stage's
total input size by a target of 128-200MB per partition.
Then verify: after the change, check whether task duration within that stage became more even (fewer long-tail tasks) and whether spill-to-disk bytes dropped. If spill is still high, you undershot; if most partitions are now near-empty, you overshot.
Executor memory: size it to the largest partition you expect, not the average
Executor OOMs almost always come from the largest partition in a stage, not the average one — so sizing
memory off average partition size under-provisions by construction. A reasonable floor is
largest_partition_size * 3-4 (headroom for deserialization overhead and any in-memory
aggregation), capped by what your cluster's instance type can actually offer per executor.
Core count: more isn't free past a point
Executor cores control how many tasks run concurrently within one executor's memory budget — pushing this too high means those tasks are all competing for the same heap, which shows up as more GC time, not more throughput. 4-5 cores per executor is a reasonable ceiling for most JVM-heavy workloads; going past that usually trades GC pressure for parallelism you don't actually gain.
The metric that actually tells you if you're right: GC time as a percentage of total task time. Under ~10% is healthy. Climbing past 15-20% means memory pressure is eating into real compute time, regardless of what your shuffle-partition count or core count happen to be set to.
See what this looks like on your own pipeline.
Upload a real Spark event log, dbt run_results.json, or Flink metrics export and get
concrete, approve-before-apply recommendations back — not another rule of thumb.