Spring Kafka批量错误处理器实现示例咨询:跳过poison pill无效记录
Spring Kafka 批量错误处理器实现跳过毒丸消息方案
前置要求
- Spring Kafka 版本 >= 2.3(2.8+推荐使用
DefaultErrorHandler,2.3~2.7版本使用RetryingBatchErrorHandler) - 已开启批量消费配置:
spring.kafka.listener.type=batch,或监听容器工厂设置setBatchListener(true)
核心实现
1. 配置错误处理器
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.listener.CommonErrorHandler; import org.springframework.kafka.listener.ConsumerRecordRecoverer; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.util.backoff.FixedBackOff; @Configuration public class KafkaBatchErrorConfig { @Bean public CommonErrorHandler skipPoisonPillErrorHandler() { // FixedBackOff参数:重试间隔0ms,最大重试次数0,即遇到异常直接跳过不重试 DefaultErrorHandler errorHandler = new DefaultErrorHandler(new FixedBackOff(0, 0)); // 配置不需要重试的异常,可按需添加自定义的反序列化、格式校验异常类 errorHandler.addNotRetryableExceptions( org.springframework.kafka.support.serializer.DeserializationException.class, IllegalArgumentException.class ); // 自动提交已恢复的异常记录的offset errorHandler.setCommitRecovered(true); // 自定义异常恢复逻辑,可用于日志打印、告警、毒丸消息落库等 errorHandler.setRecoveryCallback(new ConsumerRecordRecoverer() { @Override public void accept(ConsumerRecord<?, ?> record, Exception e) { System.out.printf("跳过毒丸消息:topic=%s, partition=%d, offset=%d, 异常原因=%s%n", record.topic(), record.partition(), record.offset(), e.getMessage()); } }); return errorHandler; } }
2.7及更早版本替换
DefaultErrorHandler为RetryingBatchErrorHandler即可,配置逻辑完全一致。
2. 绑定错误处理器到批量监听容器工厂
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.context.annotation.Bean; @Bean public ConcurrentKafkaListenerContainerFactory<String, Object> batchListenerContainerFactory( ConsumerFactory<String, Object> consumerFactory, CommonErrorHandler skipPoisonPillErrorHandler ) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 开启批量消费 factory.setBatchListener(true); // 绑定自定义错误处理器 factory.setCommonErrorHandler(skipPoisonPillErrorHandler); return factory; }
3. 批量消费方法示例
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import java.util.List; @Component public class MyBatchConsumer { @KafkaListener(topics = "你的业务topic名称", containerFactory = "batchListenerContainerFactory") public void consumeBatch(List<ConsumerRecord<String, Object>> records) { // 业务逻辑无需额外捕获单条记录异常 records.forEach(record -> { // 你的消息校验、业务处理逻辑 Object message = record.value(); if (message == null) { // 抛出异常后会被错误处理器自动捕获,跳过当前记录 throw new IllegalArgumentException("消息内容为空,判定为毒丸消息"); } // 正常业务处理代码 }); } }
工作逻辑说明
- 批量消费过程中任意一条记录抛出异常时,错误处理器会自动定位异常记录,执行自定义恢复逻辑,提交该记录offset后继续处理批次内剩余记录
- 不会因为单条毒丸消息阻塞整个消费队列,也不需要重新消费整个批次的正常消息
- 如需对部分业务异常做重试,只需调整
FixedBackOff的重试间隔、重试次数参数,或通过addRetryableExceptions方法添加可重试异常类即可
内容的提问来源于stack exchange,提问作者MR_K
相关产品推荐
相关产品推荐

