I have a list of tasks that are meant to be executed sequentially, returning true if the task succeeded, or false if it failed (non-exceptionally - there is a third case where it encountered some sort of exception, which should also abort processing, but I want to treat that as a distinct third case).
These tasks are represented as a list of Mono<Boolean> instances. There could be any number of them. From this list, I want to produce a Mono<Boolean> that has the following characteristics:
- It executes each task sequentially.
- If any task in the list returns
falseor encounters an exception, it should abort processing of any further tasks. - It returns
trueif all tasks succeeded;falseif any task failed non-exceptionally; it should emit an error if any tasks encountered an exception.
Here is an example list of 3 tasks - the first one succeeds, and the second one fails - so I don't want the third one to run:
Mono<Boolean> task1 = Mono.create(sink -> {
try {
System.out.println("Executing task 1...");
// Simulate some work being done
Thread.sleep(1000);
System.out.println("Finished executing task 1.");
// This task is successful, return true
sink.success(true);
} catch (InterruptedException e) {
sink.error(e);
}
});
Mono<Boolean> task2 = Mono.create(sink -> {
try {
System.out.println("Executing task 2...");
// Simulate some work being done
Thread.sleep(1000);
System.out.println("Finished executing task 2.");
// This task is NOT successful, return false
sink.success(false);
} catch (InterruptedException e) {
sink.error(e);
}
});
Mono<Boolean> task3 = Mono.create(sink -> {
try {
System.out.println("Executing task 3...");
// Simulate some work being done
Thread.sleep(1000);
System.out.println("Finished executing task 3.");
// This task is successful, return true
sink.success(true);
} catch (InterruptedException e) {
sink.error(e);
}
});
List<Mono<Boolean>> tasks = Arrays.asList(task1, task2, task3);
My first thought was to convert the list of Monos to a Flux, and then use Flux#reduce:
Mono<Boolean> process = Flux.concat(tasks)
.reduce(true, (accum, value) -> accum && value);
process.subscribe(System.out::println);
This correctly produces false, which indicates that one of the steps in the process did not succeed (step 2). However, given the output, it is clear that task 3 is still being executed, which is not correct:
Executing task 1...
Finished executing task 1.
Executing task 2...
Finished executing task 2.
Executing task 3...
Finished executing task 3.
false
I can get the desired behavior if I go back to the synchronous world and use Mono#block, but this doesn't feel very idiomatic:
Mono<Boolean> process = Mono.create(sink -> {
Boolean result = true;
for (Mono<Boolean> task: tasks) {
result = task.block();
if (!result) {
break;
}
}
sink.success(result);
});
process.subscribe(System.out::println);
Is there a better way to do this using the built-in operators like reduce, map, flatMap, etc.?