spring-kafka中RetryingBatchErrorHandler恢复器如何提交偏移量
spring-kafka 2.3.8 恢复器偏移量提交实现方案
1. 默认自动提交(无自定义提交需求优先用这个)
RetryingBatchErrorHandler本身内置了恢复后的偏移量自动提交逻辑,不需要你在恢复器里手写提交代码,满足以下两个配置即可生效:
- 容器的
AckMode配置为AckMode.BATCH或AckMode.RECORD(spring-kafka默认就是AckMode.BATCH) - 不要关闭错误处理器的自动提交开关(2.3.8版本该开关默认开启)
你现有代码不需要修改恢复器内部逻辑,重试耗尽调用恢复器记录完日志后,框架会自动提交对应偏移量。
2. 自定义手动提交方案(有特殊提交需求时使用)
如果需要在恢复器里手动控制提交逻辑(比如携带自定义偏移量元数据、自定义提交时机等),可以通过KafkaUtils获取当前线程绑定的消费者实例,手动执行提交,修改后代码如下:
import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import org.springframework.kafka.support.KafkaUtils; import org.springframework.kafka.listener.ContainerCustomizer; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.RetryingBatchErrorHandler; import org.springframework.kafka.listener.ConsumerRecordRecoverer; import org.springframework.util.backoff.FixedBackOff; import java.util.Collections; import lombok.extern.slf4j.Slf4j; @Slf4j public class Customizer implements ContainerCustomizer { private static ConsumerRecordRecoverer createConsumerRecordRecoverer() { return (consumerRecord, e) -> { log.info("Number of attempts exhausted. partition: {}, offset: {}", consumerRecord.partition(), consumerRecord.offset()); // 手动提交偏移量逻辑 Consumer<?, ?> consumer = KafkaUtils.getConsumerResource(); if (consumer != null) { TopicPartition tp = new TopicPartition(consumerRecord.topic(), consumerRecord.partition()); // *注意:kafka提交的是下一条要消费的偏移量,所以要+1* OffsetAndMetadata offsetMeta = new OffsetAndMetadata(consumerRecord.offset() + 1); // 同步提交,需要异步的话换成commitAsync即可 consumer.commitSync(Collections.singletonMap(tp, offsetMeta)); log.info("偏移量手动提交完成,topic:{}, partition:{}, 提交offset:{}", consumerRecord.topic(), consumerRecord.partition(), consumerRecord.offset() + 1); } }; } @Override public void configure(AbstractMessageListenerContainer container) { // 关闭自动提交避免重复提交,手动提交时配置 container.getContainerProperties().setEnableAutoCommit(false); container.setBatchErrorHandler(new RetryingBatchErrorHandler(new FixedBackOff(5000L, 3L), createConsumerRecordRecoverer())); } }
注意事项
- 手动提交时必须关闭消费者的
enableAutoCommit配置,避免自动提交和手动提交冲突 - 如果是批处理场景下有多条失败记录,恢复器会逐次调用,你可以批量收集偏移量后统一提交,减少提交开销
- 同步提交会阻塞线程直到提交结果返回,异步提交不会阻塞但可能丢失提交失败异常,根据业务可靠性要求选择即可
内容的提问来源于stack exchange,提问作者Scott Stark
相关产品推荐
相关产品推荐

