Apache Flink Avro FileSink长时间处于in-progress状态问题排查
问题根因
- 代码执行顺序错误:
env.execute()是Flink作业的触发阻塞方法,会等待作业全量执行完成后才会返回执行后续语句,你定义的Thread.sleep(300000)位于env.execute()之后,作业运行期间完全不会执行,有限流作业跑完所有测试数据后直接退出,没有预留checkpoint触发的时间窗口。 - 有限数据源的checkpoint触发逻辑问题:你使用
fromCollection生成的是仅包含6条数据的有限流,数据量极小,作业会在你配置的5s首次checkpoint触发之前就执行完毕,而你选用的OnCheckpointRollingPolicy滚动策略要求必须完成一次成功的checkpoint才会将in-progress状态的临时文件转换为正式可读写的avro文件,没有成功的checkpoint触发自然不会执行文件转正逻辑。 - 若你使用的是1.14以下版本的Flink,默认不会在有限流作业结束前触发最终checkpoint,即使调整时间也可能出现作业跑完没触发最终checkpoint的情况。
修复方案
- 调整代码逻辑,将有限数据源改为无限流,确保作业不会提前退出,预留checkpoint触发时间,示例代码:
// 替换原来的fromCollection逻辑 DataStream<User> source = env.addSource(new SourceFunction<User>() { private volatile boolean isRunning = true; @Override public void run(SourceContext<User> ctx) throws Exception { while (isRunning) { ctx.collect(getUser()); Thread.sleep(1000); } } @Override public void cancel() { isRunning = false; } });
- 如果你需要保留有限流逻辑,开启Flink最终checkpoint配置(适用于Flink 1.14+版本),确保作业结束前触发最后一次checkpoint完成文件提交:
env.getCheckpointConfig().setFinalCheckpointEnabled(true);
- 也可以根据测试场景缩小checkpoint间隔,比如将5s改为100ms,保证有限流跑完前能触发至少一次checkpoint:
env.enableCheckpointing(100);
内容的提问来源于stack exchange,提问作者user51
相关产品推荐
相关产品推荐

