Spring WebFlux单元测试:WebTestClient遇EmitterProcessor阻塞问题
我完全懂你遇到的这个头疼问题——实际运行时,返回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

