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

测试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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 23:50:42