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

Spring Cloud RabbitMQ 函数式路由表达式失效问题排查

Spring Cloud Stream路由失效问题(Spring Boot 3迁移后)

从Spring Boot 2.3.3+Spring Cloud Hoxton.SR8迁移到Spring Boot 3.3.1+Spring Cloud 2023.0.2,应用基于Spring Cloud Stream+RabbitMQ构建。已适配函数式编程模型,但出现核心路由异常:

  • 两个消费绑定共享同一目标队列名称storage,需通过amqp_receivedRoutingKey请求头区分路由到processDelivery或processCargo函数
  • 迁移后单条消息会同时被两个函数接收,导致原有单元测试失败,且无法修改上游发送端(storage-service)的队列目标名称

当前配置信息

cloud:
  function:
    definition: processBaseDelivery;processDelivery;processCargo
    #routing-expression: "headers.amqp_receivedRoutingKey.toString().contains('delivery') ? 'processDelivery' : 'processCargo'"
    routing-expression: "headers['amqp_receivedRoutingKey'] == 'delivery.INSERT' ? 'processDelivery' : 'processCargo'"
  stream:
    bindings:
      processBaseDelivery-in-0:
        destination: transport.baseDeliveryResponse
        group: delivery-service
      processDelivery-in-0:
        destination: storage
        group: delivery-service.delivery
      processCargo-in-0:
        destination: storage
        group: delivery-service.cargo
    rabbit:
      bindings:
        processBaseDelivery-in-0:
          consumer:
            transacted: true
            requeue-rejected: true
            bindingRoutingKey: delivery.response
        processDelivery-in-0:
          consumer:
            max-concurrency: 1
            autoBindDlq: true
            transacted: true
            requeue-rejected: false
            bindingRoutingKey: delivery.*
        processCargo-in-0:
          consumer:
            max-concurrency: 1
            autoBindDlq: true
            transacted: true
            requeue-rejected: false
            bindingRoutingKey: cargo.*

单元测试代码

@SpringBootTest
@ExtendWith(MockitoExtension.class)
@ActiveProfiles("it")
@Import(TestChannelBinderConfiguration.class)
public class MessageProcessorServiceTest {

    @MockBean
    private CargoService cargoService;

    @MockBean
    private DeliveryService deliveryService;

    @Autowired
    private InputDestination inputDestination;

    @Test
    public void processDelivery_delete() {
        final DeliveryMessage deliveryMessage = new DeliveryMessage();
        deliveryMessage.setId(1L);
        deliveryMessage.setOperation(Operation.DELETE);

        inputDestination.send(
            MessageBuilder.withPayload(deliveryMessage)
                .setHeader(AmqpHeaders.RECEIVED_ROUTING_KEY, "delivery.INSERT")
                .build(), "lagerlegging");

        Mockito.verify(deliveryService, new Times(1))
            .deleteDelivery(eq(1L));
        Mockito.verifyNoInteractions(cargoService);
    }

}

已尝试的函数定义形式

// 形式1:直接接收Payload
public Consumer<DeliveryMessage> processDelivery()

// 形式2:接收完整Message对象
public Consumer<Message<DeliveryMessage>> processDelivery()

解决方案

1. 修正路由表达式配置层级

Spring Cloud Stream 3.x+版本中,routing-expression不属于cloud.function配置层级,需配置在cloud.stream.bindings.<binding-name>.consumer下(针对单个绑定)或cloud.stream.default.consumer下(全局生效)。原配置的层级错误导致路由规则未生效。

方式一:全局路由表达式(推荐)

cloud:
  function:
    definition: processBaseDelivery;processDelivery;processCargo
  stream:
    default:
      consumer:
        routing-expression: "headers['amqp_receivedRoutingKey'].startsWith('delivery.') ? 'processDelivery' : 'processCargo'"
    bindings:
      # 原有bindings配置保持不变
      processBaseDelivery-in-0:
        destination: transport.baseDeliveryResponse
        group: delivery-service
      processDelivery-in-0:
        destination: storage
        group: delivery-service.delivery
      processCargo-in-0:
        destination: storage
        group: delivery-service.cargo
    # 原有rabbit配置保持不变

方式二:绑定级路由过滤

针对两个共享storage目标的绑定,分别设置路由过滤规则,确保每个绑定只接收对应路由键的消息:

cloud:
  function:
    definition: processBaseDelivery;processDelivery;processCargo
  stream:
    bindings:
      processBaseDelivery-in-0:
        destination: transport.baseDeliveryResponse
        group: delivery-service
      processDelivery-in-0:
        destination: storage
        group: delivery-service.delivery
        consumer:
          routing-expression: "headers['amqp_receivedRoutingKey'].startsWith('delivery.')"
      processCargo-in-0:
        destination: storage
        group: delivery-service.cargo
        consumer:
          routing-expression: "headers['amqp_receivedRoutingKey'].startsWith('cargo.')"
    # 原有rabbit配置保持不变

2. 修正单元测试的消息目标

单元测试中消息发送的目标是lagerlegging,但配置中绑定的目标是storage,导致消息无法进入正确的路由链路,需修改为:

inputDestination.send(
    MessageBuilder.withPayload(deliveryMessage)
        .setHeader(AmqpHeaders.RECEIVED_ROUTING_KEY, "delivery.INSERT")
        .build(), "storage"); // 替换原有的lagerlegging

3. 统一函数定义形式

推荐使用Consumer<Message<DeliveryMessage>>形式,确保可以直接访问消息头,避免Spring自动转换Payload时丢失头信息:

@Bean
public Consumer<Message<DeliveryMessage>> processDelivery() {
    return message -> {
        DeliveryMessage payload = message.getPayload();
        deliveryService.deleteDelivery(payload.getId());
    };
}

4. 验证RabbitMQ绑定规则

原有RabbitMQ绑定的bindingRoutingKey配置(delivery.*/cargo.*)是正确的,可与Stream路由表达式配合使用,确保消息在RabbitMQ层面就完成初步过滤,减少不必要的消息流转。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 08:49:54