DATACOUCH-605 - Support ScanConsistency in n1ql queries

This commit is contained in:
mikereiche
2020-10-01 21:35:19 -07:00
parent fab72ee9c9
commit 82edad5759
17 changed files with 117 additions and 28 deletions

View File

@@ -55,6 +55,7 @@ class CouchbaseTemplateKeyValueIntegrationTests extends ClusterAwareIntegrationT
private static CouchbaseClientFactory couchbaseClientFactory;
private CouchbaseTemplate couchbaseTemplate;
private ReactiveCouchbaseTemplate reactiveCouchbaseTemplate;
@BeforeAll
static void beforeAll() {
@@ -71,6 +72,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
@@ -89,6 +91,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();
}