Spark流作业changelog文件不存在报错,请求社区协助排查解决
Spark流作业Kubernetes部署后无法读取S3中RocksDB状态文件的解决方案
报错核心信息
Caused by: java.lang.IllegalStateException: Error reading streaming state file of s3a://table-onlineofflineinternet-store/checkpointing/abcd/temp_table/state/0/6/563042.changelog does not exist. If the stream job is restarted with a new or updated state operation, please create a new checkpoint location or clear the existing checkpoint location. at org.apache.spark.sql.execution.streaming.state.StateStoreChangelogReader.liftedTree1$1(StateStoreChangelog.scala:136) at org.apache.spark.sql.execution.streaming.state.StateStoreChangelogReader.<init>(StateStoreChangelog.scala:129) at org.apache.spark.sql.execution.streaming.state.RocksDBFileManager.getChangelogReader(RocksDBFileManager.scala:159) at org.apache.spark.sql.execution.streaming.state.RocksDB.$anonfun$replayChangelog$1(RocksDB.scala:196) at scala.runtime.java8.JFunction1$mcVJ$sp.apply(JFunction1$mcVJ$sp.java:23) at scala.collection.immutable.NumericRange.foreach(NumericRange.scala:75) at org.apache.spark.sql.execution.streaming.state.RocksDB.replayChangelog(RocksDB.scala:193) at org.apache.spark.sql.execution.streaming.state.RocksDB.load(RocksDB.scala:166) at org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider.getReadStore(RocksDBStateStoreProvider.scala:200) at org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider.getReadStore(RocksDBStateStoreProvider.scala:30) at org.apache.spark.sql.execution.streaming.state.StateStore$.getReadOnly(StateStore.scala:492) at org.apache.spark.sql.execution.streaming.state.ReadStateStoreRDD.compute(StateStoreRDD.scala:92) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364)
解决步骤
1. 确认状态逻辑变更并重置Checkpoint
如果最近修改了流作业中涉及状态的逻辑(比如修改聚合字段、调整Watermark、新增/删除状态相关算子),必须执行以下操作之一:
- 修改代码中
checkpointLocation参数,指定S3上的全新路径 - 清空原Checkpoint路径的所有内容(例如使用命令:
aws s3 rm --recursive s3a://table-onlineofflineinternet-store/checkpointing/abcd/temp_table/)
Spark的状态存储与作业逻辑强绑定,逻辑变更后旧Checkpoint无法兼容,强行复用必然导致状态文件找不到。
2. 排查S3存储一致性问题
Kubernetes环境下可能因S3的最终一致性特性,或客户端配置问题导致文件延迟可见:
- 在Spark配置中添加或修改:
开启S3强一致性并增加重试次数,让作业等待文件同步spark.hadoop.fs.s3a.consistent=true spark.hadoop.fs.s3a.retry.max=10 - 手动登录S3控制台确认
563042.changelog文件是否存在。若确实不存在,说明作业异常退出时未正确写入状态,只能清空Checkpoint重新启动
3. 验证Kubernetes权限配置
检查K8s中作业运行的服务账号权限:
- 确认服务账号绑定的IAM角色拥有S3目标路径的
GetObject、ListBucket权限 - 排查是否因K8s集群权限变更导致作业无法读取S3文件
4. 检查Spark版本兼容性
确保Kubernetes环境使用的Spark版本与之前作业正常运行的版本完全一致:
- 不同Spark版本的RocksDB状态存储格式可能存在差异,版本不匹配会导致无法读取状态文件
- 若必须升级Spark版本,需重新初始化Checkpoint,不能复用旧路径
内容的提问来源于stack exchange,提问作者Aditya Verma
相关产品推荐
相关产品推荐

