Flux to List<Objects> without blocking

Viewed 11529

Looking for converting Flux to List<Object>. Getting error if I use block(). So, need to conver without blocking calls.

Flux.from(Collection.find())

Using reactive programming, but graphql expects List<objects> and erroring with returning Flux.

Code with Block()

public List<Test> findAll() {
        return Flux.from(testCollection.find()).collectList().block();

}

Error :-

block()/blockFirst()/blockLast() are blocking, which is not supported in thread reactor-http-kqueue-7

Here, I need to return List<Test> as I can not send Flux<Test> for some reason.

3 Answers

As stated in the comments, You can't. The reactive pattern is to stay in a flow.

So,

Mono<GraphqlResponse> = Flux.just("A", "B" "C")
  .collectList()
  .map(this::someMethod);

GraphqlResponse someMethod(List<String> abcs) {
    return graphQl.doSomething(abcs);
}

I assume the Collection class is from some reactor library doing some tcp/http call? collection.find returns a Flux/Mono/Publisher I assume? Then thats not because of collectList doesn't allow you to do that, thats because you are trying to run block on a thread which is a NonBlocking thread, I assume collection.find publish the elements on a thread which is an instance of NonBlocking, from its name reactor-http-kqueue-7 which is probably a netty thread.

You can take a look at BlockingSingleSubscriber.blockingGet method which tells you why

    final T blockingGet() {
        if (Schedulers.isInNonBlockingThread()) {
            throw new IllegalStateException("block()/blockFirst()/blockLast() are blocking, which is not supported in thread " + Thread.currentThread().getName());
        }
...

If you have to get the result from the caller thread (which calls Flux.from), then you can do

Flux.from(testCondition.find())
.publishOn(Schedulers.boundedElastic())
.collectList()
.block()

following example is converting flux<Object> to List<Object> but remember that converting a Flux to a List makes the whole thing NOT reactive:

public static void main(String[] args) {
    Flux<String> flux = Flux.just("test1", "test2", "test3");
    List<String> list = new ArrayList<>();
    flux.collectList().subscribe(list::addAll);
    list.forEach(System.out::println);
}
Related