I have a primitive Flux of Strings, and run this code in the main() method.
package com.example;
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;
import reactor.util.Logger;
import reactor.util.Loggers;
import java.util.Arrays;
import java.util.List;
public class Parallel {
private static final Logger log = Loggers.getLogger(Parallel.class.getName());
private static List<String> COLORS = Arrays.asList("red", "white", "blue");
public static void main(String[] args) throws InterruptedException {
Flux<String> flux = Flux.fromIterable(COLORS);
flux
.log()
.map(String::toUpperCase)
.subscribeOn(Schedulers.newParallel("sub"))
.publishOn(Schedulers.newParallel("pub", 1))
.subscribe(value -> {
log.info("==============Consumed: " + value);
});
}
}
If you try to run this code the app never stops running and you need to stop it manually.
If I replace .newParallel() with .parallel() everything works as expected and the app finishes normally.
Why can't it finish running on its own? Why does it hang? What is the reason for this behavior?
If you run this code as a JUnit test it works fine and it doesn't hang.