RocksDB StateStore抗OOM测试异常:两种状态存储均触发OOM求助
问题背景
我编译了Spark 3.2.0版本,在Manjaro Linux系统(2核4线程、16GB内存)上用Java编写了基于Structured Streaming官方示例的WordCount程序,搭配Python服务器向其发送大量随机词,目的是验证RocksDB StateStore能否抵御OOM异常。但实际测试中,使用RocksDBStateStore和HDFSStateStore时都触发了OOM异常。原本预期只有HDFSStateStore会触发OOM,请问问题出在哪里?
SparkSession配置
SparkSession spark = SparkSession .builder() .appName("JavaStructuredNetworkWordCount") .config("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider") .config("spark.local.dir", "/home/username/sparkTemp") .config("spark.executor.memory", "15g") .config("spark.driver.memory", "15g") .config("spark.memory.offHeap.enabled", true) .config("spark.memory.offHeap.use", true) .config("spark.memory.offHeap.size", "50g") .config("spark.executor.memoryOverhead", "50g") .config("spark.sql.shuffle.partitions", 8) .config("spark.sql.streaming.stateStore.rocksdb.trackTotalNumberOfRows", false) .getOrCreate();
执行命令
/path/to/spark-submit --master local[*] --deploy-mode client --class org.example.Test4 --name Run /path/to/Test4-1.0-SNAPSHOT.jar --driver-memory 15g --executor-memory 15g
问题原因分析
1. 内存配置严重超出物理硬件上限
你的机器总内存仅16GB,但配置中spark.driver.memory和spark.executor.memory均设为15g,再加上spark.executor.memoryOverhead50g、spark.memory.offHeap.size50g,这些参数的总和远远超过了机器实际可提供的内存资源。即使在local[*]模式下,driver与executor共享同一进程,叠加后的内存需求直接突破系统内存上限,必然触发OOM——这和使用哪种StateStore无关,是整体内存配置逻辑错误导致的。
2. Local模式下的内存参数误解
使用--master local[*]时,Spark不会启动独立的executor进程,driver进程同时承担executor的角色。此时spark.executor.memory和--executor-memory参数完全无效,真正生效的只有spark.driver.memory和--driver-memory。你将driver内存设为15g,几乎占满了机器的物理内存,剩余空间还要分给系统进程、RocksDB本地缓存、JVM自身开销,完全没有冗余,很容易触发OOM。
3. RocksDB内存未做限制
RocksDB依赖自身的block cache缓存数据,默认情况下会尽可能占用可用内存。你没有配置spark.sql.streaming.stateStore.rocksdb.blockCacheSize来限制其缓存大小,导致RocksDB无节制占用内存,进一步加剧了内存耗尽的问题。
修正方案
- 匹配物理内存调整核心参数:将
spark.driver.memory设为8g,spark.memory.offHeap.size设为4g,spark.executor.memoryOverhead设为2g,预留至少4GB内存给系统和其他进程。 - 清理无效配置:去掉
spark.executor.memory和--executor-memory参数,避免在local模式下产生配置混淆。 - 限制RocksDB缓存:添加
spark.sql.streaming.stateStore.rocksdb.blockCacheSize=2g配置,控制RocksDB的内存占用。 - 调整并行度适配硬件:将
spark.sql.shuffle.partitions改为4,和机器4线程的硬件规格匹配,减少内存碎片化和不必要的资源开销。
内容的提问来源于stack exchange,提问作者user3551056

