DATACOUCH-585 - Support ScanConsistency for count() queries.

This commit is contained in:
mikereiche
2020-07-16 17:31:37 -07:00
parent e2444403a7
commit 05fb9eadbe
5 changed files with 38 additions and 12 deletions

View File

@@ -181,7 +181,7 @@ public class SimpleReactiveCouchbaseRepository<T, ID> implements ReactiveCouchba
@SuppressWarnings("unchecked")
@Override
public Mono<Long> count() {
return operations.findByQuery(entityInformation.getJavaType()).count();
return operations.findByQuery(entityInformation.getJavaType()).consistentWith(buildQueryScanConsistency()).count();
}
@SuppressWarnings("unchecked")

View File

@@ -18,6 +18,7 @@ package org.springframework.data.couchbase.domain;
import java.util.List;
import org.jetbrains.annotations.NotNull;
import org.springframework.data.couchbase.repository.Query;
import org.springframework.data.couchbase.repository.ScanConsistency;
import org.springframework.data.repository.PagingAndSortingRepository;
@@ -38,6 +39,11 @@ public interface AirportRepository extends PagingAndSortingRepository<Airport, S
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Iterable<Airport> findAll();
@Override
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Airport save(Airport airport);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
List<Airport> findAllByIata(String iata);
@Query("#{#n1ql.selectEntity} where iata = $1")
@@ -45,8 +51,14 @@ public interface AirportRepository extends PagingAndSortingRepository<Airport, S
long countByIataIn(String... iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
long countByIcaoAndIataIn(String icao, String... iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
long countByIcaoOrIataIn(String icao, String... iata);
@Override
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
long count();
}

View File

@@ -38,9 +38,23 @@ public interface ReactiveAirportRepository extends ReactiveSortingRepository<Air
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Flux<Airport> findAll();
@Override
Mono<Airport> save(Airport a);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Flux<Airport> findAllByIata(String iata);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Mono<Long> countByIataIn(String... iatas);
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Mono<Long> countByIcaoAndIataIn(String icao, String... iatas);
@Override
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Mono<Long> count();
@Override
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Mono<Airport> findById(String var1);
}

View File

@@ -134,9 +134,8 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr
iatas[i].toLowerCase(Locale.ROOT) /* lcao */);
airportRepository.save(airport);
}
sleep(1000);
long airportCount = 0;
airportCount = airportRepository.count();
long airportCount = airportRepository.count();
assertEquals(7, airportCount);
airportCount = airportRepository.countByIataIn("JFK", "IAD", "SFO");

View File

@@ -29,6 +29,7 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Configuration;
import org.springframework.dao.DataRetrievalFailureException;
import org.springframework.data.couchbase.CouchbaseClientFactory;
import org.springframework.data.couchbase.config.AbstractCouchbaseConfiguration;
import org.springframework.data.couchbase.domain.Airport;
@@ -100,19 +101,15 @@ public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends ClusterAwa
String[] iatas = { "JFK", "IAD", "SFO", "SJC", "SEA", "LAX", "PHX" };
Future[] future = new Future[iatas.length];
ExecutorService executorService = Executors.newFixedThreadPool(iatas.length);
try {
Callable<Boolean>[] suppliers = new Callable[iatas.length];
for (int i = 0; i < iatas.length; i++) {
Airport airport = new Airport("airports::" + iatas[i], iatas[i] /*iata*/, iatas[i].toLowerCase() /* lcao */);
airportRepository.save(airport).block();
}
try {
Thread.sleep(1000);
} catch (InterruptedException ie) {}
Long airportCount = null;
airportCount = airportRepository.count().block();
assertEquals(7, airportCount);
Long airportCount = airportCount = airportRepository.count().block();
assertEquals(iatas.length, airportCount);
airportCount = airportRepository.countByIataIn("JFK", "IAD", "SFO").block();
assertEquals(3, airportCount);
@@ -126,7 +123,11 @@ public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends ClusterAwa
} finally {
for (int i = 0; i < iatas.length; i++) {
Airport airport = new Airport("airports::" + iatas[i], iatas[i] /*iata*/, iatas[i] /* lcao */);
airportRepository.delete(airport).block();
try {
airportRepository.delete(airport).block();
} catch (DataRetrievalFailureException drfe) {
System.out.println("Failed to delete: " + airport);
}
}
}
}