You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Executor统一内存管理下纯SQL执行时OOM异常原因问询

Why Does Spark SQL OOM Happen Without Using User Memory?

Great question—let’s dig into this because Spark’s unified memory model has some tricky edge cases that can trip you up even when you’re sticking strictly to Spark SQL queries and not touching user memory explicitly.

First, a quick recap to set the stage (since you already know this, but it helps frame the issues): The unified memory pool (total Executor memory × spark.memory.fraction, default 0.6) is split dynamically between execution and storage. Storage has a protected fraction (spark.memory.storageFraction, default 0.5) that execution can’t evict. User memory is static (total memory × (1 - 0.6)) and unmonitored, but you’re saying you’re not using it—so why the OOM?

Here are the most common culprits:

1. Execution Memory Limits Hit Even With Dynamic Sharing

Even though execution can steal from storage (except the protected 50%), there’s still a hard cap on how much execution memory can use, and SQL operations can easily push past it:

  • Large shuffle partitions: When running joins, group bys, or window functions, Spark shuffles data across Executors. If a single shuffle partition is too big (e.g., due to poor partitioning), it can’t fit into the available execution memory. Even if storage has free space, processing the partition requires first loading it into memory to serialize for disk overflow—and if the data is too large, this triggers OOM before overflow can happen.
  • Protected storage memory blocking execution: If storage is using most of its protected fraction (e.g., from cached broadcast tables or internal optimizations), the execution pool gets squeezed. For example, with 10GB total Executor memory, the unified pool is 6GB. If storage uses its full 3GB protected limit, execution only has 3GB left. A large shuffle partition needing 4GB will OOM here, even with unused user memory.

2. Hidden Storage Memory Usage From SQL Operations

You might not call cache() or persist() manually, but Spark SQL uses storage memory under the hood for optimizations:

  • Broadcast joins: When Spark uses a broadcast join, it caches the smaller table in storage memory. If that table is larger than expected (e.g., you thought it was 500MB but it’s actually 2GB), it eats into the storage pool, leaving less room for execution.
  • Internal query caching: Spark’s Catalyst optimizer may cache intermediate results (like from Exchange operators) to avoid recomputing them, using storage memory without your explicit input.
  • Window functions: Some window operations require caching partitioned data in storage to compute aggregates, which can consume significant memory if window partitions are large.

3. Java Object Overhead and Serialization Inefficiencies

Spark SQL processes data as Row objects, which come with inherent memory overhead (object headers, references, etc.) that adds up fast:

  • For wide tables (many columns), each Row has more metadata and references, pushing memory usage beyond the raw data size.
  • Using default Java serialization instead of Kryo inflates object sizes drastically. Even if your raw data seems manageable, the inflated object size can push you over memory limits.

4. Non-Memory-Pool Heap Usage

The Executor’s heap isn’t just split into execution, storage, and user pools—other critical overhead eats into the total heap:

  • Thread stacks: Each Task runs in its own thread, with a default 1MB stack per thread. For queries with many concurrent Tasks, this adds up quickly.
  • GC overhead: The JVM needs space for garbage collection (eden, survivor, old gen spaces). If the heap is nearly full, GC can’t free up space fast enough, leading to OOM.
  • Spark metadata: Shuffle indices, task status objects, and internal bookkeeping take up memory outside the unified pool. For large queries with thousands of tasks, this overhead can be significant.

5. Data Skew Causes Single-Task Overload

This is one of the most common SQL-specific OOM triggers. If your query has data skew (e.g., a join key with millions of records), a single Task ends up processing all that data. Even if total Executor memory is sufficient, that one Task’s memory demand can exceed the execution pool’s capacity. For example, if most keys have 10k records but one key has 10 million, the Task handling that key will try to load all 10 million records into memory—leading to OOM, even with unused user memory.

As you can see, there are plenty of scenarios where OOM occurs without touching user memory at all. User memory OOMs usually happen when using custom UDFs, storing large objects in driver/Executor-side code, or caching data outside Spark’s managed pools—but pure SQL queries can hit OOM from any of the issues above.

内容的提问来源于stack exchange,提问作者Mikel San Vicente

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 21:54:05