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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 10:36:04