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

Spring Cloud Stream Kafka事务中如何获取重试次数

在Spring Cloud Stream Kafka事务重试中获取重试次数并路由到不同Topic

利用Spring Retry上下文获取重试次数

Spring Cloud Stream的重试机制基于Spring Retry实现,你可以直接通过RetryContext获取当前重试次数,再根据次数将消息路由到不同Topic。

1. 基于@Retryable和@Recover的实现

先配置重试模板,定义最大重试次数和间隔:

@Configuration
@EnableRetry
public class RetryConfig {
    @Bean
    public RetryTemplate retryTemplate() {
        SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
        retryPolicy.setMaxAttempts(3); // 最多重试2次(算上首次尝试共3次)

        FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
        backOffPolicy.setBackOffPeriod(2000); // 每次重试间隔2秒

        RetryTemplate retryTemplate = new RetryTemplate();
        retryTemplate.setRetryPolicy(retryPolicy);
        retryTemplate.setBackOffPolicy(backOffPolicy);
        return retryTemplate;
    }
}

然后在消费服务中,用@Retryable标记需要重试的方法,在@Recover方法里处理重试后的路由逻辑:

@Service
public class KafkaConsumerService {

    @Autowired
    private StreamBridge streamBridge;

    @Retryable(value = {SpecificBizException.class}, retryTemplate = "retryTemplate")
    @StreamListener(target = "input-topic")
    public void processMessage(Message<String> message) {
        // 你的业务逻辑,抛出特定异常触发重试
        throw new SpecificBizException("业务处理失败");
    }

    @Recover
    public void handleRetryFailure(SpecificBizException e, Message<String> message, RetryContext context) {
        int retryCount = context.getRetryCount();
        // 根据重试次数选择目标Topic
        String targetTopic = switch (retryCount) {
            case 1 -> "retry-first-topic";
            case 2 -> "retry-second-topic";
            default -> "dead-letter-topic";
        };
        // 发送到对应Topic
        streamBridge.send(targetTopic, message);
    }
}

2. 函数式编程模式下的实现

如果用Spring Cloud Stream的函数式消费方式,可以通过RetryListener拦截重试过程,获取次数并路由:

@Configuration
public class KafkaFunctionConfig {

    @Autowired
    private StreamBridge streamBridge;

    @Bean
    public Consumer<Message<String>> messageConsumer() {
        return message -> {
            // 业务逻辑,抛出异常触发重试
            throw new SpecificBizException("处理失败");
        };
    }

    @Bean
    public RetryTemplate retryTemplate() {
        SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
        retryPolicy.setMaxAttempts(3);
        FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
        backOffPolicy.setBackOffPeriod(2000);

        RetryTemplate retryTemplate = new RetryTemplate();
        retryTemplate.setRetryPolicy(retryPolicy);
        retryTemplate.setBackOffPolicy(backOffPolicy);

        // 注册重试监听器,在重试结束后处理路由
        retryTemplate.registerListener(new RetryListener() {
            @Override
            public <T, E extends Throwable> boolean open(RetryContext context, RetryCallback<T, E> callback) {
                return true;
            }

            @Override
            public <T, E extends Throwable> void close(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {
                if (throwable instanceof SpecificBizException) {
                    int retryCount = context.getRetryCount();
                    Message<String> message = (Message<String>) context.getAttribute("current-message");
                    String targetTopic = switch (retryCount) {
                        case 1 -> "retry-1-topic";
                        case 2 -> "retry-2-topic";
                        default -> "dead-letter-topic";
                    };
                    streamBridge.send(targetTopic, message);
                }
            }

            @Override
            public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {
                // 把当前消息存到上下文,方便后续获取
                context.setAttribute("current-message", callback.getArgs()[0]);
            }
        });
        return retryTemplate;
    }
}

原生重试配置下的快捷方式

如果直接用Spring Cloud Stream的原生重试配置(比如在application.yml里配置spring.cloud.stream.bindings.input-topic.consumer.max-attempts),可以直接从消息头X-SCS-RETRY-COUNT获取重试次数:

@StreamListener(target = "input-topic")
public void processMessage(Message<String> message) {
    Integer retryCount = (Integer) message.getHeaders().get("X-SCS-RETRY-COUNT");
    if (retryCount != null) {
        // 根据次数处理路由逻辑
    }
    // 业务逻辑
}

关键注意点

  • 重试次数从0开始计数:首次尝试失败后,第一次重试的count是1,第二次是2,直到达到最大次数后进入恢复逻辑。
  • 事务场景下,要确保重试操作和消息路由在同一个事务边界内,或者根据业务需求调整事务策略,避免数据不一致。
  • 确保已经添加@EnableRetry注解开启Spring Retry功能,否则重试逻辑不会生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 04:25:31