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

Spring Reactive单元测试遇java.io.IOException: Connection closed prematurely异常

解决Spring Reactive流式服务单元测试的Connection closed prematurely异常

这个问题我之前在做Spring Reactive测试时也碰到过,大概率是测试框架没处理好无限流响应,提前把连接给关了。我给你梳理下原因和具体的解决办法:

首先先确认你的控制器代码是不是补全了(看起来你最后一行代码没写完),正确的写法应该是这样:

@GetMapping(value = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<Event> getEvents() {
    Flux<Event> eventFlux = Flux.fromStream(Stream.generate(() -> new Event(new Random().nextLong(), "Hello Event")));
    Flux<Long> emitFlux = Flux.interval(Duration.ofSeconds(1));
    // 取zip后的第一个元素,也就是生成的Event
    return Flux.zip(eventFlux, emitFlux).map(Tuple2::getT1);
}

核心原因

你的接口返回的是无限流(每秒生成一个事件,永远不会停止),而默认的测试逻辑会在认为响应"完成"时就关闭连接,但无限流永远不会主动完成,所以测试框架会提前断开连接,抛出Connection closed prematurely异常。

解决方案:用WebTestClient正确处理流式响应

Spring Reactive官方推荐用WebTestClient来测试响应式服务,它能很好地处理流式响应。下面给你两种常用的测试方式:

方式1:指定接收的事件数量,让测试等待完成

这种方式适合你只需要验证前N个事件的场景,测试会等待指定数量的事件接收完成后再结束,不会提前关连接:

@WebFluxTest(EventController.class)
class EventControllerTest {

    @Autowired
    private WebTestClient webTestClient;

    @Test
    void testGetEvents() {
        webTestClient.get().uri("/events")
                .accept(MediaType.TEXT_EVENT_STREAM)
                .exchange()
                .expectStatus().isOk()
                // 等待接收3个事件,验证数量和内容
                .expectBodyList(Event.class)
                .hasSize(3)
                .consumeWith(response -> {
                    List<Event> events = response.getResponseBody();
                    assert events != null;
                    events.forEach(event -> {
                        assertNotNull(event.getId());
                        assertEquals("Hello Event", event.getMessage());
                    });
                });
    }
}

方式2:手动控制流的终止和超时

如果需要更灵活的控制(比如测试更长时间的流),可以手动取前N个元素并设置超时时间,确保测试有足够时间接收事件:

@Test
void testGetEventsWithTimeout() {
    webTestClient.get().uri("/events")
            .accept(MediaType.TEXT_EVENT_STREAM)
            .exchange()
            .expectStatus().isOk()
            .returnResult(Event.class)
            .getResponseBody()
            // 只取前3个事件
            .take(3)
            // 设置5秒超时,确保有足够时间接收(每秒1个,3个需要3秒,留冗余时间)
            .blockLast(Duration.ofSeconds(5));
}

额外注意点

  1. 不要用MockMvc测试响应式流式服务:MockMvc是为同步Servlet设计的,处理Reactive的无限流会出问题,一定要用WebTestClient。
  2. 确保Event类可序列化:WebTestClient需要解析响应为Event对象,所以你的Event类必须有:
    • 无参构造函数
    • 对应的getter/setter方法
      示例:
    public class Event {
        private Long id;
        private String message;
    
        public Event() {} // 必须要有无参构造
    
        public Event(Long id, String message) {
            this.id = id;
            this.message = message;
        }
    
        // getter和setter
        public Long getId() { return id; }
        public void setId(Long id) { this.id = id; }
        public String getMessage() { return message; }
        public void setMessage(String message) { this.message = message; }
    }
    

按照上面的方法调整测试代码,应该就能解决这个连接提前关闭的问题了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:51:30