I'm having a hard time understanding the behaviour of the following code sample;
Flowable<String> f = Flowable.just(1)
.flatMap(it -> Flowable.create(e -> {
for(int i = 1; i < 1001; ++i) {
log.info("Emitting: " + i);
if(i % 10 == 0) {
Thread.sleep(1000);
}
e.onNext(i);
}
e.onComplete();
}, BackpressureStrategy.BUFFER))
.map(String::valueOf)
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.newThread());
f.subscribe(val -> {
Thread.sleep(100);
log.info("Observing: " + val);
});
Thread.sleep(1000000);
The code works OK until 128 items are observed by the subscribe call. Emit and observe are in parallel. But after that, Flowable continues to emit items (which are queued somewhere obviously) but no item is observed until all 1000 of items are emitted. After all 1000 items are emitted, then the rest of items (> 128) are observed at once.
This looks to be related to backpressure bufferSize of 128 but still I would expect the emit and observe to be in parallel for the whole 1000 items, because observer is obviously not slower than the emitter. Is there something I'm missing here? What should I do to fix the code?