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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 18:22:08