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

如何在无消息中间件时测试同目的地收发的Spring Cloud Stream应用

解决方案

问题根源

你的orderSubmittedEventHandler是Function类型,处理完OrderSubmittedEvent后会将OrderCreatedEvent发送到functionRouter-out-0绑定的order目的地;而orderCreatedEventHandler是Consumer类型,监听的也是order目的地的同一分组org.example.orderprocessing。

在Spring Cloud Stream的测试绑定器(Test Binder)中,消息发送到目的地后会直接被同分组的Consumer消费,不会进入OutputDestination,因此output.receive()返回null,导致空指针异常。

方案一:拆分绑定目的地或分组

修改application.yml,将Function的输出目的地与Consumer的监听目的地分离,或者给Consumer配置不同的分组:

spring:
  cloud:
    stream:
      function:
        routing:
          enabled: true
      bindings:
        functionRouter-in-0:
          destination: order-submitted
          group: org.example.orderprocessing
        functionRouter-out-0:
          destination: order-created
          content-type: application/json;type=org.example.order.OrderCreatedEvent
        orderCreatedEventHandler-in-0:
          destination: order-created
          group: org.example.orderprocessing

测试时,input.send()将消息发送到order-submitted,output.receive()即可从order-created获取到Function输出的OrderCreatedEvent。

方案二:直接验证Consumer的调用行为(推荐)

既然OrderCreatedEvent会被orderCreatedEventHandler消费,无需通过OutputDestination捕获消息,直接通过Mock验证Consumer是否被调用即可:

@SpringBootTest
@EnableTestBinder
@ExtendWith(MockitoExtension.class)
public class OrderProcessingApplicationTest {

    @Autowired
    private InputDestination input;

    @MockBean
    private Consumer<OrderCreatedEvent> orderCreatedEventHandler;

    @Test
    void testOrderProcessingFlow() {
        // 构造测试事件
        UUID orderId = UUID.randomUUID();
        UUID customerId = UUID.randomUUID();
        List<LineItem> lineItems = List.of(new LineItem("foo", 10, 1), new LineItem("bar", 20, 2));
        OrderSubmittedEvent submittedEvent = new OrderSubmittedEvent(orderId, customerId, lineItems);

        // 发送消息
        input.send(MessageBuilder.withPayload(submittedEvent)
                .setHeader("contentType", "application/json;type=org.example.shoppingcart.OrderSubmittedEvent")
                .build());

        // 验证Consumer是否在超时时间内被调用
        ArgumentCaptor<OrderCreatedEvent> eventCaptor = ArgumentCaptor.forClass(OrderCreatedEvent.class);
        verify(orderCreatedEventHandler, timeout(1000)).accept(eventCaptor.capture());

        // 验证事件内容正确性
        OrderCreatedEvent createdEvent = eventCaptor.getValue();
        assertThat(createdEvent.customerId()).isEqualTo(customerId);
        assertThat(createdEvent.lineItems()).isEqualTo(lineItems);
    }
}

需确保项目中引入了Mockito依赖,Spring Boot项目可直接使用spring-boot-starter-test,其中包含Mockito相关包。

方案三:直接操作绑定通道

通过TestChannelBinderConfiguration自定义测试配置,直接操作Function的输入输出通道,绕过目的地路由:

@SpringBootTest
@EnableTestBinder
@Import(TestChannelBinderConfiguration.class)
public class OrderProcessingApplicationTest {

    @Autowired
    private MessageChannel functionRouterIn0;

    @Autowired
    private PollableChannel functionRouterOut0;

    @Test
    void testMessageFlow() throws InterruptedException {
        // 构造并发送测试消息
        OrderSubmittedEvent submittedEvent = new OrderSubmittedEvent(UUID.randomUUID(), UUID.randomUUID(), List.of(new LineItem("foo", 10, 1)));
        Message<OrderSubmittedEvent> message = MessageBuilder.withPayload(submittedEvent)
                .setHeader("contentType", "application/json;type=org.example.shoppingcart.OrderSubmittedEvent")
                .build();
        functionRouterIn0.send(message);

        // 从输出通道接收消息并验证
        Message<?> outputMessage = functionRouterOut0.receive(1000);
        assertThat(outputMessage).isNotNull();
        assertThat(outputMessage.getPayload()).isInstanceOf(OrderCreatedEvent.class);
    }
}

这种方式适合单独测试Function的转换逻辑,无需关注后续的Consumer消费流程。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 12:25:58