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

如何在Testcontainers的RabbitMQ中向Spring Cloud Stream指定队列发消息做集成测试?

Spring Cloud Stream集成测试:向指定队列发送消息的解决方案

问题描述

我需要为基于Spring Cloud Function的pack和label方法(通过pipe组合为pack|label)编写集成测试。已知应用配置的输入绑定监听order-accepted交换机下的order-accepted.dispatcher-service队列,处理完成后Spring Cloud Stream会自动将消息发送到order-dispatched交换机。我的需求是向order-accepted.dispatcher-service队列发送测试消息,并验证order-dispatched交换机的输出结果,但不清楚该用什么路由键来向目标队列发送消息?

应用配置(application.yml)

server:
  port: 9003

spring:
  application:
    name: dispatcher-service
  cloud:
    function:
      definition: pack|label
    stream:
      bindings:
        packlabel-in-0:
          destination: order-accepted
          group: ${spring.application.name}
        packlabel-out-0:
          destination: order-dispatched

  rabbitmq:
    host: localhost
    port: 5672
    username: user
    password: password
    connection-timeout: 5s

核心解决思路

Spring Cloud Stream与RabbitMQ集成时,默认遵循以下规则:

  • 队列命名规则:{destination}.{group},对应你的场景就是order-accepted.dispatcher-service
  • 交换机与队列的默认绑定路由键为#(匹配所有消息)

你有两种可靠的方式向目标队列发送消息:

  1. 直接向队列发送:跳过交换机,直接将消息投递到目标队列,无需关心路由键
  2. 通过交换机发送:向order-accepted交换机发送消息,使用路由键#,消息会被路由到绑定的order-accepted.dispatcher-service队列

另外,测试时无需手动创建输入队列和绑定——应用启动时Spring Cloud Stream已经自动完成了这些操作,你只需要直接使用即可。

修正后的集成测试代码

import org.junit.jupiter.api.Test;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.core.ParameterizedTypeReference;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.containers.RabbitMQContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;

import java.util.concurrent.TimeUnit;

import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;

@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
@Testcontainers
class DispatcherServiceApplicationTests {

    @Container
    static RabbitMQContainer rabbitMQ = new RabbitMQContainer(DockerImageName.parse("rabbitmq:3.10-management"));

    @DynamicPropertySource
    static void rabbitMQProperties(DynamicPropertyRegistry registry){
        registry.add("spring.rabbitmq.host", rabbitMQ::getHost);
        registry.add("spring.rabbitmq.port", rabbitMQ::getAmqpPort);
        registry.add("spring.rabbitmq.username", rabbitMQ::getAdminUsername);
        registry.add("spring.rabbitmq.password", rabbitMQ::getAdminPassword);
    }

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Autowired
    private RabbitAdmin rabbitAdmin;

    @Test
    void contextLoads() {
    }

    @Test
    void packAndLabel(){
        long orderId = 121;
        String targetInputQueue = "order-accepted.dispatcher-service";
        String outputExchange = "order-dispatched";

        // 确认输入队列已被Spring Cloud Stream自动创建
        assertThat(rabbitAdmin.getQueueProperties(targetInputQueue)).isNotNull();

        // 创建临时队列绑定到输出交换机,用于捕获处理后的消息
        Queue outputQueue = rabbitAdmin.declareQueue();
        assertThat(outputQueue).isNotNull();
        Binding outputBinding = new Binding(
                outputQueue.getName(),
                Binding.DestinationType.QUEUE,
                outputExchange,
                "#", // 匹配输出交换机的所有消息
                null
        );
        rabbitAdmin.declareBinding(outputBinding);

        // 方式1:直接向目标队列发送消息(推荐,无需关心路由键)
        rabbitTemplate.convertAndSend(targetInputQueue, new OrderAcceptedMessage(orderId));

        // 方式2:通过交换机+路由键发送(与Spring Cloud Stream的默认绑定规则匹配)
        // rabbitTemplate.convertAndSend("order-accepted", "#", new OrderAcceptedMessage(orderId));

        // 等待并验证输出消息
        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> {
            OrderDispatchedMessage message = rabbitTemplate.receiveAndConvert(
                    outputQueue.getName(),
                    10000,
                    new ParameterizedTypeReference<OrderDispatchedMessage>(){}
            );

            assertThat(message).isNotNull();
            assertThat(message.orderId()).isEqualTo(orderId);
        });
    }
}

代码说明

  • 输入队列使用:直接使用Spring Cloud Stream自动创建的order-accepted.dispatcher-service队列,无需手动创建
  • 消息发送:提供两种可选方式,直接发送到队列更直观可靠
  • 输出消息捕获:创建临时队列绑定到order-dispatched交换机,确保能接收到所有输出消息
  • 验证逻辑:使用Awaitility等待消息异步处理完成,然后断言消息内容符合预期

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 19:44:57