diff --git a/src/main/java/org/springframework/data/couchbase/repository/query/CouchbaseQueryMethod.java b/src/main/java/org/springframework/data/couchbase/repository/query/CouchbaseQueryMethod.java index 0cbe6f94..43223010 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/query/CouchbaseQueryMethod.java +++ b/src/main/java/org/springframework/data/couchbase/repository/query/CouchbaseQueryMethod.java @@ -163,6 +163,16 @@ public class CouchbaseQueryMethod extends QueryMethod { return StringUtils.hasText(query) ? query : null; } + /** + * indicates if the method begins with "count" + * + * @return true if the method begins with "count", indicating that .count() should be called instead of one() or + * all(). + */ + public boolean isCountQuery() { + return getName().toLowerCase().startsWith("count"); + } + @Override public String toString() { return super.toString(); diff --git a/src/main/java/org/springframework/data/couchbase/repository/query/N1qlRepositoryQueryExecutor.java b/src/main/java/org/springframework/data/couchbase/repository/query/N1qlRepositoryQueryExecutor.java index b90a0b2a..d9d3eed3 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/query/N1qlRepositoryQueryExecutor.java +++ b/src/main/java/org/springframework/data/couchbase/repository/query/N1qlRepositoryQueryExecutor.java @@ -59,15 +59,16 @@ public class N1qlRepositoryQueryExecutor { Query query; ExecutableFindByQueryOperation.ExecutableFindByQuery q; if (queryMethod.hasN1qlAnnotation()) { - query = new StringN1qlQueryCreator(accessor, queryMethod, operations.getConverter(), - operations.getBucketName(), QueryMethodEvaluationContextProvider.DEFAULT, - namedQueries).createQuery(); + query = new StringN1qlQueryCreator(accessor, queryMethod, operations.getConverter(), operations.getBucketName(), + QueryMethodEvaluationContextProvider.DEFAULT, namedQueries).createQuery(); } else { 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); - if (queryMethod.isCollectionQuery()) { + if (queryMethod.isCountQuery()) { + return q.count(); + } else if (queryMethod.isCollectionQuery()) { return q.all(); } else { return q.oneValue(); diff --git a/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveN1qlRepositoryQueryExecutor.java b/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveN1qlRepositoryQueryExecutor.java index e000f25c..2a540980 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveN1qlRepositoryQueryExecutor.java +++ b/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveN1qlRepositoryQueryExecutor.java @@ -63,15 +63,16 @@ public class ReactiveN1qlRepositoryQueryExecutor { Query query; ReactiveFindByQueryOperation.ReactiveFindByQuery q; if (queryMethod.hasN1qlAnnotation()) { - query = new StringN1qlQueryCreator(accessor, queryMethod, operations.getConverter(), - operations.getBucketName(), QueryMethodEvaluationContextProvider.DEFAULT, - namedQueries).createQuery(); + query = new StringN1qlQueryCreator(accessor, queryMethod, operations.getConverter(), operations.getBucketName(), + QueryMethodEvaluationContextProvider.DEFAULT, namedQueries).createQuery(); } else { 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); - if (queryMethod.isCollectionQuery()) { + if (queryMethod.isCountQuery()) { + return q.count(); + } else if (queryMethod.isCollectionQuery()) { return q.all(); } else { return q.one(); 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 36c82c94..c9dba069 100644 --- a/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java +++ b/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java @@ -25,6 +25,12 @@ import org.springframework.stereotype.Repository; import com.couchbase.client.java.query.QueryScanConsistency; +/** + * template class for Reactive Couchbase operations + * + * @author Michael Nitschinger + * @author Michael Reiche + */ @Repository public interface AirportRepository extends PagingAndSortingRepository { @@ -37,4 +43,8 @@ public interface AirportRepository extends PagingAndSortingRepository getAllByIata(String iata); + long countByIataIn(String... iata); + + long countByIcaoAndIataIn(String icao, String... iata); + } diff --git a/src/test/java/org/springframework/data/couchbase/domain/Config.java b/src/test/java/org/springframework/data/couchbase/domain/Config.java index b548ff49..6fd05245 100644 --- a/src/test/java/org/springframework/data/couchbase/domain/Config.java +++ b/src/test/java/org/springframework/data/couchbase/domain/Config.java @@ -69,7 +69,7 @@ public class Config extends AbstractCouchbaseConfiguration { } } - String clusterGet( String methodName, String defaultValue){ + String clusterGet(String methodName, String defaultValue) { if (clusterAware != null) { try { return (String) clusterAware.getMethod(methodName).invoke(null); @@ -82,22 +82,22 @@ public class Config extends AbstractCouchbaseConfiguration { @Override public String getConnectionString() { - return clusterGet( "connectionString", connectionString ); + return clusterGet("connectionString", connectionString); } @Override public String getUserName() { - return clusterGet( "username", username ); + return clusterGet("username", username); } @Override public String getPassword() { - return clusterGet( "password", password ); + return clusterGet("password", password); } @Override public String getBucketName() { - return clusterGet( "bucketName", bucketname ); + return clusterGet("bucketName", bucketname); } @Bean(name = "auditorAwareRef") @@ -113,23 +113,30 @@ public class Config extends AbstractCouchbaseConfiguration { @Override public void configureReactiveRepositoryOperationsMapping(ReactiveRepositoryOperationsMapping baseMapping) { try { - ReactiveCouchbaseTemplate personTemplate = myReactiveCouchbaseTemplate(myCouchbaseClientFactory("protected"),new MappingCouchbaseConverter()); - baseMapping.mapEntity(Person.class, personTemplate); // Person goes in "protected" bucket - ReactiveCouchbaseTemplate userTemplate = myReactiveCouchbaseTemplate(myCouchbaseClientFactory("mybucket"),new MappingCouchbaseConverter()); - baseMapping.mapEntity(User.class, userTemplate); // User goes in "mybucket" - // everything else goes in getBucketName() ( which is travel-sample ) + // comment out references to 'protected' and 'mybucket' - they are only to show how multi-bucket would work + // ReactiveCouchbaseTemplate personTemplate = + // myReactiveCouchbaseTemplate(myCouchbaseClientFactory("protected"),new MappingCouchbaseConverter()); + // baseMapping.mapEntity(Person.class, personTemplate); // Person goes in "protected" bucket + // ReactiveCouchbaseTemplate userTemplate = myReactiveCouchbaseTemplate(myCouchbaseClientFactory("mybucket"),new + // MappingCouchbaseConverter()); + // baseMapping.mapEntity(User.class, userTemplate); // User goes in "mybucket" + // everything else goes in getBucketName() ( which is travel-sample ) } catch (Exception e) { throw e; } } + @Override public void configureRepositoryOperationsMapping(RepositoryOperationsMapping baseMapping) { try { - CouchbaseTemplate personTemplate = myCouchbaseTemplate(myCouchbaseClientFactory("protected"),new MappingCouchbaseConverter()); - baseMapping.mapEntity(Person.class, personTemplate); // Person goes in "protected" bucket - CouchbaseTemplate userTemplate = myCouchbaseTemplate(myCouchbaseClientFactory("mybucket"),new MappingCouchbaseConverter()); - baseMapping.mapEntity(User.class, userTemplate); // User goes in "mybucket" - // everything else goes in getBucketName() ( which is travel-sample ) + // comment out references to 'protected' and 'mybucket' - they are only to show how multi-bucket would work + // CouchbaseTemplate personTemplate = myCouchbaseTemplate(myCouchbaseClientFactory("protected"),new + // MappingCouchbaseConverter()); + // baseMapping.mapEntity(Person.class, personTemplate); // Person goes in "protected" bucket + // CouchbaseTemplate userTemplate = myCouchbaseTemplate(myCouchbaseClientFactory("mybucket"),new + // MappingCouchbaseConverter()); + // baseMapping.mapEntity(User.class, userTemplate); // User goes in "mybucket" + // everything else goes in getBucketName() ( which is travel-sample ) } catch (Exception e) { throw e; } @@ -152,7 +159,7 @@ public class Config extends AbstractCouchbaseConfiguration { // do not use couchbaseClientFactory for the name of this method, otherwise the value of that bean will // will be used instead of this call being made ( bucketname is an arg here, instead of using bucketName() ) public CouchbaseClientFactory myCouchbaseClientFactory(String bucketName) { - return new SimpleCouchbaseClientFactory(getConnectionString(),authenticator(), bucketName ); + return new SimpleCouchbaseClientFactory(getConnectionString(), authenticator(), bucketName); } // convenience constructor for tests diff --git a/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java b/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java index 16008071..8f6dc7bb 100644 --- a/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java +++ b/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java @@ -21,9 +21,16 @@ import reactor.core.publisher.Flux; import org.springframework.data.couchbase.repository.ScanConsistency; import org.springframework.data.repository.reactive.ReactiveSortingRepository; import org.springframework.stereotype.Repository; +import reactor.core.publisher.Mono; import com.couchbase.client.java.query.QueryScanConsistency; +/** + * template class for Reactive Couchbase operations + * + * @author Michael Nitschinger + * @author Michael Reiche + */ @Repository public interface ReactiveAirportRepository extends ReactiveSortingRepository { @@ -33,4 +40,7 @@ public interface ReactiveAirportRepository extends ReactiveSortingRepository findAllByIata(String iata); + Mono countByIataIn(String... iatas); + + Mono countByIcaoAndIataIn(String icao, String... iatas); } 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 157bb9af..9834f726 100644 --- a/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java +++ b/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java @@ -104,6 +104,15 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr airportCount = airportRepository.count(); assertEquals(7, airportCount); + airportCount = airportRepository.countByIataIn("JFK", "IAD", "SFO"); + assertEquals(3, airportCount); + + airportCount = airportRepository.countByIcaoAndIataIn("jfk", "JFK", "IAD", "SFO", "XXX"); + assertEquals(1, airportCount); + + airportCount = airportRepository.countByIataIn("XXX"); + assertEquals(0, airportCount); + } finally { for (int i = 0; i < iatas.length; i++) { Airport airport = new Airport("airports::" + iatas[i], iatas[i] /*iata*/, iatas[i] /* lcao */); @@ -113,7 +122,7 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr } @Test - void threadSafeParametersTest() { + void threadSafeParametersTest() throws Exception { String[] iatas = { "JFK", "IAD", "SFO", "SJC", "SEA", "LAX", "PHX" }; Future[] future = new Future[iatas.length]; ExecutorService executorService = Executors.newFixedThreadPool(iatas.length); @@ -121,7 +130,7 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr 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] /* lcao */); + Airport airport = new Airport("airports::" + iatas[i], iatas[i] /*iata*/, iatas[i].toLowerCase() /* lcao */); airportRepository.save(airport); final int idx = i; suppliers[i] = () -> { @@ -143,8 +152,7 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr for (int i = 0; i < iatas.length; i++) { future[i].get(); } - } catch (Exception e) { - throw new RuntimeException(e); + } finally { executorService.shutdown(); for (int i = 0; i < iatas.length; i++) { @@ -155,7 +163,7 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr } @Test - void threadSafeStringParametersTest() { + void threadSafeStringParametersTest() throws Exception { String[] iatas = { "JFK", "IAD", "SFO", "SJC", "SEA", "LAX", "PHX" }; Future[] future = new Future[iatas.length]; ExecutorService executorService = Executors.newFixedThreadPool(iatas.length); @@ -185,8 +193,6 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr for (int i = 0; i < iatas.length; i++) { future[i].get(); } - } catch (Exception e) { - throw new RuntimeException(e); } finally { executorService.shutdown(); for (int i = 0; i < iatas.length; i++) { 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 09bf4465..5275e7e8 100644 --- a/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java +++ b/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java @@ -19,6 +19,10 @@ package org.springframework.data.couchbase.repository; import static org.junit.jupiter.api.Assertions.*; import java.util.List; +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 org.junit.jupiter.api.BeforeEach; @@ -38,6 +42,12 @@ import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; import com.couchbase.client.core.error.IndexExistsException; +/** + * template class for Reactive Couchbase operations + * + * @author Michael Nitschinger + * @author Michael Reiche + */ @SpringJUnitConfig(ReactiveCouchbaseRepositoryQueryIntegrationTests.Config.class) @IgnoreWhen(missesCapabilities = Capabilities.QUERY, clusterTypes = ClusterType.MOCKED) public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegrationTests { @@ -73,6 +83,42 @@ public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends ClusterAwa System.err.println(airports); } + @Test + void count() { + 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); + + airportCount = airportRepository.countByIataIn("JFK", "IAD", "SFO").block(); + assertEquals(3, airportCount); + + airportCount = airportRepository.countByIcaoAndIataIn("jfk", "JFK", "IAD", "SFO", "XXX").block(); + assertEquals(1, airportCount); + + airportCount = airportRepository.countByIataIn("XXX").block(); + assertEquals(0, airportCount); + + } 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); + } + } + } + @Configuration @EnableReactiveCouchbaseRepositories("org.springframework.data.couchbase") static class Config extends AbstractCouchbaseConfiguration {