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

Spring WebFlux单元测试:WebTestClient遇EmitterProcessor阻塞问题

WebTestClient Blocks with EmitterProcessor in Spring WebFlux Controller Tests

我完全懂你遇到的这个头疼问题——实际运行时,返回EmitterProcessor的Controller能正常让客户端订阅接收更新,但用WebTestClient做单元测试时,exchange()直接阻塞超时,而返回普通Flux的测试却完全正常。咱们来拆解一下原因和解决办法:

问题根源

首先得搞清楚EmitterProcessor和普通Flux(比如Flux.just)的本质区别:

  • Flux.just("hello")是冷发布者:订阅后会立即发射所有数据,然后发送完成信号,整个流是有限的。WebTestClient的exchange()能快速拿到响应状态、头信息,然后顺利进入后续的断言步骤。
  • EmitterProcessor.create()是热发布者:默认创建的是一个空的、不会自动完成的流。如果没有调用onNext()发射数据,服务器会一直保持SSE连接,但不会发送任何内容。WebTestClient默认的响应超时是5秒,在这段时间内如果没有收到任何数据,就会触发你看到的IllegalStateException超时错误。

你提到的“注释掉processor.onNext("hello");测试失败,放开就成功”也印证了这一点:有数据发射时,服务器发送了第一个SSE事件,WebTestClient的exchange()能顺利完成阻塞,继续执行StepVerifier的断言;没有数据时,连接一直空转直到超时。

解决方案

根据你的测试场景,这里有几种可行的解决方式:

1. 在测试中手动触发数据发射

如果你的Controller设计是依赖外部触发数据(比如实际业务中其他组件调用processor.onNext()),那在测试里需要拿到Processor的引用,手动发送测试数据。可以修改Controller暴露一个测试用的方法:

@RestController
static class StringController {
    private EmitterProcessor<String> processor = EmitterProcessor.create();

    @GetMapping(value = "/processor", produces = TEXT_EVENT_STREAM_VALUE)
    Flux<String> getProcessor() {
        return processor;
    }

    // 仅用于测试:手动发送消息
    public void triggerMessage(String message) {
        processor.onNext(message);
    }
}

然后在测试用例中,先获取Controller实例,再在exchange()之后触发数据:

@Test
public void entityStreamProcessor() {
    StringController controller = new StringController();
    WebTestClient client = WebTestClient.bindToController(controller)
            .configureClient()
            .build();

    FluxExchangeResult<String> result = client
            .get().uri("/processor")
            .accept(TEXT_EVENT_STREAM)
            .exchange()
            .expectStatus().isOk()
            .expectHeader().contentTypeCompatibleWith(TEXT_EVENT_STREAM)
            .returnResult(String.class);

    // 手动触发数据发射
    controller.triggerMessage("hello");

    StepVerifier.create(result.getResponseBody())
            .expectNext("hello")
            .thenCancel()
            .verify();
}

2. 给Processor初始化默认数据

如果你的业务场景允许,创建EmitterProcessor时可以传入初始数据,这样每次请求都会有默认数据发射,测试就能顺利进行:

@GetMapping(value = "/processor", produces = TEXT_EVENT_STREAM_VALUE)
Flux<String> getProcessor() {
    // 初始化时传入默认数据
    return EmitterProcessor.create(Collections.singletonList("hello"));
}

3. 调整WebTestClient的超时时间(不推荐作为根本解决办法)

如果只是临时调试,可以延长WebTestClient的响应超时,但这只能缓解超时问题,没法解决“没有数据就无法通过断言”的核心问题:

private final WebTestClient client = WebTestClient.bindToController(new StringController())
        .configureClient()
        .responseTimeout(Duration.ofSeconds(10)) // 延长超时时间
        .build();

为什么实际运行正常?

实际生产环境中,客户端(比如浏览器、前端应用)订阅SSE流后,会保持长连接等待后续数据,不会像WebTestClient那样有严格的超时限制(或者客户端的超时设置足够长)。只要后续有数据发射,客户端就能收到,所以运行正常;但测试中因为默认超时较短,且没有数据触发,就会出现阻塞超时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:18:15