Re: Counter-intuitive behavior of CompleteableFuture
Viktor Klang via Concurrency-interest <[email protected]> Thu, 20 May 2021 13:22:17 +0000
| Newsgroups | gmane.comp.java.jsr.166-concurrency |
|---|---|
| Message-ID | <CANPzfU8npJE8v9OpEqB++-D3mk9Cg2NC6n6ETkG_QBTRM05KeA@mail.gmail.com> |
minimalCompletionStage to the rescue :)
public static void main(String[] args) {
CompletableFuture<String> initial = new CompletableFuture<>();
CompletionStage<String> withSteps = initial
.minimalCompletionStage <---
This will prevent consumers from having toCompletableFuture returning `this`
.thenApply(s -> {
System.out.println("STEP1: " + s);
return s;
})
.thenApply(s -> {
System.out.println("STEP2: " + s);
return s;
});
CompletionStage<String> timeout = withSteps;
timeout.whenComplete((s, throwable) ->
System.out.println("TIMEOUT: " + throwable));
timeout.completeExceptionally(new RuntimeException("TIMEOUT")); <---
now won't compile
initial.complete("SUCCESS");
}
On Thu, May 20, 2021 at 1:14 PM Michał Górniewski <
[email protected]> wrote:
> This looks good from an API design perspective, since it clearly
> informs users that it really shouldn't call complete() on this Future
> :).
>
> But still, in Java you can do:
>
> timeout.toCompletableFuture().completeExceptionally(new
> RuntimeException("TIMEOUT"));
>
> So toCompletableFuture() would need to create new, depended,
> CompletableStage to be safe.
>
> Best regards,
> Michał
>
> czw., 20 maj 2021 o 14:59 Viktor Klang <[email protected]>
> napisał(a):
> >
> > A solution which has been used very successfully in Scala is the
> distinction between Promise (write a single value 0..1 times) and Future
> (read a single value 0..N times), which means that only the thread who owns
> a reference to the Promise can write values, which means that Futures can
> be freely shared since consumers cannot manipulate the result "further up
> the chain".
> >
> > You can achieve something similar in Java by separating
> CompletableFuture from CompletionStage.
> >
> > public static void main(String[] args) {
> > CompletableFuture<String> initial = new CompletableFuture<>();
> >
> > CompletionStage<String> withSteps = initial
> > .thenApply(s -> {
> > System.out.println("STEP1: " + s);
> > return s;
> > })
> > .thenApply(s -> {
> > System.out.println("STEP2: " + s);
> > return s;
> > });
> >
> > CompletionStage<String> timeout = withSteps;
> >
> > timeout.whenComplete((s, throwable) ->
> > System.out.println("TIMEOUT: " + throwable));
> >
> > timeout.completeExceptionally(new RuntimeException("TIMEOUT")); <---
> now won't compile
> > initial.complete("SUCCESS");
> > }
> >
> >
> > On Thu, May 20, 2021 at 12:50 PM Michał Górniewski via
> Concurrency-interest <[email protected]> wrote:
> >>
> >> What bothers me here is API of CompletableFuture itself, e.g:
> >>
> >> public CompletableFuture<ProcessingOutput>
> >> executeProcessing(Function<Path, ProcessingOutput> processor) {
> >> Path data = prepareData();
> >> return CompletableFuture.supplyAsync(() -> processor.apply(data),
> >> processingExecutor)
> >> .whenComplete((o, ex) -> cleanup(data));
> >> }
> >>
> >> Callers of this method may just break cleanup, without even knowing
> >> about this, but simply calling complete() on returned CF. But it seems
> >> that nothing can be done about it.
> >>
> >> czw., 20 maj 2021 o 14:40 Benjamin Manes <[email protected]>
> napisał(a):
> >> >
> >> > The timeout is applied to the STEP2 future which completes it
> exceptionally. This causes it to not run the mapping function for a result,
> as it was completed prior to that being triggered. You can use copy() to
> decouple the source from the timeout if you want both to complete; e.g.
> cache a long running value but fail a consumer if taking to long.
> >> >
> >> > On Thu, May 20, 2021 at 5:31 AM Michał Górniewski via
> Concurrency-interest <[email protected]> wrote:
> >> >>
> >> >> With following code:
> >> >>
> >> >> ------------------------------
> >> >> public static void main(String[] args) {
> >> >> CompletableFuture<String> initial = new CompletableFuture<>();
> >> >>
> >> >> CompletableFuture<String> withSteps = initial
> >> >> .thenApply(s -> {
> >> >> System.out.println("STEP1: " + s);
> >> >> return s;
> >> >> })
> >> >> .thenApply(s -> {
> >> >> System.out.println("STEP2: " + s);
> >> >> return s;
> >> >> });
> >> >>
> >> >> CompletableFuture<String> timeout = withSteps;
> >> >>
> >> >> timeout.whenComplete((s, throwable) ->
> >> >> System.out.println("TIMEOUT: " + throwable));
> >> >>
> >> >> timeout.completeExceptionally(new RuntimeException("TIMEOUT"));
> >> >> initial.complete("SUCCESS");
> >> >> }
> >> >> ------------------------------
> >> >>
> >> >> I got output like this:
> >> >>
> >> >> ------------------------------
> >> >> TIMEOUT: java.lang.RuntimeException: TIMEOUT
> >> >> STEP1: SUCCESS
> >> >> ------------------------------
> >> >>
> >> >> I don't understand why "STEP2" is not executed in this case. Is this
> expected?
> >> >> For me this API is confusing and may create issues like this:
> >> >> https://github.com/resilience4j/resilience4j/issues/1427
> >> >>
> >> >> Thanks,
> >> >> Michał Górniewski
> >> >> _______________________________________________
> >> >> Concurrency-interest mailing list
> >> >> [email protected]
> >> >> http://cs.oswego.edu/mailman/listinfo/concurrency-interest
> >> _______________________________________________
> >> Concurrency-interest mailing list
> >> [email protected]
> >> http://cs.oswego.edu/mailman/listinfo/concurrency-interest
> >
> >
> >
> > --
> > Cheers,
> > √
>
--
Cheers,
√
_______________________________________________
Concurrency-interest mailing list
[email protected]
http://cs.oswego.edu/mailman/listinfo/concurrency-interest