Spring Cloud Stream StreamBridge 事务异常未回滚及配置差异咨询
Spring Cloud Stream 监听器全链路事务配置问题
需求背景
需要实现 Spring Cloud Stream 监听器的全链路事务处理:函数内通过 StreamBridge 手动发送的所有消息,若后续抛出异常应当全部回滚,而非直接提交。
环境依赖
spring : 2.5.5 spring cloud stream : 3.1.4 spring cloud stream rabbit binder : 3.1.4
初始配置与代码
YAML配置
spring: cloud: function: definition: test stream: rabbit: default: producer: transacted: true consumer: transacted: true bindings: test-in-0: consumer: queueNameGroupOnly: true receive-timeout: 500 transacted: true test-out-0: producer: queueNameGroupOnly: true transacted: true other-out-0: producer: queueNameGroupOnly: true transacted: true bindings: test-in-0: destination: test.request group: test.request consumer: requiredGroups: test.request maxAttempts: 1 test-out-0: destination: test.response group: test.response producer: requiredGroups: test.response other-out-0: destination: test.other.request group: test.other.request producer: requiredGroups: test.other.request
Java代码
函数定义
@Configuration public class TestSender { @Bean public Function<Message<TestRequest>, Message<String>> test(Service service) { return (request) -> service.run(request.getPayload().getContent()); } }
业务逻辑类
@Component @Transactional public class Service { private static final Logger LOGGER = LoggerFactory.getLogger(Service.class); StreamBridge bridge; IWorker worker; public Service(StreamBridge bridge, IWorker worker) { this.bridge = bridge; this.worker = worker; } @Transactional public Message<String> run(String message) { LOGGER.info("Processing {}", message); bridge.send("other-out-0", MessageBuilder.withPayload("test") .setHeader("toto", "titi").build()); if (message.equals("error")) { throw new RuntimeException("test error"); } return MessageBuilder.withPayload("test") .setHeader("toto", "titi").build(); } }
启动类
@SpringBootApplication public class EmptyWorkerApplication { private static final Logger LOGGER = LoggerFactory.getLogger(EmptyWorkerApplication.class); public static void main(String[] args) { SpringApplication.run(EmptyWorkerApplication.class, args); } @Bean public ApplicationRunner runner(RabbitTemplate template) { return args -> { LOGGER.info("Sending messages ..."); template.convertAndSend("test.request", "#", org.springframework.amqp.core.MessageBuilder.withBody( "{\"content\":\"toto\"}".getBytes(StandardCharsets.UTF_8)) .setContentType("application/json") .build()); template.convertAndSend("test.request", "#", org.springframework.amqp.core.MessageBuilder.withBody( "{\"content\":\"error\"}".getBytes(StandardCharsets.UTF_8)) .setContentType("application/json") .build()); template.convertAndSend("test.request", "#", org.springframework.amqp.core.MessageBuilder.withBody( "{\"content\":\"titi\"}".getBytes(StandardCharsets.UTF_8)) .setContentType("application/json") .build()); }; } }
事务管理器配置
@Configuration @EnableTransactionManagement public class TransactionManagerConfiguration { @Bean(name = "transactionManager") public RabbitTransactionManager rabbitTransactionManager(ConnectionFactory cf) { RabbitTransactionManager manager = new RabbitTransactionManager(cf); return manager; } }
初始问题现象
运行后Rabbit队列test.other.request最终有3条消息,但预期只有2条(error场景消息应该回滚)。
后续测试调整
调整后代码
@Component("myfunction") public class Myfunction implements Consumer<String> { private final StreamBridge streamBridge; public Myfunction(StreamBridge streamBridge) { this.streamBridge = streamBridge; } @Override @Transactional public void accept(String request) { this.streamBridge.send("myfunction-out-0", request); if (request.equals("error")) { throw new RuntimeException("test error"); } } }
@SpringBootApplication public class EmptyWorkerApplication { public static void main(String[] args) { SpringApplication.run(EmptyWorkerApplication.class, args); } @Bean public RabbitTransactionManager rabbitTransactionManager(ConnectionFactory cf) { RabbitTransactionManager manager = new RabbitTransactionManager(cf); return manager; } @Bean public ApplicationRunner runner(RabbitTemplate template) { return args -> { template.convertAndSend("test.request", "#", org.springframework.amqp.core.MessageBuilder.withBody( "test".getBytes(StandardCharsets.UTF_8)) .setContentType("text/plain") .build()); template.convertAndSend("test.request", "#", org.springframework.amqp.core.MessageBuilder.withBody( "error".getBytes(StandardCharsets.UTF_8)) .setContentType("text/plain") .build()); template.convertAndSend("test.request", "#", org.springframework.amqp.core.MessageBuilder.withBody( "test".getBytes(StandardCharsets.UTF_8)) .setContentType("text/plain") .build()); }; } }
调整后配置
spring: rabbitmq: host: xx port: xx username: xx password: xx virtual-host: xx cloud: function: definition: myfunction stream: rabbit: bindings: myfunction-in-0: queueNameGroupOnly: true myfunction-out-0: queueNameGroupOnly: true transacted: true bindings: myfunction-in-0: destination: test.request group: test.request consumer: requiredGroups: test.request autoBindDlq: true maxAttempts: 1 myfunction-out-0: destination: test.response group: test.response producer: requiredGroups: test.response
问题根因与配置差异说明
修复方案
最终通过调整配置解决问题:错误配置为spring.cloud.stream.rabbit.bindings.myfunction-in-0.consumer.transacted=true,正确配置为spring.cloud.stream.rabbit.bindings.myfunction-in.consumer.transacted=true。
两种配置的差异
Spring Cloud Stream 3.x版本的绑定配置分为两个层级,规则如下:
- 带
-<index>后缀的命名(如myfunction-in-0)是运行时绑定实例标识,仅可用于spring.cloud.stream.bindings节点下配置通用绑定属性,包括目的地、分组、重试次数、并发数等。Spring Cloud Stream Rabbit Binder的扩展属性(如transacted、queueNameGroupOnly)不会识别带索引后缀的绑定名。 - 不带索引的命名(如
myfunction-in)是绑定逻辑名称,spring.cloud.stream.rabbit.bindings节点下的所有Binder扩展配置,需要匹配逻辑名称,无需添加输入输出索引,框架会自动将扩展属性映射到对应逻辑下的所有索引实例。
除此之外,全链路事务生效需要同时满足三个条件:
- 消费者端绑定开启
transacted属性 - 消费逻辑方法添加
@Transactional注解 - 配置
RabbitTransactionManager作为事务管理器
内容的提问来源于stack exchange,提问作者Sébastien
相关产品推荐
相关产品推荐

