Infra
Mid-levelDriver memory vs executor memory: they're not the same tuning problem
“Just increase the memory” works differently depending on which memory - and increasing the wrong one just delays the same crash.
The driver runs a fundamentally different job than an executor: it builds and schedules the DAG, not process partitions of data - which means its memory pressure comes from different places entirely.
What actually runs on the driver
Plan building and DAG scheduling, coordinating task assignment to executors, collecting accumulator
values, and anything that pulls data back to the driver process explicitly - collect(),
toPandas(), take() on a large result. None of that is "processing partitions
of your dataset," which is the executors' job entirely.
Why driver OOM has different real causes
The most common one by far: someone called .collect() or .toPandas() on a
DataFrame that's much bigger than expected, pulling the whole thing into the single driver process's
memory instead of leaving it distributed. Second most common: a broadcast join (see the broadcast-join
article) where the "small" side wasn't actually small - the driver collects that data before pushing it
to executors, so it OOMs first. Third: a job tracking an enormous number of partitions or a very deep
DAG, where scheduling overhead alone - not any actual data - eats the driver's heap.
Why bumping executor memory doesn't touch any of this
spark.executor.memory and spark.driver.memory are genuinely separate
settings controlling separate JVM processes - increasing one does nothing for pressure on the other.
A driver OOM from an unexpected collect() needs either a code fix (aggregate before
collecting, or don't collect at all) or spark.driver.memory specifically; throwing more
executor memory at it is addressing a process that was never the one running out.
Quick diagnostic: an executor OOM shows up as a specific task failing with an
ExecutorLostFailure or similar in the event log. A driver OOM kills the whole application
at once, often with no single failing task to point to - a different shape entirely, and the first
clue you're looking at the wrong knob.
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.