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,提问作者Руслан Цегельников

