From c843293fa52cca77d6d0860845420167664903ee Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Mon, 18 Jun 2018 11:33:20 +0200 Subject: [PATCH] DATACASS-529 - Polishing. Introduce ReactiveResultSet.availableRows() to fetch rows without transparent paging. Refactor QueryUtils to extract a Slice from an Iterable. Add slice(Query, Class) to ReactiveCassandraOperations to expose a consistent API. Convert space indentation to tab indentation. Add tests. Add since tags. Add documentation. Reformat code. Original pull request: #128. --- .../data/cassandra/ReactiveResultSet.java | 33 ++-- .../data/cassandra/core/QueryUtils.java | 34 +++- .../core/ReactiveCassandraOperations.java | 26 +++- .../core/ReactiveCassandraTemplate.java | 84 ++++++---- .../DefaultBridgedReactiveSession.java | 37 +++-- .../query/AbstractReactiveCassandraQuery.java | 67 ++++---- .../ReactiveCassandraQueryExecution.java | 4 +- ...tiveCassandraTemplateIntegrationTests.java | 147 +++++++++++++++++- ...efaultBridgedReactiveSessionUnitTests.java | 64 ++++---- ...veCassandraRepositoryIntegrationTests.java | 16 +- src/main/asciidoc/new-features.adoc | 3 +- .../reference/cassandra-repositories.adoc | 4 +- .../reactive-cassandra-repositories.adoc | 29 ++++ 13 files changed, 397 insertions(+), 151 deletions(-) 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) execute(parameters)) + return getQueryMethod().isCollectionQuery() ? Flux.defer(() -> (Publisher) execute(parameters)) : Mono.defer(() -> (Mono) execute(parameters)); } private Object executeNow(Object[] parameters) { - ReactiveCassandraParameterAccessor parameterAccessor = - new ReactiveCassandraParameterAccessor(getQueryMethod(), parameters); + ReactiveCassandraParameterAccessor parameterAccessor = new ReactiveCassandraParameterAccessor(getQueryMethod(), + parameters); CassandraParameterAccessor convertingParameterAccessor = new ConvertingParameterAccessor( - toConverter(getReactiveCassandraOperations()), parameterAccessor); + getRequiredConverter(getReactiveCassandraOperations()), parameterAccessor); Statement statement = createQuery(convertingParameterAccessor); ResultProcessor resultProcessor = getQueryMethod().getResultProcessor() .withDynamicProjection(convertingParameterAccessor); - ReactiveCassandraQueryExecution queryExecution = getExecution(parameterAccessor,new ResultProcessingConverter(resultProcessor, - toMappingContext(getReactiveCassandraOperations()), getEntityInstantiators())); + ReactiveCassandraQueryExecution queryExecution = getExecution(parameterAccessor, new ResultProcessingConverter( + resultProcessor, getRequiredMappingContext(getReactiveCassandraOperations()), getEntityInstantiators())); Class resultType = resolveResultType(resultProcessor); @@ -130,7 +113,7 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository private Class resolveResultType(ResultProcessor resultProcessor) { CassandraReturnedType returnedType = new CassandraReturnedType(resultProcessor.getReturnedType(), - toConverter(getReactiveCassandraOperations()).getCustomConversions()); + getRequiredConverter(getReactiveCassandraOperations()).getCustomConversions()); return (returnedType.isProjecting() ? returnedType.getDomainType() : returnedType.getReturnedType()); } @@ -142,25 +125,30 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository */ protected abstract Statement createQuery(CassandraParameterAccessor accessor); + protected ReactiveCassandraOperations getReactiveCassandraOperations() { + return this.operations; + } + /** * Returns the execution instance to use. + * * @param parameterAccessor must not be {@literal null}. - * @param resultProcessing must not be {@literal null}. @return + * @param resultProcessing must not be {@literal null}. */ - private ReactiveCassandraQueryExecution getExecution(CassandraParameterAccessor parameterAccessor, + private ReactiveCassandraQueryExecution getExecution(ReactiveCassandraParameterAccessor parameterAccessor, Converter resultProcessing) { return new ResultProcessingExecution(getExecutionToWrap(parameterAccessor), resultProcessing); } - private ReactiveCassandraQueryExecution getExecutionToWrap() { + private ReactiveCassandraQueryExecution getExecutionToWrap(CassandraParameterAccessor parameterAccessor) { if (getQueryMethod().isSliceQuery()) { return new SlicedExecution(getReactiveCassandraOperations(), parameterAccessor.getPageable()); - }else if (getQueryMethod().isCollectionQuery()) { + } else if (getQueryMethod().isCollectionQuery()) { return new CollectionExecution(getReactiveCassandraOperations()); } else if (isCountQuery()) { - return ((statement, type) -> - new SingleEntityExecution(getReactiveCassandraOperations(), false).execute(statement, Long.class)); + return ((statement, type) -> new SingleEntityExecution(getReactiveCassandraOperations(), false).execute(statement, + Long.class)); } else if (isExistsQuery()) { return new ExistsExecution(getReactiveCassandraOperations()); } else { @@ -191,4 +179,15 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository * @since 2.0.4 */ protected abstract boolean isLimiting(); + + private static CassandraConverter getRequiredConverter(ReactiveCassandraOperations operations) { + + Assert.notNull(operations, "ReactiveCassandraOperations must not be null"); + + return operations.getConverter(); + } + + private static CassandraMappingContext getRequiredMappingContext(ReactiveCassandraOperations operations) { + return getRequiredConverter(operations).getMappingContext(); + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java index 00796ccbb..16e09651a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java @@ -29,7 +29,6 @@ import org.springframework.data.cassandra.core.mapping.CassandraPersistentProper import org.springframework.data.cassandra.core.query.CassandraPageRequest; import org.springframework.data.convert.EntityInstantiators; import org.springframework.data.domain.Pageable; -import org.springframework.data.domain.Slice; import org.springframework.data.mapping.context.MappingContext; import org.springframework.data.repository.query.ResultProcessor; import org.springframework.data.repository.query.ReturnedType; @@ -50,9 +49,10 @@ interface ReactiveCassandraQueryExecution { Object execute(Statement statement, Class type); /** - * {@link CassandraQueryExecution} for a {@link Slice}. + * {@link ReactiveCassandraQueryExecution} for a {@link org.springframework.data.domain.Slice}. * * @author Hleb Albau + * @since 2.1 */ @RequiredArgsConstructor final class SlicedExecution implements ReactiveCassandraQueryExecution { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java index e0fdce7e4..6d8355b24 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java @@ -15,20 +15,28 @@ */ package org.springframework.data.cassandra.core; -import static org.assertj.core.api.Assertions.assertThat; -import static org.springframework.data.cassandra.core.query.Criteria.where; +import static org.assertj.core.api.Assertions.*; +import static org.springframework.data.cassandra.core.query.Criteria.*; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; import reactor.test.StepVerifier; import reactor.test.StepVerifier.FirstStep; +import java.lang.reflect.UndeclaredThrowableException; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +import java.util.function.Predicate; + import org.junit.Before; import org.junit.Test; - import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; import org.springframework.data.cassandra.core.cql.ReactiveCqlTemplate; import org.springframework.data.cassandra.core.cql.session.DefaultBridgedReactiveSession; +import org.springframework.data.cassandra.core.query.CassandraPageRequest; import org.springframework.data.cassandra.core.query.Columns; import org.springframework.data.cassandra.core.query.Criteria; import org.springframework.data.cassandra.core.query.Query; @@ -37,8 +45,14 @@ import org.springframework.data.cassandra.domain.User; import org.springframework.data.cassandra.domain.UserToken; import org.springframework.data.cassandra.repository.support.SchemaTestUtils; import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTest; +import org.springframework.data.domain.Slice; import org.springframework.data.domain.Sort; +import org.springframework.util.Assert; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.Host; +import com.datastax.driver.core.LatencyTracker; +import com.datastax.driver.core.Statement; import com.datastax.driver.core.utils.UUIDs; /** @@ -294,7 +308,134 @@ public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceC assertThat(template.selectOne(query, UserToken.class).block()).isEqualTo(token1); } + @Test // DATACASS-529 + public void pagedSelectShouldIssueMultipleStatements() { + + Set expectedIds = new LinkedHashSet<>(); + + for (int count = 0; count < 100; count++) { + User user = new User("heisenberg" + count, "Walter", "White"); + expectedIds.add(user.getId()); + template.insert(user).as(StepVerifier::create).expectNextCount(1).verifyComplete(); + } + + QueryListener listener = new QueryListener(); + this.cluster.register(listener); + + Query query = Query.empty().pageRequest(CassandraPageRequest.first(10)); + + template.select(query, User.class).as(StepVerifier::create).expectNextCount(100).verifyComplete(); + + listener.await(it -> it.size() == 11); + assertThat(listener.statements).hasSize(11); + + this.cluster.unregister(listener); + } + + @Test // DATACASS-529 + public void shouldIssueSinglePageRequestForSlice() { + + Set expectedIds = new LinkedHashSet<>(); + + for (int count = 0; count < 100; count++) { + User user = new User("heisenberg" + count, "Walter", "White"); + expectedIds.add(user.getId()); + template.insert(user).as(StepVerifier::create).expectNextCount(1).verifyComplete(); + } + + QueryListener listener = new QueryListener(); + this.cluster.register(listener); + + Query query = Query.empty().pageRequest(CassandraPageRequest.first(10)); + + Mono> slice = template.slice(query, User.class); + + slice.as(StepVerifier::create).consumeNextWith(it -> { + assertThat(it).hasSize(10); + }).verifyComplete(); + + listener.await(it -> it.size() == 1); + assertThat(listener.statements).hasSize(1); + + this.cluster.unregister(listener); + } + + @Test // DATACASS-529 + public void shouldReturnEmptySliceOnEmptyResult() { + + Query query = Query.query(where("id").is("foo")).pageRequest(CassandraPageRequest.first(10)); + + Mono> slice = template.slice(query, User.class); + + slice.as(StepVerifier::create).consumeNextWith(it -> { + assertThat(it).isEmpty(); + }).verifyComplete(); + } + private FirstStep verifyUser(String userId) { return StepVerifier.create(template.selectOneById(userId, User.class)); } + + static class QueryListener implements LatencyTracker { + + private List statements = new CopyOnWriteArrayList<>(); + + @Override + public void update(Host host, Statement statement, Exception exception, long newLatencyNanos) { + statements.add(statement); + } + + @Override + public void onRegister(Cluster cluster) {} + + @Override + public void onUnregister(Cluster cluster) {} + + /** + * Await until {@link Predicate} yields {@literal true}. Waits up to {@literal 10 SECONDS} + * + * @param predicate must not be {@literal null}. + * @throws IllegalStateException if the timeout exceeds. + * @throws UndeclaredThrowableException in case of {@link InterruptedException}. + */ + public void await(Predicate> predicate) { + await(predicate, 10, TimeUnit.SECONDS); + } + + /** + * Await until {@link Predicate} yields {@literal true}. The predicate is eagerly evaluated. If the predicate yields + * {@literal false}, micro-waits of {@code 100ms} are applied. + * + * @param predicate must not be {@literal null}. + * @param timeout + * @param unit must not be {@literal null}. + * @throws IllegalStateException if the timeout exceeds. + * @throws UndeclaredThrowableException in case of {@link InterruptedException}. + */ + public void await(Predicate> predicate, long timeout, TimeUnit unit) { + + Assert.notNull(predicate, "Predicate must not be null"); + Assert.notNull(unit, "TimeUnit must not be null"); + + long waitedNs = 0; + long timeoutMs = unit.toNanos(timeout); + long waitSegmentMs = TimeUnit.MILLISECONDS.toMillis(100); + + while (!predicate.test(statements)) { + + try { + Thread.sleep(waitSegmentMs); + waitedNs += TimeUnit.MILLISECONDS.toNanos(waitSegmentMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new UndeclaredThrowableException(e); + } + + if (waitedNs > timeoutMs) { + throw new IllegalStateException( + String.format("Timeout: Condition did not evaluate to true within %d %s!", timeout, unit)); + } + } + } + } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/DefaultBridgedReactiveSessionUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/DefaultBridgedReactiveSessionUnitTests.java index dc0e2de90..106e0fca0 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/DefaultBridgedReactiveSessionUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/DefaultBridgedReactiveSessionUnitTests.java @@ -15,30 +15,23 @@ */ package org.springframework.data.cassandra.core.cql; -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.any; -import static org.mockito.Mockito.doAnswer; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.never; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; +import static org.assertj.core.api.Assertions.*; +import static org.mockito.Mockito.*; + +import reactor.core.publisher.Flux; +import reactor.test.StepVerifier; import java.util.ArrayDeque; import java.util.Collections; import java.util.Iterator; import java.util.Queue; -import reactor.core.publisher.Flux; -import reactor.test.StepVerifier; - import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.ArgumentMatchers; import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; - import org.springframework.data.cassandra.ReactiveResultSet; import org.springframework.data.cassandra.core.cql.session.DefaultBridgedReactiveSession; @@ -75,6 +68,15 @@ public class DefaultBridgedReactiveSessionUnitTests { when(sessionMock.executeAsync(any(Statement.class))).thenReturn(future); when(sessionMock.prepareAsync(any(RegularStatement.class))).thenReturn(preparedStatementFuture); + + doAnswer(invocation -> { + + Runnable listener = invocation.getArgument(0); + + listener.run(); + + return null; + }).when(future).addListener(any(), any()); } @Test // DATACASS-335 @@ -169,15 +171,6 @@ public class DefaultBridgedReactiveSessionUnitTests { when(resultSet.getAvailableWithoutFetching()).thenReturn(10); when(resultSet.iterator()).thenReturn(rows); - doAnswer(invocation -> { - - Runnable listener = invocation.getArgument(0); - - listener.run(); - - return null; - }).when(future).addListener(any(), any()); - when(future.get()).thenReturn(resultSet); when(resultSet.isFullyFetched()).thenReturn(true); @@ -190,6 +183,25 @@ public class DefaultBridgedReactiveSessionUnitTests { verify(resultSet, never()).fetchMoreResults(); } + @Test // DATACASS-529 + public void shouldReadAvailableResults() throws Exception { + + Iterator rows = mockIterator(); + + ResultSet resultSet = mock(ResultSet.class); + when(resultSet.iterator()).thenReturn(rows); + when(resultSet.getAvailableWithoutFetching()).thenReturn(10); + when(future.get()).thenReturn(resultSet); + + Flux flux = reactiveSession.execute(new SimpleStatement("")).flatMapMany(ReactiveResultSet::availableRows); + + StepVerifier.create(flux).expectNextCount(10).verifyComplete(); + + verify(rows, times(10)).next(); + verify(future, times(1)).addListener(any(), any()); + verify(resultSet, never()).fetchMoreResults(); + } + @Test // DATACASS-509 public void shouldFetchMore() throws Exception { @@ -205,15 +217,6 @@ public class DefaultBridgedReactiveSessionUnitTests { when(emptyResultSet.iterator()).thenReturn(Collections.emptyIterator()); when(emptyResultSet.isFullyFetched()).thenReturn(true); - doAnswer(invocation -> { - - Runnable listener = invocation.getArgument(0); - - listener.run(); - - return null; - }).when(future).addListener(any(), any()); - when(future.get()).thenReturn(resultSet); when(resultSet.isFullyFetched()).thenReturn(false, true); when(resultSet.fetchMoreResults()).thenReturn(Futures.immediateFuture(emptyResultSet)); @@ -239,6 +242,7 @@ public class DefaultBridgedReactiveSessionUnitTests { when(resultSet.getAvailableWithoutFetching()).thenReturn(10); when(resultSet.iterator()).thenReturn(rows); + reset(future); doAnswer(invocation -> { runnables.offer(invocation.getArgument(0)); return null; diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ReactiveCassandraRepositoryIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ReactiveCassandraRepositoryIntegrationTests.java index c8d8bbcae..1b6bffb72 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ReactiveCassandraRepositoryIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ReactiveCassandraRepositoryIntegrationTests.java @@ -15,18 +15,17 @@ */ package org.springframework.data.cassandra.repository; -import java.util.Arrays; -import java.util.HashSet; -import java.util.Set; - import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.test.StepVerifier; +import java.util.Arrays; +import java.util.HashSet; +import java.util.Set; + import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; - import org.reactivestreams.Publisher; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanClassLoaderAware; @@ -49,6 +48,7 @@ import org.springframework.data.domain.Slice; import org.springframework.data.domain.Sort; import org.springframework.data.domain.Sort.Direction; import org.springframework.data.repository.query.QueryMethodEvaluationContextProvider; +import org.springframework.data.util.Streamable; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -139,6 +139,12 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac .verifyComplete(); } + @Test // DATACASS-529 + public void shouldFindEmpptySliceByLastName() { + StepVerifier.create(repository.findByLastname("foo", CassandraPageRequest.first(1))) + .expectNextMatches(Streamable::isEmpty).verifyComplete(); + } + @Test // DATACASS-525 public void findOneWithManyResultsShouldFail() { StepVerifier.create(repository.findOneByLastname(dave.getLastname())) diff --git a/src/main/asciidoc/new-features.adoc b/src/main/asciidoc/new-features.adoc index e96d8ff62..bc0f464f0 100644 --- a/src/main/asciidoc/new-features.adoc +++ b/src/main/asciidoc/new-features.adoc @@ -11,8 +11,9 @@ This chapter summarizes changes and new features for each release. * Cassandra Mapped Tuple support via `@Tuple`. * Support for Cassandra `time` columns via `LocalTime`. * Support for `map` columns using User-defined/converted types. -* <> +* <>. * Kotlin extensions for Template API. +* Reactive Paging support through `Mono>`. [[new-features.2-0-0]] == What's new in Spring Data for Apache Cassandra 2.0 diff --git a/src/main/asciidoc/reference/cassandra-repositories.adoc b/src/main/asciidoc/reference/cassandra-repositories.adoc index cd08fc54f..e9af5ef1b 100644 --- a/src/main/asciidoc/reference/cassandra-repositories.adoc +++ b/src/main/asciidoc/reference/cassandra-repositories.adoc @@ -117,7 +117,7 @@ as the following example does by autowiring `PersonRepository`: ==== [source,java] ---- -@RunWith(SpringJUnit4ClassRunner.class) +@RunWith(SpringRunner.class) @ContextConfiguration public class PersonRepositoryTests { @@ -139,7 +139,7 @@ Cassandra repositories support paging and sorting for paginated and sorted acces ==== [source,java] ---- -@RunWith(SpringJUnit4ClassRunner.class) +@RunWith(SpringRunner.class) @ContextConfiguration public class PersonRepositoryTests { diff --git a/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc b/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc index ab3609d55..e367ed017 100644 --- a/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc +++ b/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc @@ -137,6 +137,35 @@ public class PersonRepositoryTests { ---- ==== +Cassandra repositories support paging and sorting for paginated and sorted access to the entities. Cassandra paging requires a paging state to forward-only navigate through pages. A `Slice` keeps track of the current paging state and allows for creation of a `Pageable` to request the next page. The following example shows how to set up paging access to `Person` entities: + +.Paging access to `Person` entities +==== +[source,java] +---- +@RunWith(SpringRunner.class) +@ContextConfiguration +public class PersonRepositoryTests { + + @Autowired PersonRepository repository; + + @Test + public void readsPagesCorrectly() { + + Mono> firstBatch = repository.findAll(CassandraPageRequest.first(10)); + + Mono> nextBatch = firstBatch.flatMap(it -> repository.findAll(it.nextPageable())); + + // … + } +} +---- +==== + +The preceding example creates an application context with Spring's unit test support, which performs annotation-based +dependency injection into the test class. Inside the test cases (the test methods), we use the repository to query +the data store. We invoke the repository query method that requests all `Person` instances. + [[cassandra.reactive.repositories.features]] == Features