DATACASS-529 - Allow slice queries using reactive repositories.

We now support Slice queries using reactive Cassandra repositories. Repository query methods can declare Mono<Slice<T>> as their return type to query for slices without applying transparent paging. The resulting Mono always completes with a value. An empty query result returns a Mono emitting an empty Slice.

interface UserRepository extends CrudRepository<User, String> {

  Mono<Slice<User>> findByLastname(String firstname, Pageable pageable);
}

Mono<Slice<User>> result =repository.findByLastname("White", CassandraPageRequest.first(1));

Original pull request: #128.
This commit is contained in:
Hleb Albau
2018-05-14 22:30:24 +03:00
committed by Mark Paluch
parent 34903b7350
commit d261301e9e
6 changed files with 113 additions and 7 deletions

View File

@@ -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;

View File

@@ -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
*/
<T> Flux<T> select(Statement statement, Class<T> 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
*/
<T> Mono<Slice<T>> slice(Statement statement, Class<T> entityClass) throws DataAccessException;
/**
* Execute a {@code SELECT} query and convert the resulting item to an entity.
*

View File

@@ -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 <T> Mono<Slice<T>> slice(Statement statement, Class<T> entityClass) {
Assert.notNull(statement, "Statement must not be null");
Assert.notNull(entityClass, "Entity type must not be null");
Mono<ReactiveResultSet> resultSetMono = getReactiveCqlOperations().queryForResultSet(statement);
Mono<Integer> effectiveFetchSizeMono = getEffectiveFetchSize(statement);
Function<Row,T> 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<WriteResult>, CqlProvider {
@SuppressWarnings("ConstantConditions")
private Mono<Integer> 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<Integer>) session ->
Mono.fromSupplier(() -> session.getCluster().getConfiguration().getQueryOptions().getFetchSize())
).single();
}
@Value
static class StatementCallback implements ReactiveSessionCallback<WriteResult>, CqlProvider {
@lombok.NonNull Statement statement;

View File

@@ -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<Object, Object> resultProcessing) {
return new ResultProcessingExecution(getExecutionToWrap(), resultProcessing);
private ReactiveCassandraQueryExecution getExecution(CassandraParameterAccessor parameterAccessor,
Converter<Object, Object> 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) ->

View File

@@ -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.
*

View File

@@ -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<User> findByLastname(String lastname);
Mono<Slice<User>> findByLastname(String firstname, Pageable pageable);
Mono<User> findFirstByLastname(String lastname);
Mono<User> findOneByLastname(String lastname);