When ints is given:
List<Integer> ints = IntStream.range(0, 1000).boxed().collect(Collectors.toList());
With Java Stream API, we can reduce them
MyValue myvalue = ints
.parallelStream()
.map(x -> toMyValue(x))
.reduce((t, t2) -> t.combine(t2))
.get();
In this example, what important to me are...
- items will be reduced in multiple threads
- early mapped items will be reduced early
- not all result of
toMyValue()will be loaded at the same time
Now I want to do same processing by CompletableFuture API.
To do map, I did:
List<CompeletableFuture<MyValue>> myValueFutures = ints
.stream()
.map(x -> CompletableFuture.supplyAsync(() -> toMyValue(x), MY_THREAD_POOL))
.collect(Collectors.toList());
And now I have no idea how to reduce List<CompeletableFuture<MyValue>> myValueFutures to get single MyValue.
Parallel stream provides convenient APIs but because of these problems I want not to use Stream API:
- parallel stream is hard to stop the stage during processing.
- parallel stream's active worker count can exceed the parallelism when some workers blocked by IO. This helps to maximize cpu utilization but memory overhead can occur(even OOM).
Any way to reduce CompetableFutures? one by one with out stream reduce api?