DATACOUCH-605 - Support ScanConsistency in n1ql queries

Also fixes expiry bug
This commit is contained in:
mikereiche
2020-10-01 21:35:19 -07:00
parent 7026a39589
commit b54dfcaf1c
17 changed files with 117 additions and 28 deletions

View File

@@ -126,7 +126,7 @@ public interface ExecutableFindByQueryOperation {
*
* @param scanConsistency the custom scan consistency to use for this query.
*/
FindByQueryWithQuery<T> consistentWith(QueryScanConsistency scanConsistency);
FindByQueryConsistentWith<T> consistentWith(QueryScanConsistency scanConsistency);
}

View File

@@ -21,6 +21,7 @@ import java.util.stream.Stream;
import org.springframework.data.couchbase.core.query.Query;
import com.couchbase.client.java.query.QueryScanConsistency;
import org.springframework.data.couchbase.core.ReactiveFindByQueryOperationSupport.ReactiveFindByQuerySupport;
public class ExecutableFindByQueryOperationSupport implements ExecutableFindByQueryOperation {
@@ -50,8 +51,7 @@ public class ExecutableFindByQueryOperationSupport implements ExecutableFindByQu
this.template = template;
this.domainType = domainType;
this.query = query;
this.reactiveSupport = new ReactiveFindByQueryOperationSupport.ReactiveFindByQuerySupport<T>(
template.reactive(), domainType, query, scanConsistency);
this.reactiveSupport = new ReactiveFindByQuerySupport<T>(template.reactive(), domainType, query, scanConsistency);
this.scanConsistency = scanConsistency;
}
@@ -72,11 +72,17 @@ public class ExecutableFindByQueryOperationSupport implements ExecutableFindByQu
@Override
public TerminatingFindByQuery<T> matching(final Query query) {
return new ExecutableFindByQuerySupport<>(template, domainType, query, scanConsistency);
QueryScanConsistency scanCons;
if (query.getScanConsistency() != null) {
scanCons = query.getScanConsistency();
} else {
scanCons = scanConsistency;
}
return new ExecutableFindByQuerySupport<>(template, domainType, query, scanCons);
}
@Override
public FindByQueryWithQuery<T> consistentWith(final QueryScanConsistency scanConsistency) {
public FindByQueryConsistentWith<T> consistentWith(final QueryScanConsistency scanConsistency) {
return new ExecutableFindByQuerySupport<>(template, domainType, query, scanConsistency);
}

View File

@@ -100,7 +100,7 @@ public interface ReactiveFindByQueryOperation {
*
* @param scanConsistency the custom scan consistency to use for this query.
*/
FindByQueryWithQuery<T> consistentWith(QueryScanConsistency scanConsistency);
FindByQueryConsistentWith<T> consistentWith(QueryScanConsistency scanConsistency);
}

View File

@@ -60,11 +60,17 @@ public class ReactiveFindByQueryOperationSupport implements ReactiveFindByQueryO
@Override
public TerminatingFindByQuery<T> matching(Query query) {
return new ReactiveFindByQuerySupport<>(template, domainType, query, scanConsistency);
QueryScanConsistency scanCons;
if (query.getScanConsistency() != null) {
scanCons = query.getScanConsistency();
} else {
scanCons = scanConsistency;
}
return new ReactiveFindByQuerySupport<>(template, domainType, query, scanCons);
}
@Override
public FindByQueryWithQuery<T> consistentWith(QueryScanConsistency scanConsistency) {
public FindByQueryConsistentWith<T> consistentWith(QueryScanConsistency scanConsistency) {
return new ReactiveFindByQuerySupport<>(template, domainType, query, scanConsistency);
}

View File

@@ -15,7 +15,6 @@
*/
package org.springframework.data.couchbase.core;
import com.couchbase.client.java.kv.UpsertOptions;
import org.springframework.data.couchbase.core.mapping.Document;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@@ -98,7 +97,7 @@ public class ReactiveInsertByIdOperationSupport implements ReactiveInsertByIdOpe
} else if (durabilityLevel != DurabilityLevel.NONE) {
options.durability(durabilityLevel);
}
if (expiry != null) {
if (expiry != null && ! expiry.isZero()) {
options.expiry(expiry);
} else if (domainType.isAnnotationPresent(Document.class)) {
Document documentAnn = domainType.getAnnotation(Document.class);

View File

@@ -15,7 +15,6 @@
*/
package org.springframework.data.couchbase.core;
import com.couchbase.client.java.kv.Expiry;
import org.springframework.data.couchbase.core.mapping.Document;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@@ -98,7 +97,7 @@ public class ReactiveUpsertByIdOperationSupport implements ReactiveUpsertByIdOpe
} else if (durabilityLevel != DurabilityLevel.NONE) {
options.durability(durabilityLevel);
}
if (expiry != null) {
if (expiry != null && !expiry.isZero()) {
options.expiry(expiry);
} else if (domainType.isAnnotationPresent(Document.class)) {
Document documentAnn = domainType.getAnnotation(Document.class);

View File

@@ -47,6 +47,7 @@ public class Query {
private long skip;
private int limit;
private Sort sort = Sort.unsorted();
private QueryScanConsistency queryScanConsistency;
static private final Pattern WHERE_PATTERN = Pattern.compile("\\sWHERE\\s");
@@ -127,6 +128,27 @@ public class Query {
return with(pageable.getSort());
}
/**
* queryScanConsistency
*
* @return queryScanConsistency
*/
public QueryScanConsistency getScanConsistency() {
return queryScanConsistency;
}
/**
* Sets the given scan consistency on the {@link Query} instance.
*
* @param queryScanConsistency
* @return this
*/
public Query scanConsistency(final QueryScanConsistency queryScanConsistency) {
this.queryScanConsistency = queryScanConsistency;
return this;
}
/**
* Adds a {@link Sort} to the {@link Query} instance.
*
@@ -280,7 +302,7 @@ public class Query {
}
/**
* build QueryOptions forom parameters and scanConsistency
* build QueryOptions from parameters and scanConsistency
*
* @param scanConsistency
* @return QueryOptions

View File

@@ -18,6 +18,7 @@ package org.springframework.data.couchbase.repository;
import java.util.List;
import com.couchbase.client.java.query.QueryScanConsistency;
import org.springframework.data.domain.Sort;
import org.springframework.data.repository.NoRepositoryBean;
import org.springframework.data.repository.PagingAndSortingRepository;
@@ -34,6 +35,8 @@ public interface CouchbaseRepository<T, ID> extends PagingAndSortingRepository<T
@Override
List<T> findAll(Sort sort);
List<T> findAll(QueryScanConsistency queryScanConsistency);
@Override
List<T> findAll();

View File

@@ -25,6 +25,7 @@ import org.springframework.data.couchbase.core.query.Dimensional;
import org.springframework.data.couchbase.core.query.View;
import org.springframework.data.couchbase.core.query.WithConsistency;
import org.springframework.data.couchbase.repository.Query;
import org.springframework.data.couchbase.repository.ScanConsistency;
import org.springframework.data.mapping.context.MappingContext;
import org.springframework.data.projection.ProjectionFactory;
import org.springframework.data.repository.core.RepositoryMetadata;
@@ -152,6 +153,24 @@ public class CouchbaseQueryMethod extends QueryMethod {
return method.getAnnotation(WithConsistency.class);
}
/**
* If the method has a @ScanConsistency annotation
*
* @return true if this has the @ScanConsistency annotation
*/
public boolean hasScanConsistencyAnnotation() {
return getScanConsistencyAnnotation() != null;
}
/**
* ScanConsistency annotation
*
* @return the @ScanConsistency annotation
*/
public ScanConsistency getScanConsistencyAnnotation() {
return method.getAnnotation(ScanConsistency.class);
}
/**
* Returns the query string declared in a {@link Query} annotation or {@literal null} if neither the annotation found
* nor the attribute was specified.

View File

@@ -15,9 +15,7 @@
*/
package org.springframework.data.couchbase.repository.query;
import java.util.ArrayList;
import java.util.List;
import com.couchbase.client.java.query.QueryScanConsistency;
import org.springframework.data.couchbase.core.CouchbaseOperations;
import org.springframework.data.couchbase.core.ExecutableFindByQueryOperation;
import org.springframework.data.couchbase.core.query.Query;
@@ -65,7 +63,8 @@ public class N1qlRepositoryQueryExecutor {
final PartTree tree = new PartTree(queryMethod.getName(), domainClass);
query = new N1qlQueryCreator(tree, accessor, queryMethod, operations.getConverter()).createQuery();
}
q = (ExecutableFindByQueryOperation.ExecutableFindByQuery) operations.findByQuery(domainClass).matching(query);
q = (ExecutableFindByQueryOperation.ExecutableFindByQuery) operations.findByQuery(domainClass)
.consistentWith(buildQueryScanConsistency()).matching(query);
if (queryMethod.isCountQuery()) {
return q.count();
} else if (queryMethod.isCollectionQuery()) {
@@ -76,4 +75,14 @@ public class N1qlRepositoryQueryExecutor {
}
private QueryScanConsistency buildQueryScanConsistency() {
QueryScanConsistency scanConsistency = QueryScanConsistency.NOT_BOUNDED;
if (queryMethod.hasConsistencyAnnotation()) {
scanConsistency = queryMethod.getConsistencyAnnotation().value();
} else if (queryMethod.hasScanConsistencyAnnotation()) {
scanConsistency = queryMethod.getScanConsistencyAnnotation().query();
}
return scanConsistency;
}
}

View File

@@ -15,6 +15,7 @@
*/
package org.springframework.data.couchbase.repository.query;
import com.couchbase.client.java.query.QueryScanConsistency;
import org.springframework.data.couchbase.core.ExecutableFindByQueryOperation;
import org.springframework.data.couchbase.core.ReactiveFindByQueryOperation;
import org.springframework.data.mapping.context.MappingContext;
@@ -69,7 +70,8 @@ public class ReactiveN1qlRepositoryQueryExecutor {
final PartTree tree = new PartTree(queryMethod.getName(), domainClass);
query = new N1qlQueryCreator(tree, accessor, queryMethod, operations.getConverter()).createQuery();
}
q = (ReactiveFindByQueryOperation.ReactiveFindByQuery) operations.findByQuery(domainClass).matching(query);
q = (ReactiveFindByQueryOperation.ReactiveFindByQuery) operations.findByQuery(domainClass)
.consistentWith(buildQueryScanConsistency()).matching(query);
if (queryMethod.isCountQuery()) {
return q.count();
} else if (queryMethod.isCollectionQuery()) {
@@ -79,4 +81,14 @@ public class ReactiveN1qlRepositoryQueryExecutor {
}
}
private QueryScanConsistency buildQueryScanConsistency() {
QueryScanConsistency scanConsistency = QueryScanConsistency.NOT_BOUNDED;
if (queryMethod.hasConsistencyAnnotation()) {
scanConsistency = queryMethod.getConsistencyAnnotation().value();
} else if (queryMethod.hasScanConsistencyAnnotation()) {
scanConsistency = queryMethod.getScanConsistencyAnnotation().query();
}
return scanConsistency;
}
}

View File

@@ -147,6 +147,11 @@ public class SimpleCouchbaseRepository<T, ID> implements CouchbaseRepository<T,
return findAll(new Query().with(sort));
}
@Override
public List<T> findAll(final QueryScanConsistency queryScanConsistency) {
return findAll(new Query().scanConsistency(queryScanConsistency));
}
@Override
public Page<T> findAll(final Pageable pageable) {
List<T> results = findAll(new Query().with(pageable));

View File

@@ -56,6 +56,7 @@ class CouchbaseTemplateKeyValueIntegrationTests extends ClusterAwareIntegrationT
private static CouchbaseClientFactory couchbaseClientFactory;
private CouchbaseTemplate couchbaseTemplate;
private ReactiveCouchbaseTemplate reactiveCouchbaseTemplate;
@BeforeAll
static void beforeAll() {
@@ -72,6 +73,7 @@ class CouchbaseTemplateKeyValueIntegrationTests extends ClusterAwareIntegrationT
void beforeEach() {
ApplicationContext ac = new AnnotationConfigApplicationContext(Config.class);
couchbaseTemplate = (CouchbaseTemplate) ac.getBean(COUCHBASE_TEMPLATE);
reactiveCouchbaseTemplate = (ReactiveCouchbaseTemplate) ac.getBean(REACTIVE_COUCHBASE_TEMPLATE);
}
@Test
@@ -90,6 +92,7 @@ class CouchbaseTemplateKeyValueIntegrationTests extends ClusterAwareIntegrationT
assertEquals(user, found);
couchbaseTemplate.removeById().one(user.getId());
reactiveCouchbaseTemplate.replaceById(User.class).withDurability(PersistTo.ACTIVE, ReplicateTo.THREE).one(user);
}
@Test

View File

@@ -40,15 +40,19 @@ public interface AirportRepository extends PagingAndSortingRepository<Airport, S
Iterable<Airport> findAll();
@Override
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Airport save(Airport airport);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
List<Airport> findAllByIata(String iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Airport findByIata(String iata);
@Query("#{#n1ql.selectEntity} where iata = $1")
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
List<Airport> getAllByIata(String iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
long countByIataIn(String... iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)

View File

@@ -16,6 +16,7 @@
package org.springframework.data.couchbase.domain;
import org.springframework.data.couchbase.repository.Query;
import reactor.core.publisher.Flux;
import org.springframework.data.couchbase.repository.ScanConsistency;
@@ -44,6 +45,10 @@ public interface ReactiveAirportRepository extends ReactiveSortingRepository<Air
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Flux<Airport> findAllByIata(String iata);
@Query("#{#n1ql.selectEntity} where iata = $1")
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Flux<Airport> getAllByIata(String iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Mono<Long> countByIataIn(String... iatas);

View File

@@ -25,8 +25,6 @@ import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.stream.Collectors;
import java.util.stream.StreamSupport;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -93,7 +91,6 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr
airportRepository.save(vie);
xxx = new Airport("airports::xxx", "xxx", "xxxx");
airportRepository.save(xxx);
sleep(1000);
List<Airport> airports;
airports = airportRepository.findAllByIata("1\" or iata=iata or iata=\"1");
assertEquals(0, airports.size());
@@ -112,7 +109,6 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr
try {
vie = new Airport("airports::vie", "vie", "loww");
airportRepository.save(vie);
sleep(1000);
List<Airport> airports = airportRepository.findAllByIata("vie");
assertEquals(vie.getId(), airports.get(0).getId());
} finally {
@@ -170,7 +166,7 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr
iatas[i].toLowerCase(Locale.ROOT) /* lcao */);
airportRepository.save(airport);
}
sleep(1000);
for (int k = 0; k < 50; k++) {
Callable<Boolean>[] suppliers = new Callable[iatas.length];
for (int i = 0; i < iatas.length; i++) {
@@ -211,7 +207,7 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr
Airport airport = new Airport("airports::" + iatas[i], iatas[i] /*iata*/, iatas[i].toLowerCase() /* lcao */);
airportRepository.save(airport);
}
sleep(1000);
for (int k = 0; k < 100; k++) {
Callable<Boolean>[] suppliers = new Callable[iatas.length];
for (int i = 0; i < iatas.length; i++) {

View File

@@ -88,9 +88,10 @@ public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends ClusterAwa
try {
vie = new Airport("airports::vie", "vie", "loww");
airportRepository.save(vie).block();
List<Airport> airports = airportRepository.findAllByIata("vie").collectList().block();
// TODO
System.err.println(airports);
List<Airport> airports1 = airportRepository.findAllByIata("vie").collectList().block();
assertEquals(1,airports1.size());
List<Airport> airports2 = airportRepository.findAllByIata("vie").collectList().block();
assertEquals(1,airports2.size());
} finally {
airportRepository.delete(vie).block();
}