如何将delayElements()与Flux合并一起使用?

abi*_*kay 1 java project-reactor spring-webflux

我正在学习教程,我相信我的代码与讲师的代码相同,但我不明白为什么delayElements()不起作用。

这是调用者方法:

public static void main(String[] args) {
    FluxAndMonoGeneratorService fluxAndMonoGeneratorService = new FluxAndMonoGeneratorService();
    fluxAndMonoGeneratorService.explore_merge()
            .doOnComplete(() -> System.out.println("Completed !"))
            .onErrorReturn("asdasd")
            .subscribe(System.out::println);
}
Run Code Online (Sandbox Code Playgroud)

如果我将没有延迟元素的方法编写为:

public Flux<String> explore_merge() {

        Flux<String> abcFlux = Flux.just("A", "B", "C");
        Flux<String> defFlux = Flux.just("D", "E", "F");

        return Flux.merge(abcFlux, defFlux);
    }
Run Code Online (Sandbox Code Playgroud)

然后控制台中的输出是(如预期):

00:53:19.443 [main] DEBUG reactor.util.Loggers$LoggerFactory - Using Slf4j logging framework
A
B
C
D
E
F
Completed !

BUILD SUCCESSFUL in 1s
Run Code Online (Sandbox Code Playgroud)

但我想使用delayElements()来测试merge()方法:

public Flux<String> explore_merge() {

        Flux<String> abcFlux = Flux.just("A", "B", "C").delayElements(Duration.ofMillis(151));
        Flux<String> defFlux = Flux.just("D", "E", "F").delayElements(Duration.ofMillis(100));

        return Flux.merge(abcFlux, defFlux);
    }
Run Code Online (Sandbox Code Playgroud)

什么也没有发生,onComplete 和 onErrorReturn 都没有发生,并且输出什么也没有:

0:55:22: Executing ':reactive-programming-using-reactor:FluxAndMonoGeneratorService.main()'...

> Task :reactive-programming-using-reactor:generateLombokConfig UP-TO-DATE
> Task :reactive-programming-using-reactor:compileJava
> Task :reactive-programming-using-reactor:processResources NO-SOURCE
> Task :reactive-programming-using-reactor:classes

> Task :reactive-programming-using-reactor:FluxAndMonoGeneratorService.main()
00:55:23.715 [main] DEBUG reactor.util.Loggers$LoggerFactory - Using Slf4j logging framework

BUILD SUCCESSFUL in 1s
Run Code Online (Sandbox Code Playgroud)

这是什么原因呢?(我的意思是至少 onError 我期待着,但什么也没有......)

注意:mergeWith()也不适用于这个delayElements()

Ale*_*lex 5

subscribe不是阻塞操作,delayElements将被调度到另一个线程上(默认parallel调度程序)。结果,您的程序在元素发出之前退出。这是一个测试

@Test
void mergeWithDelayElements() {
    Flux<String> abcFlux = Flux.just("A", "B", "C").delayElements(Duration.ofMillis(151));
    Flux<String> defFlux = Flux.just("D", "E", "F").delayElements(Duration.ofMillis(100));

    StepVerifier.create(Flux.merge(abcFlux, defFlux))
            .expectNext("D", "A", "E", "B", "F", "C")
            .verifyComplete();
}
Run Code Online (Sandbox Code Playgroud)