How can this project-reactor behavior explained?

Viewed 56

Here is a simple project-reactor code snippet:

Consumer<String> slowConsumer = x -> {
    try {
        TimeUnit.MILLISECONDS.sleep(50);
    } catch(Exception ignore) {
    }
};
Flux<String> publisher = Flux
        .just(1, 2, 3)
        .parallel(2)
        .runOn(Schedulers.newParallel("writer", 2))
        .flatMap(rail -> Flux.range(1, 300).map(i -> String.format("Rail %d -> Row %d", rail, i)).log())
        .sequential()
        .publishOn(Schedulers.newSingle("reader"))
        .doOnNext(slowConsumer);
publisher.subscribe();

My expectation is that everything that happens in "flatMap" should be executed within "writer" threads. However, this is what is logged:

onNext(Rail 2 -> Row 300) [Flux.MapFuseable.1] [writer-2]
...
onNext(Rail 1 -> Row 300) [Flux.MapFuseable.2] [writer-1]
...
onNext(Rail 3 -> Row 256) [Flux.MapFuseable.3] [writer-1]
request(256) [Flux.MapFuseable.3] [reader-1]
onNext(Rail 3 -> Row 257) [Flux.MapFuseable.3] [reader-1]
...
onNext(Rail 3 -> Row 300) [Flux.MapFuseable.3] [reader-1]

Can someone explain this? How come "reader" thread is processing the tail of the last rail? What am I missing?

0 Answers
Related