Spark

Mid-level

How 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.

Read in:
Shuffle partitions Executor memory Core count
What to size first, in order — not a measured breakdown.

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.

# from the Spark UI's Stages tab, "Shuffle Read" for the relevant stage shuffle_partitions = shuffle_read_bytes / (150 * 1024 * 1024)

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.