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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 17:09:01