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

Spring-Kafka多监听器与多死信队列配置问题求助

Spring Boot 2.7.9 Kafka监听器错误处理配置方案

针对你的需求,我们可以通过DefaultErrorHandler(替代废弃的ErrorHandler)结合DeadLetterPublishingRecoverer和自定义日志恢复逻辑来实现。以下是修改后的完整配置代码:

@Configuration
@EnableKafka
public class ConsumerKafkaConfig {

    private final KafkaProperties kafkaProperties;

    public ConsumerKafkaConfig(KafkaProperties kafkaProperties) {
        this.kafkaProperties = kafkaProperties;
    }

    // 新增:创建ProducerFactory和KafkaTemplate,用于死信队列消息发送
    @Bean
    public ProducerFactory<String, Object> producerFactory() {
        return new DefaultKafkaProducerFactory<>(kafkaProperties.buildProducerProperties());
    }

    @Bean
    public KafkaTemplate<String, Object> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }

    @Bean
    public ConsumerFactory<String, Object> firstEventFactory() {
        Map<String, Object> props = kafkaProperties.buildConsumerProperties();
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> firstContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(firstEventFactory());

        // 配置死信队列错误处理器
        DeadLetterPublishingRecoverer dlqRecoverer = new DeadLetterPublishingRecoverer(kafkaTemplate(),
                (record, ex) -> new TopicPartition("my-first-dlq", record.partition()));
        // 重试0次后直接发送死信,可根据需求调整FixedBackOff参数
        DefaultErrorHandler dlqErrorHandler = new DefaultErrorHandler(dlqRecoverer, new FixedBackOff(0L, 0));
        factory.setCommonErrorHandler(dlqErrorHandler);

        return factory;
    }

    @Bean
    public ConsumerFactory<Object, Object> secondConsumerEventFactory() {
        final Map<String, Object> properties = kafkaProperties.buildConsumerProperties();
        properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
        return new DefaultKafkaConsumerFactory<>(properties);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<Object, Object> secondConsumerEventContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<Object, Object> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(secondConsumerEventFactory());

        // 配置仅记录错误的处理器
        RecoveryCallback<Void> loggingCallback = (context) -> {
            ConsumerRecord<?, ?> record = context.getConsumerRecord();
            LoggerFactory.getLogger(ConsumerKafkaConfig.class)
                    .error("消息处理失败,topic: {}, partition: {}, offset: {}",
                            record.topic(), record.partition(), record.offset(),
                            context.getThrowable());
            return null;
        };
        // 重试0次后仅记录日志
        DefaultErrorHandler loggingErrorHandler = new DefaultErrorHandler(loggingCallback, new FixedBackOff(0L, 0));
        factory.setCommonErrorHandler(loggingErrorHandler);

        return factory;
    }

    @Bean
    public ConsumerFactory<String, Object> thirdConsumerEventFactory() {
        final Map<String, Object> properties = kafkaProperties.buildConsumerProperties();
        properties.put(AbstractKafkaSchemaSerDeConfig.VALUE_SUBJECT_NAME_STRATEGY, TopicRecordNameStrategy.class.getName());
        properties.put(AbstractKafkaSchemaSerDeConfig.AUTO_REGISTER_SCHEMAS, false);
        properties.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true);
        properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
        return new DefaultKafkaConsumerFactory<>(properties);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> thirdConsumerEventContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(thirdConsumerEventFactory());
        return factory;
    }

}

关键配置说明

  • firstContainerFactory死信处理:
    1. 通过DeadLetterPublishingRecoverer指定死信队列主题my-first-dlq,保持原消息的分区一致性
    2. DefaultErrorHandler设置重试次数为0(直接进入死信),若需要重试可调整FixedBackOff的参数(第一个参数是重试间隔,第二个是重试次数)
  • secondContainerFactory日志处理:
    1. 自定义RecoveryCallback实现错误日志的详细记录,包含消息的topic、partition、offset等关键信息
    2. 同样设置重试次数为0,确保仅记录错误不重复消费

内容的提问来源于stack exchange,提问作者Pierre Jones

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 18:24:59