Spark Structured Streaming使用RocksDB状态存储磁盘未清理问题
Spark 3.2.x RocksDB状态后端磁盘冗余膨胀解决方案
问题根因
这是Spark 3.2.1版本RocksDB状态后端默认配置适配缺陷导致的:状态TTL过期逻辑本身正常,但RocksDB的SST文件压实、墓碑标记清理默认没有和流处理的状态提交逻辑绑定,过期状态删除后留下的逻辑墓碑、旧版本SST文件、WAL残留文件不会被自动物理删除,才会出现有效状态仅25MB但磁盘占用达45GB的情况,和状态本身的大小没有直接关系。
阈值触发清理的配置方案
所有配置均在spark-submit提交参数中添加即可,无需修改业务代码:
- 配置大小阈值自动触发压实
直接设置spark.sql.streaming.stateStore.rocksdb.compactionOnCommitSizeThresholdMB参数为你期望的阈值(单位MB),例如设置为1024即代表单RocksDB状态实例大小超过1GB时,在当前批次提交完成后自动触发全量压实,物理删除冗余的SST文件、墓碑标记、临时WAL文件,将磁盘占用压缩到和实际有效状态匹配的大小。 - 修复3.2.1版本压实不清理过期状态的bug
Spark 3.2.1默认没有开启状态专属的压实过滤器,需要通过spark.sql.streaming.stateStore.rocksdb.options传入配置:
配置后压实过程会自动识别已经过了TTL的状态数据,直接物理删除,不会残留墓碑标记占空间。注意必须使用Spark内置的这个过滤器工厂,不要自定义RocksDB原生过滤规则,避免误删正常状态。compaction_filter_factory=org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreCompactionFilterFactory - 源头减少冗余文件生成
调整RocksDB内存刷盘参数,降低临时文件生成量:spark.sql.streaming.stateStore.rocksdb.maxWriteBufferNumber=2(默认值为4,下调后减少内存中驻留的写缓冲数量,减少刷盘生成的临时SST文件)spark.sql.streaming.stateStore.rocksdb.writeBufferSizeMB=32(默认值为64,下调单写缓冲大小,避免单次刷盘生成过大的待压实文件)
- 兜底周期压实配置
设置spark.sql.streaming.stateStore.rocksdb.compactionIntervalNumCommits=50,即每50次批次提交不管状态大小有没有到阈值,都触发一次增量压实,避免小批量更新累积产生大量零散小文件。
存量冗余数据快速清理
如果当前已经存在大量冗余磁盘占用,无需停作业或删除checkpoint目录:
- 临时设置
spark.sql.streaming.stateStore.rocksdb.forceCompactionOnNextCommit=true,下一个批次提交完成后会自动触发一次全量压实,跑完后磁盘占用会回落到和实际状态匹配的量级,一般为几十到几百MB。 - 开启本地临时状态自动清理:设置
spark.sql.streaming.stateStore.rocksdb.cleanupLocalStateOnCheckpoint=true,修复3.2.1版本中checkpoint上传完成后本地状态临时副本不自动删除的bug。
注意:禁止手动删除checkpoint目录下的RocksDB相关文件,会直接导致状态损坏、作业启动失败。所有压实清理操作均为在线执行,不会阻塞流作业正常处理,仅会占用少量CPU与磁盘IO资源,可根据业务峰谷调整强制压实的触发时机。
内容的提问来源于stack exchange,提问作者Learnis
相关产品推荐
相关产品推荐

