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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 09:36:03