基于RocksDB的Spark Structured Streaming应用在K8S上OOMKilled问题
Spark Structured Streaming on Kubernetes 内存溢出(OOMKilled)问题排查与调优
问题背景
我在Kubernetes(K8S)上通过spark-submit --master k8s://...部署Spark Structured Streaming应用,采用RocksDB状态后端,配置spark.kubernetes.memoryOverheadFactor=2(因Spark自身会消耗配置的Executor内存,RocksDB需额外内存用于工作)。
应用运行时内存消耗缓慢增长至配置上限,随后维持在接近上限水平,一段时间后超出限制并被OOMKilled。具体表现:
- Driver内存上限2.5G,实际消耗约1.25G,无异常;
- 单个Executor内存上限5G,会不时被杀死,新启动的Executor内存从低位开始,但仍会缓慢增长至上限并再次被杀死。
额外信息:
- 无法获知应用实际状态字节数,Spark指标"State bytes"显示值远大于可用内存,但确定处理该工作负载的最小内存远低于配置上限;
- 单Executor应用中,Pod被杀死后从检查点恢复状态重启,仅消耗之前内存的一小部分即可运行;
- 当Pod内存接近上限时,会逐渐增长并偶尔回落,可见Spark或RocksDB中存在某种机制**"知晓"内存上限并尝试控制,但偶尔失效**。
相关配置
Spark 3.5.2 on Kubernetes Resource configuration (for single executor app) --conf spark.kubernetes.memoryOverheadFactor=2 \ --driver-cores=1 \ --driver-memory 500m \ --executor-memory 2g \ --executor-cores 1 \ --conf spark.executor.instances=1 Spark app conf "spark.sql.streaming.stateStore.providerClass" -> "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider", "spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled" -> "false" "spark.sql.streaming.stateStore.rocksdb.compactOnCommit" -> "true", "spark.sql.streaming.stateStore.rocksdb.boundedMemoryUsage" -> "true", "spark.sql.streaming.stateStore.rocksdb.blockSizeKB" -> "64", "spark.sql.streaming.stateStore.rocksdb.blockCacheSizeMB" -> "0",
已尝试调整上述配置项,仅改变内存增长速度和崩溃时间,基本现象未变。现咨询:
- 该内存使用控制机制的工作原理是什么?
- 有没有K8S上内存控制的文档或源码参考?
- 如何调优该机制,使内存使用维持在上限的90%左右?
解答
1. 内存使用控制机制的工作原理
你观察到的内存控制来自两层逻辑:
- RocksDB + Spark的内存绑定机制:当
spark.sql.streaming.stateStore.rocksdb.boundedMemoryUsage=true时,Spark会基于Executor的总内存配额(executor-memory+ 内存Overhead)计算RocksDB的可用内存上限,限制其内存表(MemTable)、索引缓存等组件的内存占用。RocksDB会在MemTable达到阈值时触发刷盘,同时通过定期压缩(Compaction)清理磁盘旧数据,减少内存中的冗余索引和缓存。 - K8S内存配额与Spark Overhead管理:
spark.kubernetes.memoryOverheadFactor=2会让内存Overhead等于executor-memory的2倍,这部分内存用于存放非堆内存、RocksDB的Native内存等。Spark会尝试控制这部分内存增长,但RocksDB的Native内存存在滞后性——比如压缩操作延迟、MemTable刷盘不及时,会导致内存短暂超出预期。
重启后内存占用低,是因为从检查点恢复时RocksDB仅加载必要索引和近期数据,不会一次性加载全部状态;而运行中随着数据处理,MemTable持续积累、缓存逐步建立,内存才会缓慢增长。
2. K8S上内存控制的文档与源码参考
- 官方文档:Spark官方文档中「Running Spark on Kubernetes」章节包含内存配置说明,「Structured Streaming State Management」章节详细讲解RocksDB状态后端的内存控制逻辑;
- 源码模块:
- K8S内存管理核心逻辑在
org.apache.spark.deploy.k8s包下的KubernetesExecutorBuilder类,负责计算内存Overhead和资源配额; - RocksDB内存控制逻辑在
org.apache.spark.sql.execution.streaming.state.RocksDBStateStore类,尤其是boundedMemoryUsage相关的内存计算与限制逻辑。
- K8S内存管理核心逻辑在
3. 调优建议(使内存维持在上限90%左右)
(1)精准配置内存Overhead
放弃memoryOverheadFactor,改用固定值spark.kubernetes.executor.memoryOverhead,避免Overhead过大导致K8S内存配额超出预期。比如你当前Executor总内存上限为5G,可设置executor-memory=1.5G,spark.kubernetes.executor.memoryOverhead=3.5G,让Spark更精准地计算RocksDB的内存上限。
(2)优化RocksDB内存管理
- 开启
spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled=true:启用日志检查点后,RocksDB可避免全量刷盘,降低MemTable内存压力,同时加快崩溃恢复速度; - 降低
spark.sql.streaming.stateStore.rocksdb.writeBufferSizeMB:缩小MemTable大小,触发更频繁的刷盘,避免内存积累。建议设置为64或128(默认256); - 开启
spark.sql.streaming.stateStore.rocksdb.asyncWrite:异步刷盘MemTable,避免刷盘阻塞数据处理,减少内存峰值; - 切换压缩策略:设置
spark.sql.streaming.stateStore.rocksdb.compactionStyle=universal,通用压缩策略更适配流式场景的动态数据,减少磁盘和内存占用。
(3)调整Spark内存分配比例
- 降低
executor-memory占比,增加Overhead:RocksDB主要使用Native内存(属于Overhead),比如设置executor-memory=1G,memoryOverhead=4G,总配额5G,给RocksDB分配更多Native内存空间; - 调整
spark.executor.memoryFraction:减少Spark堆内存占比,比如设置为0.3,让更多内存留给Overhead(RocksDB)。
(4)主动触发压缩与监控
- 在应用中添加定时逻辑,主动调用Spark状态API触发RocksDB压缩,避免内存积累到上限;
- 监控RocksDB核心指标(如
rocksdb.mem-table.size、rocksdb.block-cache.usage),当内存达到80%上限时,触发手动压缩或调整刷盘阈值。
内容的提问来源于stack exchange,提问作者andreikop
相关产品推荐
相关产品推荐

