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 fba88ab65..0792ecf8a 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 @@ -38,6 +38,7 @@ import org.springframework.data.domain.Pageable; import org.springframework.data.domain.Slice; import org.springframework.data.domain.SliceImpl; import org.springframework.util.Assert; +import reactor.core.publisher.Mono; import com.datastax.driver.core.PagingState; import com.datastax.driver.core.ResultSet; 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 66f754d15..fb9495653 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 @@ -25,6 +25,8 @@ import org.springframework.data.cassandra.core.cql.ReactiveCqlOperations; import org.springframework.data.cassandra.core.cql.WriteOptions; 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; @@ -97,6 +99,19 @@ 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. + * + * @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()} + * @throws DataAccessException if there is any problem executing the query. + * @since 2.0 + */ + Mono> slice(Statement statement, Class entityClass) throws DataAccessException; + /** * Execute a {@code SELECT} query and convert the resulting item to an entity. * 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 cf0440040..6ccca5a05 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,6 +15,9 @@ */ 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; @@ -33,6 +36,7 @@ import org.springframework.data.cassandra.core.convert.CassandraConverter; import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; import org.springframework.data.cassandra.core.convert.QueryMapper; import org.springframework.data.cassandra.core.convert.UpdateMapper; +import org.springframework.data.cassandra.core.cql.CassandraAccessor; import org.springframework.data.cassandra.core.cql.CqlIdentifier; import org.springframework.data.cassandra.core.cql.CqlProvider; import org.springframework.data.cassandra.core.cql.QueryOptions; @@ -51,6 +55,7 @@ import org.springframework.data.cassandra.core.mapping.event.AfterSaveEvent; 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.mapping.context.MappingContext; import org.springframework.data.projection.ProjectionFactory; import org.springframework.data.projection.SpelAwareProxyProjectionFactory; @@ -63,6 +68,7 @@ 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; @@ -230,6 +236,23 @@ public class ReactiveCassandraTemplate implements ReactiveCassandraOperations, A return getReactiveCqlOperations().query(cql, (row, rowNum) -> mapper.apply(row)); } + /* (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) { + + 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); + + return Mono.zip(resultSetMono, effectiveFetchSizeMono) + .flatMap(tuple -> QueryUtils.readSlice(tuple.getT1(), tuple.getT2(), rowMapper)); + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.ReactiveCassandraOperations#selectOne(com.datastax.driver.core.Statement, java.lang.Class) */ @@ -660,8 +683,27 @@ public class ReactiveCassandraTemplate implements ReactiveCassandraOperations, A return converter; } - @Value - static class StatementCallback implements ReactiveSessionCallback, CqlProvider { + @SuppressWarnings("ConstantConditions") + private Mono getEffectiveFetchSize(Statement statement) { + + 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()); + } + } + + return getReactiveCqlOperations().execute((ReactiveSessionCallback) session -> + Mono.fromSupplier(() -> session.getCluster().getConfiguration().getQueryOptions().getFetchSize()) + ).single(); + } + + @Value + static class StatementCallback implements ReactiveSessionCallback, CqlProvider { @lombok.NonNull Statement statement; 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 9c267cae7..1358a6348 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,6 +15,7 @@ */ 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; @@ -118,7 +119,7 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository ResultProcessor resultProcessor = getQueryMethod().getResultProcessor() .withDynamicProjection(convertingParameterAccessor); - ReactiveCassandraQueryExecution queryExecution = getExecution(new ResultProcessingConverter(resultProcessor, + ReactiveCassandraQueryExecution queryExecution = getExecution(parameterAccessor,new ResultProcessingConverter(resultProcessor, toMappingContext(getReactiveCassandraOperations()), getEntityInstantiators())); Class resultType = resolveResultType(resultProcessor); @@ -143,16 +144,19 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository /** * Returns the execution instance to use. - * + * @param parameterAccessor must not be {@literal null}. * @param resultProcessing must not be {@literal null}. @return */ - private ReactiveCassandraQueryExecution getExecution(Converter resultProcessing) { - return new ResultProcessingExecution(getExecutionToWrap(), resultProcessing); + private ReactiveCassandraQueryExecution getExecution(CassandraParameterAccessor parameterAccessor, + Converter resultProcessing) { + return new ResultProcessingExecution(getExecutionToWrap(parameterAccessor), resultProcessing); } private ReactiveCassandraQueryExecution getExecutionToWrap() { - if (getQueryMethod().isCollectionQuery()) { + if (getQueryMethod().isSliceQuery()) { + return new SlicedExecution(getReactiveCassandraOperations(), parameterAccessor.getPageable()); + }else if (getQueryMethod().isCollectionQuery()) { return new CollectionExecution(getReactiveCassandraOperations()); } else if (isCountQuery()) { return ((statement, type) -> 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 c62c40ffd..00796ccbb 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 @@ -26,7 +26,10 @@ import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.data.cassandra.core.ReactiveCassandraOperations; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; import org.springframework.data.cassandra.core.mapping.CassandraPersistentProperty; +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; @@ -46,6 +49,35 @@ interface ReactiveCassandraQueryExecution { Object execute(Statement statement, Class type); + /** + * {@link CassandraQueryExecution} for a {@link Slice}. + * + * @author Hleb Albau + */ + @RequiredArgsConstructor + final class SlicedExecution implements ReactiveCassandraQueryExecution { + + private final @NonNull ReactiveCassandraOperations operations; + private final @NonNull Pageable pageable; + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.repository.query.CassandraQueryExecution#execute(java.lang.String, java.lang.Class) + */ + @Override + public Object execute(Statement statement, Class type) { + + CassandraPageRequest.validatePageable(pageable); + + Statement statementToUse = statement.setFetchSize(pageable.getPageSize()); + + if (pageable instanceof CassandraPageRequest) { + statementToUse = statementToUse.setPagingState(((CassandraPageRequest) pageable).getPagingState()); + } + + return operations.slice(statementToUse, type); + } + } + /** * {@link ReactiveCassandraQueryExecution} for collection returning queries. * 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 c10e9ec3e..c8d8bbcae 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 @@ -36,6 +36,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Configuration; import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.data.cassandra.core.ReactiveCassandraOperations; +import org.springframework.data.cassandra.core.query.CassandraPageRequest; import org.springframework.data.cassandra.domain.Group; import org.springframework.data.cassandra.domain.GroupKey; import org.springframework.data.cassandra.domain.User; @@ -43,6 +44,8 @@ import org.springframework.data.cassandra.repository.support.IntegrationTestConf import org.springframework.data.cassandra.repository.support.ReactiveCassandraRepositoryFactory; import org.springframework.data.cassandra.repository.support.SimpleReactiveCassandraRepository; import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTest; +import org.springframework.data.domain.Pageable; +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; @@ -129,6 +132,13 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac StepVerifier.create(repository.findByLastname(dave.getLastname())).expectNextCount(2).verifyComplete(); } + @Test //DATACASS-529 + public void shouldFindSliceByLastName() { + StepVerifier.create(repository.findByLastname(carter.getLastname(), CassandraPageRequest.first(1))) + .expectNextMatches(users -> users.getSize() == 1 && users.hasNext()) + .verifyComplete(); + } + @Test // DATACASS-525 public void findOneWithManyResultsShouldFail() { StepVerifier.create(repository.findOneByLastname(dave.getLastname())) @@ -206,6 +216,8 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac Flux findByLastname(String lastname); + Mono> findByLastname(String firstname, Pageable pageable); + Mono findFirstByLastname(String lastname); Mono findOneByLastname(String lastname);