I have a potentially long-running Flux that I'd like to stop after a certain duration has passed. I've found several methods of doing this, however what I'm struggling with is how to be able to tell that the Flux timed out rather than just completed naturally.
Sample (very simple) code:
Flux.range(0, 10000)
.take(Duration.ofMillis(1))
.doOnNext(System.out::println)
.collectList()
.block();
What I'd like is something like this:
Flux.range(0, 10000)
.take(Duration.ofMillis(1))
.doOnNext(System.out::println)
.doOnError(t -> {
if (t instanceof TimeoutException) {
System.out.println("I timed out");
}
})
.collectList()
.block();
however take doesn't seem to notify on error; all I see is a terminate and complete signal which is what I get if I don't include the take in the Flux and just let it complete naturally.
I've looked briefly into the timeout operator which does throw an exception, however timeout looks like it'll only throw an exception if the Flux doesn't emit a single element within a certain time, rather than if the whole Flux doesn't complete in a certain time.
Does anyone have any tips or examples of how they've solved this?
Thanks in advance!