diff --git a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByQueryOperationSupport.java b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByQueryOperationSupport.java index 48d77f02..4d5b2e93 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByQueryOperationSupport.java +++ b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByQueryOperationSupport.java @@ -19,6 +19,7 @@ import java.util.List; import java.util.stream.Stream; import org.springframework.data.couchbase.core.ReactiveFindByQueryOperationSupport.ReactiveFindByQuerySupport; +import org.springframework.data.couchbase.core.CouchbaseQueryExecutionException; import org.springframework.data.couchbase.core.query.Query; import org.springframework.util.Assert; @@ -142,7 +143,11 @@ public class ExecutableFindByQueryOperationSupport implements ExecutableFindByQu @Override public long count() { - return reactiveSupport.count().block(); + Long l = reactiveSupport.count().block(); + if ( l == null ){ + throw new CouchbaseQueryExecutionException("count query did not return a count : "+query.export()); + } + return l; } @Override diff --git a/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByQueryOperationSupport.java b/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByQueryOperationSupport.java index d4b88e03..696c335f 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByQueryOperationSupport.java +++ b/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByQueryOperationSupport.java @@ -217,7 +217,8 @@ public class ReactiveFindByQueryOperationSupport implements ReactiveFindByQueryO } else { return throwable; } - }).flatMapMany(ReactiveQueryResult::rowsAsObject).map(row -> row.getLong(TemplateUtils.SELECT_COUNT)).next()); + }).flatMapMany(ReactiveQueryResult::rowsAsObject).map(row -> row.getLong(row.getNames().iterator().next())) + .next()); } @Override 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 589c55ce..adcbadb6 100644 --- a/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java +++ b/src/test/java/org/springframework/data/couchbase/domain/AirportRepository.java @@ -114,6 +114,12 @@ public interface AirportRepository extends CouchbaseRepository, Long countFancyExpression(@Param("projectIds") List projectIds, @Param("planIds") List planIds, @Param("active") Boolean active); + @Query("SELECT 1 FROM `#{#n1ql.bucket}` WHERE 0 = 1" ) + Long countBad(); + + @Query("SELECT count(*) FROM `#{#n1ql.bucket}`" ) + Long countGood(); + @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) Page findAllByIataNot(String iata, Pageable pageable); 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 49d12f4b..3ae45165 100644 --- a/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java +++ b/src/test/java/org/springframework/data/couchbase/repository/CouchbaseRepositoryQueryIntegrationTests.java @@ -50,6 +50,7 @@ import org.springframework.dao.DataRetrievalFailureException; import org.springframework.data.auditing.DateTimeProvider; import org.springframework.data.couchbase.CouchbaseClientFactory; import org.springframework.data.couchbase.config.AbstractCouchbaseConfiguration; +import org.springframework.data.couchbase.core.CouchbaseQueryExecutionException; import org.springframework.data.couchbase.core.CouchbaseTemplate; import org.springframework.data.couchbase.core.RemoveResult; import org.springframework.data.couchbase.core.query.N1QLExpression; @@ -474,6 +475,16 @@ public class CouchbaseRepositoryQueryIntegrationTests extends ClusterAwareIntegr } } + @Test + void badCount(){ + assertThrows(CouchbaseQueryExecutionException.class, () -> airportRepository.countBad()); + } + + @Test + void goodCount(){ + airportRepository.countGood(); + } + @Test void threadSafeParametersTest() throws Exception { String[] iatas = { "JFK", "IAD", "SFO", "SJC", "SEA", "LAX", "PHX" };