如何使用Spring Reactive Kafka实现重试与恢复逻辑
Reactor Kafka 重试与恢复实现方案
以下是基于 reactor-kafka 实现重试+死信恢复逻辑的可运行示例,适配 Spring Reactive 技术栈:
核心实现逻辑
- 区分可重试异常(如网络波动、下游服务临时不可用)和不可重试异常(如消息格式非法),仅对可重试异常触发重试
- 配置固定重试次数+退避策略,避免无意义的频繁重试
- 重试耗尽或遇到不可重试异常时,消息写入死信队列做兜底恢复,避免消息丢失
完整示例代码
1. 配置类
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.core.ReactiveKafkaConsumerTemplate; import org.springframework.kafka.core.ReactiveKafkaProducerTemplate; import reactor.kafka.receiver.ReceiverOptions; import reactor.kafka.sender.SenderOptions; import java.util.Collections; import java.util.HashMap; import java.util.Map; @Configuration public class ReactiveKafkaConfig { @Bean public ReactiveKafkaConsumerTemplate<String, String> reactiveKafkaConsumerTemplate() { Map<String, Object> consumerProps = new HashMap<>(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "business-consumer-group"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 关闭自动提交,由框架控制偏移量提交时机 consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); ReceiverOptions<String, String> receiverOptions = ReceiverOptions.<String, String>create(consumerProps) .subscription(Collections.singleton("business-topic")); return new ReactiveKafkaConsumerTemplate<>(receiverOptions); } @Bean public ReactiveKafkaProducerTemplate<String, String> dlqProducerTemplate() { Map<String, Object> producerProps = new HashMap<>(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new ReactiveKafkaProducerTemplate<>(SenderOptions.create(producerProps)); } }
2. 消费与重试逻辑
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.core.ReactiveKafkaConsumerTemplate; import org.springframework.kafka.core.ReactiveKafkaProducerTemplate; import org.springframework.stereotype.Service; import reactor.core.publisher.Mono; import reactor.util.retry.Retry; import javax.annotation.PostConstruct; import java.time.Duration; @Service public class KafkaRetryConsumerService { private static final Logger log = LoggerFactory.getLogger(KafkaRetryConsumerService.class); // 最大重试次数 private static final int MAX_RETRY_COUNT = 3; // 重试退避间隔 private static final Duration RETRY_INTERVAL = Duration.ofSeconds(2); // 死信队列主题名 private static final String DLQ_TOPIC = "business-topic-dlq"; private final ReactiveKafkaConsumerTemplate<String, String> consumerTemplate; private final ReactiveKafkaProducerTemplate<String, String> dlqProducerTemplate; public KafkaRetryConsumerService(ReactiveKafkaConsumerTemplate<String, String> consumerTemplate, ReactiveKafkaProducerTemplate<String, String> dlqProducerTemplate) { this.consumerTemplate = consumerTemplate; this.dlqProducerTemplate = dlqProducerTemplate; } @PostConstruct public void startConsume() { consumerTemplate.receiveAutoAck() .flatMap(record -> processMessage(record) // 配置重试规则 .retryWhen(Retry.backoff(MAX_RETRY_COUNT, RETRY_INTERVAL) // 仅对可重试异常触发重试 .filter(throwable -> throwable instanceof RetryableException) .onRetryExhaustedThrow((spec, signal) -> new RetryExhaustedException("重试次数耗尽", signal.failure())) ) // 异常兜底:写入死信队列 .onErrorResume(throwable -> sendToDlq(record, throwable)) ) .subscribe(); } // 业务处理逻辑 private Mono<Void> processMessage(ConsumerRecord<String, String> record) { return Mono.fromRunnable(() -> { String messageValue = record.value(); // 此处替换为实际业务逻辑 if (messageValue.contains("retry")) { throw new RetryableException("触发可重试异常"); } if (messageValue.contains("invalid")) { throw new NonRetryableException("消息格式非法,无需重试"); } log.info("消息处理成功,内容:{}", messageValue); }); } // 写入死信队列 private Mono<Void> sendToDlq(ConsumerRecord<String, String> record, Throwable error) { return dlqProducerTemplate.send(DLQ_TOPIC, record.key(), record.value()) .doOnSuccess(result -> log.info("消息已写入死信队列,偏移量:{},异常原因:{}", result.recordMetadata().offset(), error.getMessage())) .then(); } // 自定义异常类 public static class RetryableException extends RuntimeException { public RetryableException(String message) { super(message); } } public static class NonRetryableException extends RuntimeException { public NonRetryableException(String message) { super(message); } } public static class RetryExhaustedException extends RuntimeException { public RetryExhaustedException(String message, Throwable cause) { super(message, cause); } } }
关键说明
- 偏移量管理:使用
receiveAutoAck时,仅当消息处理成功(包括重试成功、写入死信成功)才会自动提交偏移量,不会出现消息丢失问题 - 重试过滤:通过
filter方法精准控制需要重试的异常类型,避免对业务非法类异常做无效重试 - 恢复逻辑:死信队列的消息可后续通过定时任务、人工后台等方式做二次处理,实现兜底恢复
内容的提问来源于stack exchange,提问作者Saurabh Agarwal
相关产品推荐
相关产品推荐

