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

