如何在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 - 交换机与队列的默认绑定路由键为
#(匹配所有消息)
你有两种可靠的方式向目标队列发送消息:
- 直接向队列发送:跳过交换机,直接将消息投递到目标队列,无需关心路由键
- 通过交换机发送:向
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
相关产品推荐
相关产品推荐

