测试Project Reactor的publish:Mock接收Flux的方法并订阅输入Flux
Mock接收Flux的方法时避免Cancel信号的解决方案
问题场景
测试Reactor代码时,Mock了接收Flux作为参数的bar.delete()方法,运行测试后发现上游Flux触发了cancel信号,原因是传入bar.delete()的Flux未被订阅。需要调整Mock逻辑,确保输入Flux被正常订阅消费。
相关代码
测试类代码
@ExtendWith(MockitoExtension.class) public class FooTest { @Mock Bar bar; @InjectMocks Foo foo; @Test void test() { when(bar.delete(any())).thenReturn(Mono.empty()); StepVerifier.create(foo.deleteRelated("123")) .expectNext() .verifyComplete(); } }
Foo类代码
public class Foo { private final Bar bar; public Foo(Bar bar) { this.bar = bar; } public Mono<Void> deleteRelated(String id) { return Mono.just(id) .flatMapMany(this::findById) .log() .publish(bar::delete) .collectList() .then(); } private Flux<String> findById(String id) { return Flux.just("a", "b", "c"); } }
Bar类代码
public class Bar { public Mono<Void> delete(Flux<String> vals) { return vals.buffer(3) .flatMap(this::deleteAFewAtATime) .collectList() .then(); } private Mono<Void> deleteAFewAtATime(List<String> key) { return Mono.empty(); } }
测试运行输出
reactor.Flux.MonoFlatMapMany.1 : onSubscribe([Synchronous Fuseable] FluxArray.ArraySubscription) reactor.Flux.MonoFlatMapMany.1 : cancel()
解决方案
1. 修正Mock逻辑,确保订阅传入的Flux
原Mock直接返回Mono.empty(),完全未处理传入的Flux,导致该Flux未被订阅,上游触发cancel信号。改用thenAnswer获取传入的Flux并确保其被订阅:
when(bar.delete(any())).thenAnswer(invocation -> { // 获取传入的Flux参数 Flux<String> inputFlux = invocation.getArgument(0); // 订阅并消费Flux,返回完成信号 return inputFlux.then(); });
2. 修正测试中的StepVerifier断言
foo.deleteRelated()返回Mono<Void>,不会发射任何元素,原测试的.expectNext()会导致断言失败,需移除该断言直接验证完成:
StepVerifier.create(foo.deleteRelated("123")) .verifyComplete();
原理说明
thenAnswer能在Mock方法被调用时获取实际传入的Flux对象,通过inputFlux.then()可订阅并消费整个Flux序列直到完成,避免上游Flux被取消。Mono<Void>类型序列不发射元素,因此StepVerifier无需添加expectNext()断言,直接验证完成即可。
内容的提问来源于stack exchange,提问作者Lucas
相关产品推荐
相关产品推荐

