DATACOUCH-574 - Support countBy...() query method calls in repository.
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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<Airport, String> {
|
||||
|
||||
@@ -37,4 +43,8 @@ public interface AirportRepository extends PagingAndSortingRepository<Airport, S
|
||||
@Query("#{#n1ql.selectEntity} where iata = $1")
|
||||
List<Airport> getAllByIata(String iata);
|
||||
|
||||
long countByIataIn(String... iata);
|
||||
|
||||
long countByIcaoAndIataIn(String icao, String... iata);
|
||||
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<Airport, String> {
|
||||
|
||||
@@ -33,4 +40,7 @@ public interface ReactiveAirportRepository extends ReactiveSortingRepository<Air
|
||||
|
||||
Flux<Airport> findAllByIata(String iata);
|
||||
|
||||
Mono<Long> countByIataIn(String... iatas);
|
||||
|
||||
Mono<Long> countByIcaoAndIataIn(String icao, String... iatas);
|
||||
}
|
||||
|
||||
@@ -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<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] /* 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++) {
|
||||
|
||||
@@ -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<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);
|
||||
|
||||
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 {
|
||||
|
||||
Reference in New Issue
Block a user