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

Flink 1.16停止作业时Kafka偏移量提交异常致数据丢失求助

问题分析与解决方案

一、AT_LEAST_ONCE策略下高负载停作业丢数据的1.16版本修复方案

该问题是Flink 1.16-1.18版本存在的偏移量提交与作业停止逻辑不一致的bug,1.19.1通过修复停止作业时的偏移量提交逻辑解决了此问题。针对1.16.x版本,可通过以下方式修复:

1. 自定义作业停止时的偏移量提交逻辑

实现CheckpointListener接口,仅在作业停止时提交已完成处理的偏移量:

public class OffsetCommitControl implements CheckpointListener {
    private final KafkaSink<?, ?, ?> kafkaSink;
    private volatile long lastCompletedCheckpointOffset;

    public OffsetCommitControl(KafkaSink<?, ?, ?> kafkaSink) {
        this.kafkaSink = kafkaSink;
    }

    @Override
    public void notifyCheckpointComplete(long checkpointId) throws Exception {
        // 记录最后完成的checkpoint对应的偏移量
        lastCompletedCheckpointOffset = getOffsetFromCheckpoint(checkpointId);
    }

    @Override
    public void notifyCheckpointAborted(long checkpointId) throws Exception {
        // 忽略失败的checkpoint
    }

    // 作业停止前调用此方法,提交已确认的偏移量
    public void commitValidOffsetOnStop() {
        kafkaSink.commitOffsets(lastCompletedCheckpointOffset);
    }
}

再通过JobListener监听作业停止事件,触发自定义偏移量提交:

env.registerJobListener(new JobListener() {
    @Override
    public void onJobSubmitted(JobClient jobClient, Throwable throwable) {}

    @Override
    public void onJobExecuted(JobExecutionResult jobExecutionResult, Throwable throwable) {
        offsetCommitControl.commitValidOffsetOnStop();
    }
});

2. 禁用自动偏移量提交,绑定checkpoint提交

修改Kafka消费者配置,关闭自动提交,强制仅在checkpoint完成后提交偏移量:

# 关闭Kafka自动偏移量提交
enable.auto.commit=false

同时配置Flink作业的checkpoint参数:

KafkaSource<String> source = KafkaSource.<String>builder()
        .setBootstrapServers("your-broker-address")
        .setTopics("your-topic")
        .setGroupId("your-group-id")
        .setStartingOffsets(OffsetsInitializer.latest())
        .setValueOnlyDeserializer(new SimpleStringSchema())
        .setProperty("enable.auto.commit", "false")
        .build();

// 配置checkpoint触发偏移量提交
env.enableCheckpointing(3000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);
env.getCheckpointConfig().setExternalizedCheckpointCleanup(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

此方式确保只有完成checkpoint的偏移量才会被提交,避免停止作业时提交未处理数据的偏移量。

3. 升级到1.16.x的补丁版本

查看Flink 1.16官方补丁说明,若存在包含该bug修复的backport小版本(如1.16.4及以后),直接升级是最简洁的解决方案。


二、EXACTLY_ONCE策略下的疑问解答

1. 慢处理下为何仍出现类似问题?

EXACTLY_ONCE模式依赖事务型Kafka Sink保证一致性,偏移量提交与事务提交绑定。出现checkpoint完成但偏移量未提交的情况,原因包括:

  • 处理速度放慢时,部分算子的checkpoint屏障延迟,但核心状态已持久化,checkpoint仍标记为COMPLETED,不过事务未满足提交条件(如未收到下游所有算子的屏障确认)。
  • 事务超时配置(transaction.timeout.ms、kafka.producer.transaction.timeout.ms)过长,导致事务未及时提交,进而偏移量也未提交。

2. 为何每3-4次checkpoint才提交一次偏移量?

这是因为EXACTLY_ONCE模式下,事务提交是基于checkpoint批次触发,而非每次checkpoint都提交。默认情况下,Flink会累积多个checkpoint的事务,直到满足以下条件之一才提交:

  • 达到kafka.producer.batch.size配置的批量大小
  • 达到linger.ms配置的等待时间
  • 作业停止或累积的事务数量达到阈值
    若checkpoint间隔远小于事务超时时间,就会出现多次checkpoint才提交一次偏移量的情况。

3. checkpoint是否在屏障到达末端前就结束?

不会。Flink标记checkpoint为COMPLETED的前提是所有算子完成该checkpoint的状态快照,且屏障已到达作业Sink端。但EXACTLY_ONCE模式下,Sink端事务不会在checkpoint完成后立即提交,而是等待下一个checkpoint屏障到达或满足事务提交条件时才提交——这是为了保证端到端一致性,避免部分提交导致的数据不一致。

内容的提问来源于stack exchange,提问作者Руслан Цегельников

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:53:15