From 05fb9eadbed210ff6bedb7102733dcf8de85d9b5 Mon Sep 17 00:00:00 2001 From: mikereiche Date: Thu, 16 Jul 2020 17:31:37 -0700 Subject: [PATCH] DATACOUCH-585 - Support ScanConsistency for count() queries. --- .../SimpleReactiveCouchbaseRepository.java | 2 +- .../couchbase/domain/AirportRepository.java | 12 ++++++++++++ .../domain/ReactiveAirportRepository.java | 14 ++++++++++++++ ...ouchbaseRepositoryQueryIntegrationTests.java | 5 ++--- ...ouchbaseRepositoryQueryIntegrationTests.java | 17 +++++++++-------- 5 files changed, 38 insertions(+), 12 deletions(-) diff --git a/src/main/java/org/springframework/data/couchbase/repository/support/SimpleReactiveCouchbaseRepository.java b/src/main/java/org/springframework/data/couchbase/repository/support/SimpleReactiveCouchbaseRepository.java index 6acb8f60..2565933c 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/support/SimpleReactiveCouchbaseRepository.java +++ b/src/main/java/org/springframework/data/couchbase/repository/support/SimpleReactiveCouchbaseRepository.java @@ -181,7 +181,7 @@ public class SimpleReactiveCouchbaseRepository implements ReactiveCouchba @SuppressWarnings("unchecked") @Override public Mono count() { - return operations.findByQuery(entityInformation.getJavaType()).count(); + return operations.findByQuery(entityInformation.getJavaType()).consistentWith(buildQueryScanConsistency()).count(); } @SuppressWarnings("unchecked") diff --git a/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java b/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java index e3590555..06ce1fec 100644 --- a/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java +++ b/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java @@ -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 findAll(); + @Override + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) + Airport save(Airport airport); + + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) List findAllByIata(String iata); @Query("#{#n1ql.selectEntity} where iata = $1") @@ -45,8 +51,14 @@ public interface AirportRepository extends PagingAndSortingRepository findAll(); + @Override + Mono save(Airport a); + + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) Flux findAllByIata(String iata); + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) Mono countByIataIn(String... iatas); + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) Mono countByIcaoAndIataIn(String icao, String... iatas); + + @Override + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) + Mono count(); + + @Override + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) + Mono findById(String var1); } diff --git a/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java b/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java index ff31d479..c4ea6cd8 100644 --- a/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java +++ b/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java @@ -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"); diff --git a/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java b/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java index f143564c..dd04aeb9 100644 --- a/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java +++ b/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java @@ -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[] 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); + } } } }