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

基于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",

已尝试调整上述配置项,仅改变内存增长速度和崩溃时间,基本现象未变。现咨询:

  1. 该内存使用控制机制的工作原理是什么?
  2. 有没有K8S上内存控制的文档或源码参考?
  3. 如何调优该机制,使内存使用维持在上限的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相关的内存计算与限制逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 05:22:10