Flink流作业无法从Checkpoint恢复且报S3分片上传无效如何解决?
根因定位
该问题是Flink StreamingFileSink对接对象存储S3时的典型状态一致性问题,触发链路如下:
- StreamingFileSink写入S3时依赖分段上传(MultiPart Upload)机制:数据先以临时分段形式上传到S3侧,仅当Checkpoint全量成功完成后,才会触发Pending文件的Commit动作,合并分段为正式Part文件,同时注销本次分段上传任务。
- 作业最后一次Checkpoint失败,报错
Checkpoint Coordinator is suspending,说明本次Checkpoint执行流程被强制中断,此时部分Subtask的分段上传已经执行完成,但还未触发Commit动作,对应待提交文件的元信息已经被写入到本次成功完成的Checkpoint状态中。 - 绝大多数S3实现(包括AWS S3、兼容S3协议的国产对象存储)都默认配置了未完成分段上传的自动清理规则,默认清理周期通常为7天,刚好匹配作业稳定运行数日后才出问题的时间特征:未完成的分段上传任务被S3自动清理,但待提交的文件元信息仍然保留在Flink的Checkpoint状态中。
- 作业恢复时,Flink会尝试提交状态中记录的所有Pending文件,但对应的分段上传任务已经被S3销毁,因此抛出
MultiPart Upload is not valid、InvalidPart类异常,最终导致恢复失败。
恢复方案
紧急恢复(可容忍极少量数据重复)
直接选择早于分段上传清理时间的历史Checkpoint触发恢复即可。
零数据丢失恢复
- 先备份当前故障的Checkpoint文件
- 使用Flink State Processor API读取故障Checkpoint,过滤掉状态中记录的无效待提交文件元信息,生成新的合法Checkpoint
- 使用新生成的Checkpoint启动作业
长期规避措施
- 调整S3 Bucket生命周期规则,将未完成分段上传的自动清理时长设置为大于Flink作业Checkpoint保留时长与最大故障恢复窗口的总和
- 开启StreamingFileSink的过期Inprogress文件自动清理配置,作业启动时自动校验并清理S3侧的无效临时文件
- 升级到Flink 1.14及以上版本,该版本对S3分段上传恢复的异常场景做了兼容性优化,可配置跳过无效待提交文件直接恢复作业
内容的提问来源于stack exchange,提问作者Yuriy Kwak
相关产品推荐
相关产品推荐

