DATACASS-359 - Polish for reactive repository query methods DTO projections support.
Original pull request: #91.
This commit is contained in:
@@ -40,9 +40,9 @@ import org.springframework.util.Assert;
|
||||
*/
|
||||
public abstract class AbstractReactiveCassandraQuery implements RepositoryQuery {
|
||||
|
||||
private final ReactiveCassandraQueryMethod method;
|
||||
private final ReactiveCassandraOperations operations;
|
||||
private final EntityInstantiators instantiators;
|
||||
private final ReactiveCassandraOperations operations;
|
||||
private final ReactiveCassandraQueryMethod method;
|
||||
|
||||
/**
|
||||
* Creates a new {@link AbstractReactiveCassandraQuery} from the given {@link CassandraQueryMethod} and
|
||||
@@ -61,7 +61,8 @@ public abstract class AbstractReactiveCassandraQuery implements RepositoryQuery
|
||||
this.instantiators = new EntityInstantiators();
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.repository.query.RepositoryQuery#getQueryMethod()
|
||||
*/
|
||||
@Override
|
||||
@@ -69,14 +70,15 @@ public abstract class AbstractReactiveCassandraQuery implements RepositoryQuery
|
||||
return method;
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.repository.query.RepositoryQuery#execute(java.lang.Object[])
|
||||
*/
|
||||
@Override
|
||||
public Object execute(Object[] parameters) {
|
||||
|
||||
return method.hasReactiveWrapperParameter() ? executeDeferred(parameters)
|
||||
: execute(new ReactiveCassandraParameterAccessor(method, parameters));
|
||||
return (method.hasReactiveWrapperParameter() ? executeDeferred(parameters)
|
||||
: execute(new ReactiveCassandraParameterAccessor(method, parameters)));
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -84,50 +86,31 @@ public abstract class AbstractReactiveCassandraQuery implements RepositoryQuery
|
||||
|
||||
ReactiveCassandraParameterAccessor accessor = new ReactiveCassandraParameterAccessor(method, parameters);
|
||||
|
||||
if (getQueryMethod().isCollectionQuery()) {
|
||||
return Flux.defer(() -> (Publisher<Object>) execute(accessor));
|
||||
}
|
||||
|
||||
return Mono.defer(() -> (Mono<Object>) execute(accessor));
|
||||
return (getQueryMethod().isCollectionQuery() ? Flux.defer(() -> (Publisher<Object>) execute(accessor))
|
||||
: Mono.defer(() -> (Mono<Object>) execute(accessor)));
|
||||
}
|
||||
|
||||
private Object execute(CassandraParameterAccessor parameterAccessor) {
|
||||
|
||||
CassandraParameterAccessor convertingParameterAccessor = new ConvertingParameterAccessor(operations.getConverter(),
|
||||
parameterAccessor);
|
||||
CassandraParameterAccessor convertingParameterAccessor =
|
||||
new ConvertingParameterAccessor(operations.getConverter(), parameterAccessor);
|
||||
|
||||
String query = createQuery(convertingParameterAccessor);
|
||||
|
||||
ResultProcessor resultProcessor = method.getResultProcessor().withDynamicProjection(convertingParameterAccessor);
|
||||
|
||||
ReactiveCassandraQueryExecution queryExecution = getExecution(convertingParameterAccessor,
|
||||
new ResultProcessingConverter(resultProcessor, operations.getConverter().getMappingContext(), instantiators));
|
||||
ReactiveCassandraQueryExecution queryExecution = getExecution(new ResultProcessingConverter(
|
||||
resultProcessor, operations.getConverter().getMappingContext(), instantiators));
|
||||
|
||||
CassandraReturnedType returnedType = new CassandraReturnedType(resultProcessor.getReturnedType(),
|
||||
operations.getConverter().getCustomConversions());
|
||||
|
||||
Class<?> resultType = (returnedType.isProjecting() ? returnedType.getDomainType() : returnedType.getReturnedType());
|
||||
Class<?> resultType = (returnedType.isProjecting() ? returnedType.getDomainType()
|
||||
: returnedType.getReturnedType());
|
||||
|
||||
return queryExecution.execute(query, resultType);
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the execution instance to use.
|
||||
*
|
||||
* @param accessor must not be {@literal null}.
|
||||
* @param resultProcessing must not be {@literal null}. @return
|
||||
*/
|
||||
private ReactiveCassandraQueryExecution getExecution(CassandraParameterAccessor accessor,
|
||||
Converter<Object, Object> resultProcessing) {
|
||||
|
||||
return new ResultProcessingExecution(getExecutionToWrap(), resultProcessing);
|
||||
}
|
||||
|
||||
private ReactiveCassandraQueryExecution getExecutionToWrap() {
|
||||
|
||||
return (method.isCollectionQuery() ? new CollectionExecution(operations) : new SingleEntityExecution(operations));
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a string query using the given {@link ParameterAccessor}
|
||||
*
|
||||
@@ -135,4 +118,17 @@ public abstract class AbstractReactiveCassandraQuery implements RepositoryQuery
|
||||
*/
|
||||
protected abstract String createQuery(CassandraParameterAccessor accessor);
|
||||
|
||||
/**
|
||||
* Returns the execution instance to use.
|
||||
*
|
||||
* @param resultProcessing must not be {@literal null}. @return
|
||||
*/
|
||||
private ReactiveCassandraQueryExecution getExecution(Converter<Object, Object> resultProcessing) {
|
||||
return new ResultProcessingExecution(getExecutionToWrap(), resultProcessing);
|
||||
}
|
||||
|
||||
/* (non-Javadoc) */
|
||||
private ReactiveCassandraQueryExecution getExecutionToWrap() {
|
||||
return (method.isCollectionQuery() ? new CollectionExecution(operations) : new SingleEntityExecution(operations));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -47,7 +47,8 @@ interface ReactiveCassandraQueryExecution {
|
||||
|
||||
private final @NonNull ReactiveCassandraOperations operations;
|
||||
|
||||
/* (non-Javadoc)
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution#execute(java.lang.String, java.lang.Class)
|
||||
*/
|
||||
@Override
|
||||
@@ -66,7 +67,8 @@ interface ReactiveCassandraQueryExecution {
|
||||
|
||||
private final @NonNull ReactiveCassandraOperations operations;
|
||||
|
||||
/* (non-Javadoc)
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution#execute(java.lang.String, java.lang.Class)
|
||||
*/
|
||||
@Override
|
||||
@@ -87,7 +89,8 @@ interface ReactiveCassandraQueryExecution {
|
||||
private final @NonNull ReactiveCassandraQueryExecution delegate;
|
||||
private final @NonNull Converter<Object, Object> converter;
|
||||
|
||||
/* (non-Javadoc)
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution#execute(java.lang.String, java.lang.Class)
|
||||
*/
|
||||
@Override
|
||||
@@ -108,7 +111,8 @@ interface ReactiveCassandraQueryExecution {
|
||||
private final @NonNull CassandraMappingContext mappingContext;
|
||||
private final @NonNull EntityInstantiators instantiators;
|
||||
|
||||
/* (non-Javadoc)
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.core.convert.converter.Converter#convert(java.lang.Object)
|
||||
*/
|
||||
@Override
|
||||
|
||||
@@ -51,7 +51,6 @@ import com.datastax.driver.core.TableMetadata;
|
||||
* Test for {@link ReactiveCassandraRepository} using reactive wrapper type conversion.
|
||||
*
|
||||
* @author Mark Paluch
|
||||
* @soundtrack Dj Marc - Euromix 97 Part 1
|
||||
*/
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@ContextConfiguration(classes = ConvertingReactiveCassandraRepositoryTests.Config.class)
|
||||
@@ -83,7 +82,6 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
TableMetadata person = keyspace.getTable("person");
|
||||
|
||||
if (person.getIndex("IX_person_lastname") == null) {
|
||||
|
||||
session.execute("CREATE INDEX IX_person_lastname ON person (lastname);");
|
||||
Thread.sleep(500);
|
||||
}
|
||||
@@ -96,6 +94,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
boyd = new Person("45", "Boyd", "Tinsley");
|
||||
|
||||
TestSubscriber<Person> subscriber = TestSubscriber.create();
|
||||
|
||||
reactiveRepository.save(Arrays.asList(oliver, dave, carter, boyd)).subscribe(subscriber);
|
||||
|
||||
subscriber.await().assertComplete().assertNoError();
|
||||
@@ -118,8 +117,8 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
@Test
|
||||
public void reactiveStreamsQueryMethodsShouldWork() {
|
||||
|
||||
TestSubscriber<Person> subscriber = TestSubscriber
|
||||
.subscribe(reactivePersonRepostitory.findByLastname(boyd.getLastname()));
|
||||
TestSubscriber<Person> subscriber = TestSubscriber.subscribe(
|
||||
reactivePersonRepostitory.findByLastname(boyd.getLastname()));
|
||||
|
||||
subscriber.awaitAndAssertNextValueCount(1).assertValues(boyd);
|
||||
}
|
||||
@@ -130,11 +129,10 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
@Test
|
||||
public void dtoProjectionShouldWork() {
|
||||
|
||||
TestSubscriber<PersonDto> subscriber = TestSubscriber
|
||||
.subscribe(reactivePersonRepostitory.findProjectedByLastname(boyd.getLastname()));
|
||||
TestSubscriber<PersonDto> subscriber = TestSubscriber.subscribe(
|
||||
reactivePersonRepostitory.findProjectedByLastname(boyd.getLastname()));
|
||||
|
||||
subscriber.awaitAndAssertNextValueCount(1).assertValuesWith(personDto -> {
|
||||
|
||||
assertThat(personDto.firstname).isEqualTo(boyd.getFirstname());
|
||||
assertThat(personDto.lastname).isEqualTo(boyd.getLastname());
|
||||
});
|
||||
@@ -147,6 +145,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
public void simpleRxJavaMethodsShouldWork() {
|
||||
|
||||
rx.observers.TestSubscriber<Boolean> subscriber = new rx.observers.TestSubscriber<>();
|
||||
|
||||
rxJava1PersonRepostitory.exists(dave.getId()).subscribe(subscriber);
|
||||
|
||||
subscriber.awaitTerminalEvent();
|
||||
@@ -162,6 +161,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
public void existsWithSingleRxJavaIdMethodsShouldWork() {
|
||||
|
||||
rx.observers.TestSubscriber<Boolean> subscriber = new rx.observers.TestSubscriber<>();
|
||||
|
||||
rxJava1PersonRepostitory.exists(Single.just(dave.getId())).subscribe(subscriber);
|
||||
|
||||
subscriber.awaitTerminalEvent();
|
||||
@@ -177,6 +177,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
public void singleRxJavaQueryMethodShouldWork() {
|
||||
|
||||
rx.observers.TestSubscriber<Person> subscriber = new rx.observers.TestSubscriber<>();
|
||||
|
||||
rxJava1PersonRepostitory.findManyByLastname(dave.getLastname()).subscribe(subscriber);
|
||||
|
||||
subscriber.awaitTerminalEvent();
|
||||
@@ -192,6 +193,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
public void singleProjectedRxJavaQueryMethodShouldWork() {
|
||||
|
||||
rx.observers.TestSubscriber<ProjectedPerson> subscriber = new rx.observers.TestSubscriber<>();
|
||||
|
||||
rxJava1PersonRepostitory.findProjectedByLastname(carter.getLastname()).subscribe(subscriber);
|
||||
|
||||
subscriber.awaitTerminalEvent();
|
||||
@@ -209,6 +211,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
|
||||
public void observableRxJavaQueryMethodShouldWork() {
|
||||
|
||||
rx.observers.TestSubscriber<Person> subscriber = new rx.observers.TestSubscriber<>();
|
||||
|
||||
rxJava1PersonRepostitory.findByLastname(boyd.getLastname()).subscribe(subscriber);
|
||||
|
||||
subscriber.awaitTerminalEvent();
|
||||
|
||||
@@ -19,8 +19,6 @@ import java.time.LocalDate;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
|
||||
import lombok.Data;
|
||||
import lombok.Getter;
|
||||
import org.springframework.data.cassandra.repository.CassandraRepository;
|
||||
import org.springframework.data.cassandra.repository.Query;
|
||||
import org.springframework.data.cassandra.test.integration.repository.querymethods.declared.Address;
|
||||
@@ -73,12 +71,11 @@ interface PersonRepository extends CassandraRepository<Person> {
|
||||
String getLastname();
|
||||
}
|
||||
|
||||
static class PersonDto {
|
||||
class PersonDto {
|
||||
|
||||
public String firstname, lastname;
|
||||
|
||||
public PersonDto(String firstname, String lastname) {
|
||||
|
||||
this.firstname = firstname;
|
||||
this.lastname = lastname;
|
||||
}
|
||||
|
||||
@@ -181,7 +181,8 @@ public class QueryDerivationIntegrationTests extends AbstractSpringDataEmbeddedC
|
||||
|
||||
Collection<PersonProjection> collection = personRepository.findPersonProjectedBy();
|
||||
|
||||
assertThat(collection).hasSize(3).extracting("lastname").contains("White", "White", "White");
|
||||
assertThat(collection).hasSize(3).extracting("firstname").contains(
|
||||
flynn.getFirstname(), skyler.getFirstname(), walter.getFirstname());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -192,7 +193,8 @@ public class QueryDerivationIntegrationTests extends AbstractSpringDataEmbeddedC
|
||||
|
||||
Collection<PersonDto> collection = personRepository.findPersonDtoBy();
|
||||
|
||||
assertThat(collection).hasSize(3).extracting("lastname").contains("White", "White", "White");
|
||||
assertThat(collection).hasSize(3).extracting("firstname").contains(
|
||||
flynn.getFirstname(), skyler.getFirstname(), walter.getFirstname());
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user