Spark & Flink
Mid-levelWhy am I getting OutOfMemory errors, and how do I actually fix them?
"Just increase executor memory" fixes the symptom about half the time and wastes money the other half. Here's how to tell which case you're in.
An OOM error tells you memory ran out. It doesn't tell you why — and the fix is completely different depending on the cause.
1. Genuine under-provisioning
The straightforward case: your largest partition/window/state size legitimately doesn't fit in the memory you gave it. This is real, and the fix really is more memory (or a smaller partition size — see the shuffle-partitions article above).
Signal: heap usage climbs steadily through the task/job's lifetime and dies near the configured limit, with no sudden spike.
2. A skewed partition, not an undersized average
If 95% of your partitions use 2GB and one uses 40GB, increasing memory for everything to survive the one outlier is expensive and often still not enough. The real fix is addressing the skew itself — salting a skewed join key, or repartitioning on a more even key.
Signal: only a small number of tasks OOM, not most of them, and they correlate with a specific key range or partition.
3. Broadcast join gone wrong
Spark's broadcast join threshold (spark.sql.autoBroadcastJoinThreshold) decides whether a
smaller table gets copied in full to every executor. If that "smaller" table grew past what your executors
can hold — or Spark's size estimate is wrong (common after a filter or join upstream that changes
cardinality without updating statistics) — every executor OOMs trying to hold a broadcast that's no longer
actually small.
Signal: OOM happens very early in a stage, before any real shuffle work, and disabling broadcast join for that query makes it succeed (slower, but succeed).
4. Unbounded state in a streaming job (Flink specifically)
In Flink, OOM often isn't about one large batch — it's state that was never supposed to grow forever doing exactly that: a window that never closes because of a misconfigured watermark, or keyed state with no TTL that accumulates keys indefinitely. This one won't show up in a quick test; it shows up days or weeks into production as heap climbs slowly and doesn't recover.
Signal: checkpoint size growing steadily over days/weeks with no corresponding growth in input volume.
Before increasing memory: check GC time (Spark) or checkpoint duration trend (Flink) first. Both are strong tells for which of the four you're actually looking at, and both are numbers opti-pipe's rule engine already reads out of your event log or checkpoint history rather than you digging through the Spark/Flink UI by hand.
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.