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

