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
相关产品推荐
相关产品推荐

