Flink 1.15.2升级至1.18时SQL作业Savepoint恢复失败咨询
问题解答
1. 是否属于需要上报的Bug?
不属于Bug。这个算子顺序调整是Flink社区为了优化性能(减少ChangelogNormalize算子处理的数据量,提前过滤字段)主动做的设计变更,属于版本间的正常优化,并非代码缺陷。
2. 保留状态完成迁移的可行方案
- 临时跳过无状态算子的状态检查:如果ChangelogNormalize算子本身没有持久化状态(这类算子通常仅做格式转换,无状态需要保留),启动1.18版本作业时添加
--allowNonRestoredState参数,即可跳过无法映射的算子状态检查,正常恢复作业且不影响业务。 - 先全量字段恢复再调整查询:先使用
SELECT *的查询语句从Savepoint恢复作业到1.18版本,此时算子顺序匹配,状态可正常恢复;待作业稳定后,再将查询修改为指定字段的版本,Flink会自动生成新执行计划并调整算子链,无需重新依赖旧Savepoint。 - 禁用对应优化规则:在作业启动配置中添加
table.optimizer.exclude-rules=CalcBeforeChangelogNormalizeRule,禁用导致算子顺序调换的优化规则,让执行计划回到1.15版本的算子顺序,即可直接从旧Savepoint恢复作业。
内容的提问来源于stack exchange,提问作者C.S.
相关产品推荐
相关产品推荐

