Spring Kafka 3.0.6阻塞重试实现代码检查:是否存在遗漏?
Spring Kafka 3.0.6 阻塞重试与死信队列代码检查及修正
当前代码目标是实现:处理消息遇到RetryableException时重试2次,重试耗尽后将消息发送至指定固定分区的死信队列。现有代码存在以下遗漏和问题:
- 异常范围未限定:
DefaultErrorHandler默认会对除FatalThrowable外的所有异常重试,不符合仅对RetryableException重试的需求。 - 重试次数配置错误:
FixedBackOff(0L, 2L)的maxAttempts参数代表总尝试次数(含第一次执行),设置为2L仅会触发1次重试,无法满足“重试2次”的要求。 - Header参数类型不匹配:消费方法中
partition、offset、timestamp定义为String,但Kafka实际传递的是Integer(partition)、Long(offset、timestamp),会触发类型转换异常。 - 死信队列Topic硬编码:死信队列名称写死为
"topicName" + ".DLT",与消费的demoTopicName不对应,易引发错误。 - 非RetryableException处理遗漏:捕获非重试异常后仅打印日志,未提交偏移量,会导致消息重复消费。
修正后的消费方法代码
@KafkaListener( autoStartup = "false", containerFactory = "concurrentKafkaListenerContainerFactory", id = "demoConsumer", groupId = "demoConsumerGroup", topics = "demoTopicName" ) public void consumeMessage( @Header(KafkaHeaders.RECEIVED_PARTITION) Integer partition, @Header(KafkaHeaders.OFFSET) Long offset, @Header(KafkaHeaders.RECEIVED_TIMESTAMP) Long timestamp, ConsumerRecord<String, String> consumerRecord, Acknowledgment acknowledgment) { try { log.info("Processing message from partition: {}, offset: {}", partition, offset); // 业务处理逻辑 acknowledgment.acknowledge(); } catch (Exception ex) { if (ex instanceof RetryableException) { throw ex; // 抛出异常触发重试逻辑 } log.error("Non-retryable exception occurred", ex); acknowledgment.acknowledge(); // 提交偏移量,避免重复消费 } }
修正后的配置类代码
@EnableKafka @Configuration @RequiredArgsConstructor public class KafkaConfig { private final KafkaTemplate<String, String> kafkaTemplate; @Bean public ConcurrentKafkaListenerContainerFactory<String, String> concurrentKafkaListenerContainerFactory(ConsumerFactory<String, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setConcurrency(1); factory.getContainerProperties().setStopImmediate(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL); factory.setCommonErrorHandler(defaultErrorHandler()); return factory; } @Bean public CommonErrorHandler defaultErrorHandler() { // 总尝试次数3次 = 1次初始执行 + 2次重试 DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer(), new FixedBackOff(0L, 3L)); // 仅对RetryableException触发重试 errorHandler.addRetryableExceptions(RetryableException.class); // 排除所有其他异常,不进行重试 errorHandler.addNotRetryableExceptions(Exception.class); return errorHandler; } @Bean public DeadLetterPublishingRecoverer recoverer() { final BiFunction<ConsumerRecord<?, ?>, Exception, TopicPartition> CUSTOMIZE_DESTINATION_RESOLVER = (cr, e) -> // 根据原Topic动态生成死信队列,固定发送至分区0 new TopicPartition(cr.topic() + ".DLT", 0); return new DeadLetterPublishingRecoverer(kafkaTemplate, CUSTOMIZE_DESTINATION_RESOLVER); } }
关键修正说明
- 异常精确控制:通过
addRetryableExceptions和addNotRetryableExceptions限定仅RetryableException触发重试。 - 重试次数校准:将
FixedBackOff的maxAttempts设为3L,确保初始执行失败后进行2次重试。 - 参数类型修正:将Header参数改为Kafka实际传递的类型,避免类型转换错误。
- 死信队列动态生成:根据原消息Topic自动生成死信队列名称,避免硬编码错误。
- 偏移量提交:非重试异常处理时提交偏移量,防止消息重复消费。
- Bean管理:将错误处理器和死信恢复器声明为Spring Bean,确保上下文正确管理。
内容的提问来源于stack exchange,提问作者Sk Monjurul Haque
相关产品推荐
相关产品推荐

