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

Flink Sink异常重启时,Source/Process算子状态是否回滚?

针对你提到的 source (kafka) -> process -> sink (db) 数据流场景,结合Checkpoint机制,直接给出明确结论:

核心结论

  • Source、Process算子的状态不会保留本次失败Checkpoint过程中生成的临时快照,作业重启时会回滚到最近一次成功完成的Checkpoint对应的状态。
  • Sink算子异常导致的作业重启不会触发额外的状态回滚,重启逻辑本身就是基于上一个有效Checkpoint恢复状态,而非基于本次失败的中间状态。

详细解释

Flink的Checkpoint遵循「全链路成功才提交」的规则:只有当所有算子都完成当前Checkpoint的快照,且所有快照都持久化到外部存储后,整个Checkpoint才会被标记为有效。

如果Sink算子异常导致本次Checkpoint失败,那么:

  1. 本次Checkpoint过程中,Source、Process算子生成的临时快照会被直接丢弃,不会被作为有效状态存储;
  2. 作业重启时,Flink会自动加载最近一次成功的Checkpoint数据,将Source的Kafka offset、Process算子的计算状态都恢复到那个时刻的状态,相当于本次失败的Checkpoint从未执行过。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 23:35:01