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

Reactor中TestPublisher与非终止Flux配合使用失效问题求助

这个问题的核心在于TestPublisher的默认严格行为和Sink的默认缓存策略之间的差异,我来帮你理清其中的误区:

首先,你遇到的测试失败原因是:TestPublisher.create()创建的是严格模式的发布者,它默认遵循Reactive Streams规范,不允许在有订阅者之前发送事件(比如调用next())。当你在StepVerifier.create(result)(也就是订阅发生)之前调用publisher.next("hello"),这个元素并不会被TestPublisher缓存下来,后续的订阅者(StepVerifier)自然接收不到它,最终导致超时失败。

而你用Sinks.many().unicast().onBackpressureBuffer()创建的Sink,默认就支持订阅前发送元素并缓存——它的onBackpressureBuffer策略会把订阅前的元素暂存起来,等有订阅者连接时再发送,所以StepVerifier能正常收到"hello"。

要让TestPublisher的测试正常工作,你需要显式关闭它的这个严格检查,创建一个允许订阅前发送事件的TestPublisher,具体修改如下:

@Test
void test_with_test_publisher() {
    // 创建允许订阅前发送事件的TestPublisher
    TestPublisher<String> publisher = TestPublisher.createNonCompliant(TestPublisher.Violation.ALLOW_EVENTS_BEFORE_SUBSCRIPTION);
    publisher.next("hello");
    var result = publisher.flux();
    StepVerifier.create(result)
            .expectNext("hello").as("initial message received")
            .verifyTimeout(Duration.ofSeconds(1));
}

这样修改后,TestPublisher会缓存订阅前发送的"hello"元素,StepVerifier订阅后就能正常接收到,测试也就通过了。

补充一点:TestPublisher的严格模式是为了帮助开发者测试符合Reactive Streams规范的发布者行为——按照规范,发布者不应该在没有订阅者的情况下发送事件。但在你的测试场景中,你需要模拟一个“提前发送元素”的场景,所以需要主动关闭这个检查。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 12:22:48