Apache Flink应用迁移旧系统监控源至Keyed状态失败求助
Apache Flink Savepoint状态加载失败排查
背景
我正在用Apache Flink开发应用替代旧系统:
- 旧系统:接收多源事件,监控事件是否符合特定条件(如传感器值超阈值),当阈值持续超出指定时长(数天至数周)时上报相关源。
- 新应用:功能一致,消费两个Kafka主题(事件数据、配置数据),通过
TimerService在监控源超出时长限制时向第三个主题写入告警,拓扑为两个源主题 → KeyedCoProcessFunction → 告警Sink主题。
状态迁移方案
要将旧系统的所有监控源迁移至Flink的Keyed状态,已导出包含初始化信息的文件,为此编写了初始化应用:
- 初始化应用拓扑:配置主题源 + 读取旧监控源文件的集合源 + 与正式应用状态定义一致的
KeyedCoProcessFunction - 计划步骤:
- 启动初始化应用,完全消费配置主题并保存状态,同时处理旧监控源列表并写入
KeyedCoProcessFunction的状态; - 手动停止初始化应用并生成Savepoint;
- 使用该Savepoint启动正式Flink应用,加载完整状态。
- 启动初始化应用,完全消费配置主题并保存状态,同时处理旧监控源列表并写入
问题现象
已为配置源和KeyedCoProcessFunction设置相同UID,但启动正式应用时Savepoint状态未正确加载,相关日志如下:
Skipping empty savepoint state for operator 1db5b24325c47f4692485c6f3204cb3b. // new event source Reset the checkpoint ID of job d92e6e970306ea6465b0f253beddebbe to 6. Restoring job d92e6e970306ea6465b0f253beddebbe from Savepoint 5 @ 0 for d92e6e970306ea6465b0f253beddebbe located at file:/tmp/savepoint-05ed46-60bd264d64df. No master state to restore Resetting coordinator to checkpoint. Closing SourceCoordinator for source Source: Configs. Source coordinator for source Source: Configs closed. Restoring SplitEnumerator of source Source: Configs from checkpoint. Resetting coordinator to checkpoint. Closing SourceCoordinator for source Source: Events. Source coordinator for source Source: Events closed. Using failover strategy org.apache.flink.runtime.executiongraph.failover.flip1.RestartPipelinedRegionFailoverStrategy@1bd0b5ba for Monitor_v1.0 (d92e6e970306ea6465b0f253beddebbe). Starting execution of job 'Monitor_v1.0' (d92e6e970306ea6465b0f253beddebbe) under job master id 00000000000000000000000000000000.
求助
有没有人遇到过类似问题?或者有什么解决思路?
内容的提问来源于stack exchange,提问作者mauam
相关产品推荐
相关产品推荐

