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

Spark Structured Streaming任务运行中途因S3 checkpoint目录delta文件缺失失败

Spark 3.4.4 结构化流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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 03:09:53