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

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版本的绑定配置分为两个层级,规则如下:

  1. 带-<index>后缀的命名(如myfunction-in-0)是运行时绑定实例标识,仅可用于spring.cloud.stream.bindings节点下配置通用绑定属性,包括目的地、分组、重试次数、并发数等。Spring Cloud Stream Rabbit Binder的扩展属性(如transacted、queueNameGroupOnly)不会识别带索引后缀的绑定名。
  2. 不带索引的命名(如myfunction-in)是绑定逻辑名称,spring.cloud.stream.rabbit.bindings节点下的所有Binder扩展配置,需要匹配逻辑名称,无需添加输入输出索引,框架会自动将扩展属性映射到对应逻辑下的所有索引实例。

除此之外,全链路事务生效需要同时满足三个条件:

  • 消费者端绑定开启transacted属性
  • 消费逻辑方法添加@Transactional注解
  • 配置RabbitTransactionManager作为事务管理器

内容的提问来源于stack exchange,提问作者Sébastien

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 02:06:04