Spark Structured Streaming任务运行中途因S3 checkpoint目录delta文件缺失失败
问题情况
用Spark 3.4.4运行结构化流应用,从Kafka读取数据做聚合后写入Cassandra,将S3作为checkpoint目录。应用上午8点启动,中午突然失败,核心问题是状态存储的delta文件s3a://checkpoint/state/0/7/1.delta找不到,导致Stage任务重试4次后中止。
关键错误栈
异常信息如下---->:
org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:333),org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:207)原因是----> + org.apache.spark.SparkException: 作业因阶段失败中止:阶段14969.0中的任务7失败4次,最近一次失败:丢失阶段14969.0中的任务7.3(TID 104482)(192.170.15.202 executor 2):java.lang.IllegalStateException: 读取HDFSStateStoreProvider[id = (op=0,part=7),dir = s3a://checkpoint/state/0/7]的delta文件s3a://checkpoint/state/0/7/1.delta时出错:s3a://checkpoint/state/0/7/1.delta不存在",
由以下原因导致:java.io.FileNotFoundException: 不存在该文件或目录:s3a://checkpoint/state/0/7/1.delta
用到的代码片段
Spark Session 配置
val spark: SparkSession = SparkSession.builder. appName(appName). config("spark.streaming.stopGracefullyOnShutdown", "true"). config("spark.cassandra.connection.ssl.enabled", "true"). config("spark.cassandra.connection.ssl.protocol", "TLS"). config("spark.sql.extensions", "com.datastax.spark.connector.CassandraSparkExtensions"). config("spark.sql.catalog.casscatalog","com.datastax.spark.connector.datasource.CassandraCatalog"). config("spark.sql.shuffle.partitions",8). config("spark.cassandra.connection.ssl.trustStore.path", cassandraDbTrustStorePath). config("spark.cassandra.connection.ssl.trustStore.password", cassandraConfigs.getString("cassandraDbTruststorePwd")). config("spark.cassandra.connection.host", cassandraConfigs.getString("cassandraHost")). config("spark.cassandra.connection.port", cassandraConfigs.getString("cassandraPort")). config("spark.cassandra.auth.username", cassandraConfigs.getString("cassandraUsername")). config("spark.cassandra.auth.password", cassandraConfigs.getString("cassandraPassword")). config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem"). config("spark.hadoop.fs.s3a.endpoint", "s3.ap-southeast-2.amazonaws.com"). config("spark.hadoop.fs.s3a.access.key", AWS_ACCESS_KEY_ID). config("spark.hadoop.fs.s3a.secret.key", AWS_SECRET_ACCESS_KEY). config("spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled", "true"). config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a", "org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory"). config("spark.hadoop.fs.s3a.committer.name", "magic"). config("spark.sql.sources.commitProtocolClass", "org.apache.spark.internal.io.cloud.PathOutputCommitProtocol"). config("spark.sql.parquet.output.committer.class", "org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter"). config("spark.sql.streaming.checkpointFileManagerClass", "org.apache.spark.internal.io.cloud.AbortableStreamBasedCheckpointFileManager"). getOrCreate()
流写入逻辑
val writeQuery = dbWrite.writeStream. queryName(queryName). outputMode("update"). foreach(new CassandraSinkForeach(nameSpace, tableName, spark, country)). option("checkpointLocation", "s3a://checkpoint/folderpath"). trigger(Trigger.ProcessingTime(triggerTime)). start()
解决办法和优化建议
1. 确认文件是否真的丢失
直接登录AWS控制台,进入对应S3桶,检查路径s3a://checkpoint/state/0/7/下的1.delta文件是否存在。如果确实不存在,说明是checkpoint写入过程中出现异常(比如S3网络抖动、写入超时导致文件未完整生成)。
2. 修复当前故障
Spark状态存储完全依赖checkpoint目录的完整性,delta文件丢失后无法直接恢复现有checkpoint,有两种处理方案:
- 方案一:删除checkpoint重跑
删除S3上的整个checkpoint目录s3a://checkpoint/folderpath/,重启应用。此方案会让应用从头消费Kafka数据,适合业务允许重新计算历史数据的场景。 - 方案二:手动尝试恢复(风险高)
查看同目录下的snapshot文件和其他完整delta文件,若存在最近的快照,可尝试删除损坏的delta文件。但该操作需要熟悉Spark状态存储结构,操作不当会导致更严重问题,谨慎使用。
3. 优化配置避免后续故障
针对Spark 3.4.4版本,调整以下配置提升S3 checkpoint的稳定性:
- 开启S3一致性校验并增加重试:
config("spark.hadoop.fs.s3a.consistent", "true") config("spark.hadoop.fs.s3a.consistent.retry.count", "10") config("spark.hadoop.fs.s3a.consistent.retry.interval", "1000ms") - 给状态存储添加重试机制:
config("spark.sql.streaming.stateStore.retryCount", "3") - 替换Magic Committer为Directory Committer:
Magic Committer更适合批量Parquet写入,checkpoint需要强一致性,改用Directory Committer:config("spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled", "false") config("spark.hadoop.fs.s3a.committer.name", "directory") - 调大网络超时参数,防止Executor心跳超时:
config("spark.network.timeout", "300s") config("spark.executor.heartbeatInterval", "60s")
4. 检查自定义Cassandra Sink
自定义的CassandraSinkForeach如果是同步阻塞实现,或缺少异常重试机制,可能导致Spark任务卡顿,间接影响checkpoint写入。建议:
- 将Sink改为异步写入,避免阻塞流处理流程
- 给Sink添加异常捕获和重试逻辑,防止因Cassandra写入失败导致任务反复重试,干扰checkpoint正常写入
内容的提问来源于stack exchange,提问作者Karthik

