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

Spark Structured Streaming使用RocksDB状态存储时ID不匹配异常求助

Spark Structured Streaming流-流Join中RocksDB唯一ID不匹配异常排查

问题场景

  • 技术栈:Java 17、Spark 3.5.0/3.4.0本地模式、Kafka数据源、RocksDB作为状态存储后端,已通过spark.sql.streaming.stateStore.providerClass配置org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider
  • 业务场景:多对一(many to one)关系的流-流关联
  • 问题现象:运行多批次后频繁触发RocksDB唯一ID不匹配异常,多次清理checkpoint目录后问题仍复现;切换至Spark 3.4.0后异常依旧;checkpoint目录位于Windows本地临时路径,未发现源数据损坏

异常信息

org.rocksdb.RocksDBException: Mismatch in unique ID on table file 27. Expected: {7568779327299478048,6781697239205038417} Actual: {7568779327299478056,2042577083335893353} in file C:\Windows\Temp\spark-fed38dd1-aae2-45ab-ad2f-8e0f02ea223b\StateStoreId(opId=0,partId=1,name=left-keyWithIndexToValue)-a692f8b0-e03c-4475-8925-0d589f67d628\workingDir-c71dca44-dcbc-4057-91d2-afce3d7cb7b4/MANIFEST-000005
    at org.rocksdb.RocksDB.open(Native Method) ~[rocksdbjni-8.3.2.jar!/:?]
    at org.rocksdb.RocksDB.open(RocksDB.java:249) ~[rocksdbjni-8.3.2.jar!/:?]
    at org.apache.spark.sql.execution.streaming.state.RocksDB.openDB(RocksDB.scala:584) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.state.RocksDB.load(RocksDB.scala:154) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider.getStore(RocksDBStateStoreProvider.scala:194) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.state.StateStore$.get(StateStore.scala:507) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.state.SymmetricHashJoinStateManager$StateStoreHandler.getStateStore(SymmetricHashJoinStateManager.scala:417) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.state.SymmetricHashJoinStateManager$KeyWithIndexToValueStore.<init>(SymmetricHashJoinStateManager.scala:600) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.state.SymmetricHashJoinStateManager.<init>(SymmetricHashJoinStateManager.scala:386) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.StreamingSymmetricHashJoinExec$OneSideHashJoiner.<init>(StreamingSymmetricHashJoinExec.scala:529) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.StreamingSymmetricHashJoinExec.processPartitions(StreamingSymmetricHashJoinExec.scala:276) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.StreamingSymmetricHashJoinExec.$anonfun$doExecute$1(StreamingSymmetricHashJoinExec.scala:241) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.StreamingSymmetricHashJoinExec.$anonfun$doExecute$1$adapted(StreamingSymmetricHashJoinExec.scala:241) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.streaming.StreamingSymmetricHashJoinHelper$StateStoreAwareZipPartitionsRDD.compute(StreamingSymmetricHashJoinHelper.scala:295) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.SQLExecutionRDD.compute(SQLExecutionRDD.scala:55) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.sql.execution.SQLExecutionRDD.compute(SQLExecutionRDD.scala:55) ~[spark-sql_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.scheduler.Task.run(Task.scala:141) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64) ~[spark-common-utils_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61) ~[spark-common-utils_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94) ~[spark-core_2.12-3.5.0.jar!/:3.5.0]
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623) [spark-core_2.12-3.5.0.jar!/:3.5.0]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) [?:?]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) [?:?]

排查与解决方案

1. 规避Windows文件系统限制

  • 更换checkpoint目录:不要使用系统临时目录(如C:\Windows\Temp),改用自定义非系统盘目录,比如D:\spark-stream-checkpoints,避免系统自动清理或权限限制导致文件损坏
  • 关闭Windows 8.3文件名格式:以管理员身份执行命令fsutil 8dot3name set C: 1,防止短文件名映射引发的文件识别异常
  • 确保权限:给Spark进程分配checkpoint目录的完全读写权限,避免写入中断导致状态文件不完整

2. 优化RocksDB状态存储配置

  • 调整写入缓冲:设置spark.sql.streaming.stateStore.rocksdb.write.buffer.size=67108864(64MB),减少磁盘flush频率,降低文件损坏概率
  • 启用校验:设置spark.sql.streaming.stateStore.rocksdb.enableChecksum=true,让RocksDB自动校验文件完整性,提前拦截损坏文件
  • 配置状态过期:针对多对一场景,给两个流设置合理的watermark,自动清理过期状态,减少状态文件体积:
    Dataset<Row> streamAWithWatermark = streamA.withWatermark("event_time", "1 hour");
    Dataset<Row> streamBWithWatermark = streamB.withWatermark("event_time", "1 hour");
    Dataset<Row> joinedStream = streamAWithWatermark.join(
        streamBWithWatermark,
        functions.col("streamA.id").equalTo(functions.col("streamB.id")),
        JoinType.INNER
    );
    

3. 检查流Join的状态管理逻辑

  • 统一字段类型:确保两个流的关联字段类型完全一致,避免隐式类型转换导致键的哈希值不一致,引发状态存储异常
  • 切换状态存储验证:临时注释掉spark.sql.streaming.stateStore.providerClass配置,使用Spark默认的内存+磁盘状态存储,验证是否为RocksDB特定问题
  • 避免异常退出:确保Spark任务正常停止,避免强制kill进程导致状态文件未正确关闭和持久化

4. 升级RocksDB JNI版本

Spark 3.5.0默认使用rocksdbjni-8.3.2,该版本在Windows环境下可能存在兼容性问题,手动替换为8.5.0或更高稳定版本,确保与Spark版本兼容

内容的提问来源于stack exchange,提问作者Hadi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 05:15:55