diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveResultSet.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveResultSet.java
index 103eff672..838eeba8d 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveResultSet.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveResultSet.java
@@ -15,10 +15,10 @@
*/
package org.springframework.data.cassandra;
-import java.util.List;
-
import reactor.core.publisher.Flux;
+import java.util.List;
+
import com.datastax.driver.core.ColumnDefinitions;
import com.datastax.driver.core.ExecutionInfo;
import com.datastax.driver.core.Row;
@@ -48,20 +48,31 @@ import com.datastax.driver.core.Row;
public interface ReactiveResultSet {
/**
- * Returns a {@link Flux} over the rows contained in this result set.
+ * Returns a {@link Flux} over the rows contained in this result set applying transparent paging.
*
* The {@link Flux} will stream over all records that in this {@link ReactiveResultSet} according to the reactive
- * demand.
+ * demand and fetch next result chunks by issuing the underlying query with the current
+ * {@link com.datastax.driver.core.PagingState} applied.
*
*
- * @return a {@link Flux} of rows that will stream over all {@link Row rows} in this {@link ReactiveResultSet}.
+ * @return a {@link Flux} of rows that will stream over all {@link Row rows} of the entire result.
*/
Flux rows();
/**
- * Returns the columns returned in this ResultSet.
+ * Returns a {@link Flux} over the rows contained in this result set chunk. This method does not apply transparent
+ * paging. Use {@link com.datastax.driver.core.PagingState} from {@link #getExecutionInfo()} to issue subsequent
+ * queries to obtain the next result chunk.
*
- * @return the columns returned in this ResultSet.
+ * @return a {@link Flux} of rows that will stream over all {@link Row rows} in this {@link ReactiveResultSet}.
+ * @since 2.1
+ */
+ Flux availableRows();
+
+ /**
+ * Returns the columns returned in this {@link ReactiveResultSet}.
+ *
+ * @return the columns returned in this {@link ReactiveResultSet}.
*/
ColumnDefinitions getColumnDefinitions();
@@ -82,7 +93,7 @@ public interface ReactiveResultSet {
boolean wasApplied();
/**
- * Returns information on the execution of the last query made for this result set.
+ * Returns information on the execution of the last query made for this {@link ReactiveResultSet}.
*
* Note that in most cases, a result set is fetched with only one query, but large result sets can be paged and thus
* be retrieved by multiple queries. In that case this method return the {@link ExecutionInfo} for the last query
@@ -91,18 +102,18 @@ public interface ReactiveResultSet {
* The returned object includes basic information such as the queried hosts, but also the Cassandra query trace if
* tracing was enabled for the query.
*
- * @return the execution info for the last query made for this result set.
+ * @return the {@link ExecutionInfo} for the last query made for this {@link ReactiveResultSet}.
*/
ExecutionInfo getExecutionInfo();
/**
- * Return the execution information for all queries made to retrieve this result set.
+ * Return the execution information for all queries made to retrieve this {@link ReactiveResultSet}.
*
* Unless the result set is large enough to get paged underneath, the returned list will be singleton. If paging has
* been used however, the returned list contains the {@link ExecutionInfo} objects for all the queries done to obtain
* this result set (at the time of the call) in the order those queries were made.
*
- * @return a list of the execution info for all the queries made for this result set.
+ * @return a list of the {@link ExecutionInfo} for all the queries made for this {@link ReactiveResultSet}.
*/
List getAllExecutionInfo();
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/QueryUtils.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/QueryUtils.java
index 0792ecf8a..6412777c5 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/QueryUtils.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/QueryUtils.java
@@ -16,6 +16,7 @@
package org.springframework.data.cassandra.core;
import java.util.ArrayList;
+import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -37,11 +38,12 @@ import org.springframework.data.domain.PageRequest;
import org.springframework.data.domain.Pageable;
import org.springframework.data.domain.Slice;
import org.springframework.data.domain.SliceImpl;
+import org.springframework.lang.Nullable;
import org.springframework.util.Assert;
-import reactor.core.publisher.Mono;
import com.datastax.driver.core.PagingState;
import com.datastax.driver.core.ResultSet;
+import com.datastax.driver.core.Row;
import com.datastax.driver.core.Statement;
import com.datastax.driver.core.querybuilder.Delete;
import com.datastax.driver.core.querybuilder.Delete.Where;
@@ -49,6 +51,7 @@ import com.datastax.driver.core.querybuilder.Insert;
import com.datastax.driver.core.querybuilder.QueryBuilder;
import com.datastax.driver.core.querybuilder.Select;
import com.datastax.driver.core.querybuilder.Update;
+import com.google.common.collect.Iterators;
/**
* Simple utility class for working with the QueryBuilder API using mapped entities.
@@ -177,15 +180,34 @@ class QueryUtils {
int toRead = resultSet.getAvailableWithoutFetching();
- List result = new ArrayList<>(toRead);
+ return readSlice(() -> Iterators.limit(resultSet.iterator(), toRead), resultSet.getExecutionInfo().getPagingState(),
+ mapper, page, pageSize);
+ }
- for (int index = 0; index < toRead; index++) {
- T element = mapper.mapRow(resultSet.one(), index);
+ /**
+ * Read a {@link Slice} of data from the {@link Iterable} of {@link Row}s for a {@link Pageable}.
+ *
+ * @param rows must not be {@literal null}.
+ * @param pagingState
+ * @param mapper must not be {@literal null}.
+ * @param page
+ * @param pageSize
+ * @return the resulting {@link Slice}.
+ * @since 2.1
+ */
+ static Slice readSlice(Iterable rows, @Nullable PagingState pagingState, RowMapper mapper, int page,
+ int pageSize) {
+
+ List result = new ArrayList<>(pageSize);
+
+ Iterator iterator = rows.iterator();
+ int index = 0;
+
+ while (iterator.hasNext()) {
+ T element = mapper.mapRow(iterator.next(), index++);
result.add(element);
}
- PagingState pagingState = resultSet.getExecutionInfo().getPagingState();
-
CassandraPageRequest pageRequest = CassandraPageRequest.of(PageRequest.of(page, pageSize), pagingState);
return new SliceImpl<>(result, pageRequest, pagingState != null);
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraOperations.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraOperations.java
index fb9495653..59e7f6114 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraOperations.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraOperations.java
@@ -23,11 +23,11 @@ import org.springframework.data.cassandra.core.convert.CassandraConverter;
import org.springframework.data.cassandra.core.cql.QueryOptions;
import org.springframework.data.cassandra.core.cql.ReactiveCqlOperations;
import org.springframework.data.cassandra.core.cql.WriteOptions;
+import org.springframework.data.cassandra.core.query.CassandraPageRequest;
import org.springframework.data.cassandra.core.query.Query;
import org.springframework.data.cassandra.core.query.Update;
import org.springframework.data.domain.Slice;
-
import com.datastax.driver.core.Statement;
/**
@@ -35,6 +35,7 @@ import com.datastax.driver.core.Statement;
* Not often used directly, but a useful option to enhance testability, as it can easily be mocked or stubbed.
*
* @author Mark Paluch
+ * @author Hleb Albau
* @since 2.0
* @see ReactiveCassandraTemplate
* @see ReactiveCqlOperations
@@ -100,15 +101,14 @@ public interface ReactiveCassandraOperations extends ReactiveFluentCassandraOper
Flux select(Statement statement, Class entityClass) throws DataAccessException;
/**
- * Execute a {@code SELECT} query with paging and convert the result set to a {@link Slice} of entities.
- *
- * A sliced query translates the effective {@link Statement#getFetchSize() fetch size} to the page size.
+ * Execute a {@code SELECT} query with paging and convert the result set to a {@link Slice} of entities. A sliced
+ * query translates the effective {@link Statement#getFetchSize() fetch size} to the page size.
*
* @param statement the CQL statement, must not be {@literal null}.
* @param entityClass The entity type must not be {@literal null}.
- * @return the result object returned by the action or {@link Mono#empty()}
+ * @return the result object returned by the action or {@link Mono#just(Object)} of an empty {@link Slice}.
* @throws DataAccessException if there is any problem executing the query.
- * @since 2.0
+ * @since 2.1
*/
Mono> slice(Statement statement, Class entityClass) throws DataAccessException;
@@ -136,12 +136,24 @@ public interface ReactiveCassandraOperations extends ReactiveFluentCassandraOper
*/
Flux select(Query query, Class entityClass) throws DataAccessException;
+ /**
+ * Execute a {@code SELECT} query with paging and convert the result set to a {@link Slice} of entities.
+ *
+ * @param query the query object used to create a CQL statement, must not be {@literal null}.
+ * @param entityClass The entity type must not be {@literal null}.
+ * @return the result object returned by the action or {@link Mono#just(Object)} of an empty {@link Slice}.
+ * @throws DataAccessException if there is any problem executing the query.
+ * @since 2.1
+ * @see CassandraPageRequest
+ */
+ Mono> slice(Query query, Class entityClass) throws DataAccessException;
+
/**
* Execute a {@code SELECT} query and convert the resulting item to an entity.
*
* @param query must not be {@literal null}.
* @param entityClass The entity type must not be {@literal null}.
- * @return the result object returned by the action or {@link Mono#empty()}
+ * @return the result object returned by the action or {@link Mono#empty()}.
* @throws DataAccessException if there is any problem issuing the execution.
*/
Mono selectOne(Query query, Class entityClass) throws DataAccessException;
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java
index 6ccca5a05..b82b8674c 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java
@@ -15,13 +15,11 @@
*/
package org.springframework.data.cassandra.core;
-import java.util.function.Function;
-
-import lombok.NonNull;
import lombok.Value;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
+import java.util.Collections;
import java.util.function.Function;
import org.reactivestreams.Publisher;
@@ -43,6 +41,7 @@ import org.springframework.data.cassandra.core.cql.QueryOptions;
import org.springframework.data.cassandra.core.cql.ReactiveCqlOperations;
import org.springframework.data.cassandra.core.cql.ReactiveCqlTemplate;
import org.springframework.data.cassandra.core.cql.ReactiveSessionCallback;
+import org.springframework.data.cassandra.core.cql.RowMapper;
import org.springframework.data.cassandra.core.cql.WriteOptions;
import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessionFactory;
import org.springframework.data.cassandra.core.mapping.CassandraMappingContext;
@@ -56,6 +55,7 @@ import org.springframework.data.cassandra.core.mapping.event.BeforeDeleteEvent;
import org.springframework.data.cassandra.core.mapping.event.BeforeSaveEvent;
import org.springframework.data.cassandra.core.query.Query;
import org.springframework.data.domain.Slice;
+import org.springframework.data.domain.SliceImpl;
import org.springframework.data.mapping.context.MappingContext;
import org.springframework.data.projection.ProjectionFactory;
import org.springframework.data.projection.SpelAwareProxyProjectionFactory;
@@ -68,7 +68,6 @@ import com.datastax.driver.core.Row;
import com.datastax.driver.core.Session;
import com.datastax.driver.core.SimpleStatement;
import com.datastax.driver.core.Statement;
-import com.datastax.driver.core.Row;
import com.datastax.driver.core.exceptions.DriverException;
import com.datastax.driver.core.querybuilder.Delete;
import com.datastax.driver.core.querybuilder.Insert;
@@ -92,6 +91,7 @@ import com.datastax.driver.core.querybuilder.Update;
* @author Mark Paluch
* @author John Blum
* @author Lukasz Antoniak
+ * @author Hleb Albau
* @since 2.0
*/
public class ReactiveCassandraTemplate implements ReactiveCassandraOperations, ApplicationEventPublisherAware {
@@ -239,19 +239,28 @@ public class ReactiveCassandraTemplate implements ReactiveCassandraOperations, A
/* (non-Javadoc)
* @see org.springframework.data.cassandra.core.CassandraOperations#slice(com.datastax.driver.core.Statement, java.lang.Class)
*/
- @Override
- public Mono> slice(Statement statement, Class entityClass) {
+ @Override
+ public Mono> slice(Statement statement, Class entityClass) {
- Assert.notNull(statement, "Statement must not be null");
- Assert.notNull(entityClass, "Entity type must not be null");
+ Assert.notNull(statement, "Statement must not be null");
+ Assert.notNull(entityClass, "Entity type must not be null");
- Mono resultSetMono = getReactiveCqlOperations().queryForResultSet(statement);
- Mono effectiveFetchSizeMono = getEffectiveFetchSize(statement);
- Function rowMapper = (row) -> getConverter().read(entityClass, row);
+ Mono resultSetMono = getReactiveCqlOperations().queryForResultSet(statement);
+ Mono effectiveFetchSizeMono = getEffectiveFetchSize(statement);
+ RowMapper rowMapper = (row, i) -> getConverter().read(entityClass, row);
- return Mono.zip(resultSetMono, effectiveFetchSizeMono)
- .flatMap(tuple -> QueryUtils.readSlice(tuple.getT1(), tuple.getT2(), rowMapper));
- }
+ return resultSetMono.zipWith(effectiveFetchSizeMono).flatMap(tuple -> {
+
+ ReactiveResultSet resultSet = tuple.getT1();
+ Integer effectiveFetchSize = tuple.getT2();
+
+ return resultSet.availableRows().collectList().map(it -> {
+ return QueryUtils.readSlice(it, resultSet.getExecutionInfo().getPagingState(), rowMapper, 1,
+ effectiveFetchSize);
+ });
+
+ }).defaultIfEmpty(new SliceImpl<>(Collections.emptyList()));
+ }
/* (non-Javadoc)
* @see org.springframework.data.cassandra.core.ReactiveCassandraOperations#selectOne(com.datastax.driver.core.Statement, java.lang.Class)
@@ -286,6 +295,20 @@ public class ReactiveCassandraTemplate implements ReactiveCassandraOperations, A
return getReactiveCqlOperations().query(select, (row, rowNum) -> mapper.apply(row));
}
+ /* (non-Javadoc)
+ * @see org.springframework.data.cassandra.core.ReactiveCassandraOperations#slice(org.springframework.data.cassandra.core.query.Query, java.lang.Class)
+ */
+ @Override
+ public Mono> slice(Query query, Class entityClass) throws DataAccessException {
+
+ Assert.notNull(query, "Query must not be null");
+ Assert.notNull(entityClass, "Entity type must not be null");
+
+ RegularStatement select = getStatementFactory().select(query, getRequiredPersistentEntity(entityClass));
+
+ return slice(select, entityClass);
+ }
+
/* (non-Javadoc)
* @see org.springframework.data.cassandra.core.ReactiveCassandraOperations#selectOne(org.springframework.data.cassandra.core.query.Query, java.lang.Class)
*/
@@ -683,27 +706,26 @@ public class ReactiveCassandraTemplate implements ReactiveCassandraOperations, A
return converter;
}
- @SuppressWarnings("ConstantConditions")
- private Mono getEffectiveFetchSize(Statement statement) {
+ @SuppressWarnings("ConstantConditions")
+ private Mono getEffectiveFetchSize(Statement statement) {
- if (statement.getFetchSize() > 0) {
- return Mono.just(statement.getFetchSize());
- }
+ if (statement.getFetchSize() > 0) {
+ return Mono.just(statement.getFetchSize());
+ }
- if (getReactiveCqlOperations() instanceof CassandraAccessor) {
- CassandraAccessor accessor = (CassandraAccessor) getReactiveCqlOperations();
- if (accessor.getFetchSize() != -1) {
- return Mono.just(accessor.getFetchSize());
- }
- }
+ if (getReactiveCqlOperations() instanceof CassandraAccessor) {
+ CassandraAccessor accessor = (CassandraAccessor) getReactiveCqlOperations();
+ if (accessor.getFetchSize() != -1) {
+ return Mono.just(accessor.getFetchSize());
+ }
+ }
- return getReactiveCqlOperations().execute((ReactiveSessionCallback) session ->
- Mono.fromSupplier(() -> session.getCluster().getConfiguration().getQueryOptions().getFetchSize())
- ).single();
- }
+ return getReactiveCqlOperations().execute((ReactiveSessionCallback) session -> Mono
+ .just(session.getCluster().getConfiguration().getQueryOptions().getFetchSize())).single();
+ }
- @Value
- static class StatementCallback implements ReactiveSessionCallback, CqlProvider {
+ @Value
+ static class StatementCallback implements ReactiveSessionCallback, CqlProvider {
@lombok.NonNull Statement statement;
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java
index de5d94420..ec11bf7fa 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java
@@ -15,33 +15,23 @@
*/
package org.springframework.data.cassandra.core.cql.session;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ExecutionException;
-
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.publisher.MonoProcessor;
import reactor.core.publisher.MonoSink;
import reactor.core.scheduler.Scheduler;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ExecutionException;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import org.springframework.data.cassandra.ReactiveResultSet;
import org.springframework.data.cassandra.ReactiveSession;
import org.springframework.util.Assert;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import com.datastax.driver.core.Cluster;
-import com.datastax.driver.core.ColumnDefinitions;
-import com.datastax.driver.core.ExecutionInfo;
-import com.datastax.driver.core.PreparedStatement;
-import com.datastax.driver.core.RegularStatement;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.Row;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.SimpleStatement;
-import com.datastax.driver.core.Statement;
+import com.datastax.driver.core.*;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
@@ -252,6 +242,7 @@ public class DefaultBridgedReactiveSession implements ReactiveSession {
this.resultSet = resultSet;
}
+
/* (non-Javadoc)
* @see org.springframework.data.cassandra.ReactiveResultSet#rows()
*/
@@ -260,7 +251,15 @@ public class DefaultBridgedReactiveSession implements ReactiveSession {
return getRows(Mono.just(this.resultSet));
}
- Flux getRows(Mono nextResults) {
+ /* (non-Javadoc)
+ * @see org.springframework.data.cassandra.ReactiveResultSet#availableRows()
+ */
+ @Override
+ public Flux availableRows() {
+ return toRows(this.resultSet);
+ }
+
+ private Flux getRows(Mono nextResults) {
return nextResults.flatMapMany(it -> {
@@ -280,7 +279,7 @@ public class DefaultBridgedReactiveSession implements ReactiveSession {
static Flux toRows(ResultSet resultSet) {
- int prefetch = Math.max(1, resultSet.getAvailableWithoutFetching());
+ int prefetch = Math.max(0, resultSet.getAvailableWithoutFetching());
return Flux.fromIterable(resultSet).take(prefetch);
}
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java
index 1358a6348..cc60175ab 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java
@@ -15,7 +15,6 @@
*/
package org.springframework.data.cassandra.repository.query;
-import org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution.SlicedExecution;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@@ -30,6 +29,7 @@ import org.springframework.data.cassandra.repository.query.ReactiveCassandraQuer
import org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution.ResultProcessingConverter;
import org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution.ResultProcessingExecution;
import org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution.SingleEntityExecution;
+import org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution.SlicedExecution;
import org.springframework.data.repository.query.ParameterAccessor;
import org.springframework.data.repository.query.RepositoryQuery;
import org.springframework.data.repository.query.ResultProcessor;
@@ -41,6 +41,7 @@ import com.datastax.driver.core.Statement;
* Base class for reactive {@link RepositoryQuery} implementations for Cassandra.
*
* @author Mark Paluch
+ * @author Hleb Albau
* @see org.springframework.data.cassandra.repository.query.CassandraRepositoryQuerySupport
* @since 2.0
*/
@@ -48,17 +49,6 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository
private final ReactiveCassandraOperations operations;
- private static CassandraConverter toConverter(ReactiveCassandraOperations operations) {
-
- Assert.notNull(operations, "ReactiveCassandraOperations must not be null");
-
- return operations.getConverter();
- }
-
- private static CassandraMappingContext toMappingContext(ReactiveCassandraOperations operations) {
- return toConverter(operations).getMappingContext();
- }
-
/**
* Create a new {@link AbstractReactiveCassandraQuery} from the given {@link CassandraQueryMethod} and
* {@link CassandraOperations}.
@@ -68,15 +58,11 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository
*/
public AbstractReactiveCassandraQuery(ReactiveCassandraQueryMethod method, ReactiveCassandraOperations operations) {
- super(method, toMappingContext(operations));
+ super(method, getRequiredMappingContext(operations));
this.operations = operations;
}
- protected ReactiveCassandraOperations getReactiveCassandraOperations() {
- return this.operations;
- }
-
/*
* (non-Javadoc)
* @see org.springframework.data.repository.query.RepositoryQuery#getQueryMethod()
@@ -93,34 +79,31 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository
@Override
public Object execute(Object[] parameters) {
- return getQueryMethod().hasReactiveWrapperParameter()
- ? executeDeferred(parameters)
- : executeNow(parameters);
+ return getQueryMethod().hasReactiveWrapperParameter() ? executeDeferred(parameters) : executeNow(parameters);
}
@SuppressWarnings("unchecked")
private Object executeDeferred(Object[] parameters) {
- return getQueryMethod().isCollectionQuery()
- ? Flux.defer(() -> (Publisher