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死信处理:
- 通过
DeadLetterPublishingRecoverer指定死信队列主题my-first-dlq,保持原消息的分区一致性 DefaultErrorHandler设置重试次数为0(直接进入死信),若需要重试可调整FixedBackOff的参数(第一个参数是重试间隔,第二个是重试次数)
- 通过
- secondContainerFactory日志处理:
- 自定义
RecoveryCallback实现错误日志的详细记录,包含消息的topic、partition、offset等关键信息 - 同样设置重试次数为0,确保仅记录错误不重复消费
- 自定义
内容的提问来源于stack exchange,提问作者Pierre Jones
相关产品推荐
相关产品推荐

