CommonLoggingErrorHandler的setAckAfterHandle行为疑问及消息丢失问题
问题描述
配置代码
@Bean public ConcurrentKafkaListenerContainerFactory<String, Object> batchContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Object> containerFactory = new ConcurrentKafkaListenerContainerFactory<>(); containerFactory.setConsumerFactory(batchConsumerFactory()); containerFactory.setBatchListener(true); ContainerProperties containerProperties = containerFactory.getContainerProperties(); containerProperties.setAckMode(ContainerProperties.AckMode.BATCH); CommonLoggingErrorHandler commonLoggingErrorHandler = new CommonLoggingErrorHandler(); commonLoggingErrorHandler.setAckAfterHandle(false); containerFactory.setCommonErrorHandler(commonLoggingErrorHandler); return containerFactory; }
监听器代码
@KafkaListener(id = "customer-topic", containerFactory = "batchContainerFactory", topics = {"#{kafkaApplicationProperties.customers}"}, groupId = "#{kafkaApplicationProperties.groupId}", clientIdPrefix = "#{kafkaApplicationProperties.clientIdPrefix}", concurrency = "#{kafkaApplicationProperties.concurrency}") public void masterDataListener(@Payload List<ConsumerRecord<String, ContractCustomerEvent>> customers) { log.info("Min offset {}, size {}", customers.stream().map(this::partitionAndOffset).sorted().findFirst().orElseThrow(), customers.size()); if (customers.stream().anyMatch(customer -> customerData.containsKey(customer.value().getId()))) { log.info("Has seen"); // 从未被调用,因为没有记录被重放 } customers.forEach(customer -> customerData.put(customer.value().getId(), customer.value())); if (counter.get().getAndIncrement() < 5) { log.info("Boom"); throw new IllegalArgumentException("Boom"); } customers.forEach(customer -> afterEx.put(customer.value().getId(), customer.value())); // afterEx和customerData内容不一致,afterEx包含更少记录,这些记录丢失了 } // 用于打印偏移量的方法 private Object getOffsets(AdminClient adminClient) { return adminClient.listConsumerGroupOffsets(kafkaApplicationProperties.getGroupId()).all() .whenComplete((result, throwable) -> log.info("######## {}", result, throwable)); }
现象与疑问
设置setAckAfterHandle(false)后,预期抛出异常时首批记录会被重放,但实际监听器继续处理不包含这些记录的批次,最终偏移量被提交导致消息丢失。
观察到的偏移量日志:
2023-10-12T12:32:51.376+02:00 INFO 5756 --- [| adminclient-2] a.r.l.k.c.KafkaTopicListenerITTest : ######## {my-group={my-customers-0=OffsetAndMetadata{offset=50000, leaderEpoch=null, metadata=''}, my-customers-1=OffsetAndMetadata{offset=50000, leaderEpoch=null, metadata=''}}}
请问这是CommonLoggingErrorHandler的预期行为吗?如果是,setAckAfterHandle属性的作用是什么?
解答
1. 这是CommonLoggingErrorHandler的预期行为
是的,这完全符合该处理器的设计逻辑。CommonLoggingErrorHandler的核心作用仅为记录异常日志,它本质是一个"无操作"(no-op)处理器,不会对消息进行重试、偏移回退等干预操作,仅完成日志打印后就会让容器继续执行默认流程。
2. setAckAfterHandle的实际作用
setAckAfterHandle(false)的作用是:错误处理器完成异常处理(此处即打印日志)后,不主动触发额外的偏移量提交。但它不会阻止容器本身的偏移提交逻辑。
结合你配置的AckMode.BATCH来看,容器的默认逻辑是:无论监听器是否抛出异常,只要批次处理流程结束(包括异常被错误处理器捕获),就会提交该批次的偏移量。setAckAfterHandle(false)只是避免错误处理器额外提交偏移,但容器本身的批次提交逻辑依然会执行,这就是偏移被提交、消息丢失的原因。
3. 实现异常时消息重放的方案
如果需要在抛出异常时重放消息,需使用具备重试/偏移回退能力的错误处理器,比如:
SeekToCurrentErrorHandler:默认会将消费者偏移回退到当前批次起始位置,触发消息重放,支持配置最大重试次数DefaultErrorHandler:Spring Kafka 2.8+推荐使用的处理器,支持更灵活的重试规则和死信队列配置
示例配置(使用SeekToCurrentErrorHandler):
@Bean public ConcurrentKafkaListenerContainerFactory<String, Object> batchContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Object> containerFactory = new ConcurrentKafkaListenerContainerFactory<>(); containerFactory.setConsumerFactory(batchConsumerFactory()); containerFactory.setBatchListener(true); ContainerProperties containerProperties = containerFactory.getContainerProperties(); containerProperties.setAckMode(ContainerProperties.AckMode.BATCH); // 替换为SeekToCurrentErrorHandler,设置1秒间隔、最多重试3次 SeekToCurrentErrorHandler errorHandler = new SeekToCurrentErrorHandler( new DeadLetterPublishingRecoverer(kafkaTemplate), new FixedBackOff(1000L, 3L) ); containerFactory.setCommonErrorHandler(errorHandler); return containerFactory; }
内容的提问来源于stack exchange,提问作者Balázs Németh
相关产品推荐
相关产品推荐

