Неверно работает Flux.merge()
Есть контроллер, который возвращает клиенту Flux результатов расчетов. Каждый расчет оформлен в виде экземпляра класса Answer, имеющего несколько полей, среди которых номер расчета, номер функции, по которой рассчитывается результат (всего этих функций две), и сам результат. Соответственно создаю два Flux, первый считает результаты и создает экземпляры Answer для первой функции, второй делает тоже самое для второй. При чем, расчет одной из функций искуственно замедлен путем Thread.sleep(4000). В клиент возвращается Flux.merge() от этих двух потоков. По идее, результаты должны приходить вперемежку, то есть по мере выполнения расчетов, но по факту приходят сначала результаты первого Flux, потом второго, словно используется concatWith, а не merge(). Буду очень благодарен за совет, почему это.
@RestController
@RequestMapping("/resolve")
public class MainController {
@GetMapping(produces = MediaType.APPLICATION_STREAM_JSON_VALUE)
public Flux<Answer> getResolving(@RequestParam(defaultValue = "0") int numberOfIterations) {
Flux<Answer> answers1 = Flux.range(1, numberOfIterations).map(integer -> {
return new Answer(1, FuncResolver.function1(integer), integer);
});
Flux<Answer> answers2 = Flux.range(1, numberOfIterations).map(integer -> {
return new Answer(2, FuncResolver.function2(integer), integer);
});
Flux<Answer> merge = Flux.merge(answers1, answers2);
return merge;
}
}
public class Answer {
private int funcNumber;
private String iterationNumber;
private String result;
public Answer(int funcNumber, int result, int iterationNumber) {
this.funcNumber = funcNumber;
this.iterationNumber = Integer.toString(iterationNumber);
this.result = Integer.toString(result);
}
public class FuncResolver {
public static int function1(int i) {
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return i * 2;
}
public static int function2(int i) {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return i * 2 + 1;
}
}