diff --git a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperation.java b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperation.java index 3fa427ab..ebe4e9a7 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperation.java +++ b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperation.java @@ -19,6 +19,8 @@ import java.util.List; import java.util.Optional; import java.util.stream.Stream; +import com.couchbase.client.java.analytics.AnalyticsScanConsistency; +import com.couchbase.client.java.query.QueryScanConsistency; import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.data.couchbase.core.query.AnalyticsQuery; import org.springframework.lang.Nullable; @@ -112,6 +114,17 @@ public interface ExecutableFindByAnalyticsOperation { } - interface ExecutableFindByAnalytics extends FindByAnalyticsWithQuery {} + interface FindByAnalyticsConsistentWith extends FindByAnalyticsWithQuery { + + /** + * Allows to override the default scan consistency. + * + * @param scanConsistency the custom scan consistency to use for this analytics query. + */ + FindByAnalyticsWithQuery consistentWith(AnalyticsScanConsistency scanConsistency); + + } + + interface ExecutableFindByAnalytics extends FindByAnalyticsConsistentWith {} } diff --git a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperationSupport.java b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperationSupport.java index d1380986..0be1cb59 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperationSupport.java +++ b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByAnalyticsOperationSupport.java @@ -18,7 +18,10 @@ package org.springframework.data.couchbase.core; import java.util.List; import java.util.stream.Stream; +import com.couchbase.client.java.analytics.AnalyticsScanConsistency; +import com.couchbase.client.java.query.QueryScanConsistency; import org.springframework.data.couchbase.core.query.AnalyticsQuery; +import org.springframework.data.couchbase.core.query.Query; public class ExecutableFindByAnalyticsOperationSupport implements ExecutableFindByAnalyticsOperation { @@ -32,7 +35,7 @@ public class ExecutableFindByAnalyticsOperationSupport implements ExecutableFind @Override public ExecutableFindByAnalytics findByAnalytics(final Class domainType) { - return new ExecutableFindByAnalyticsSupport<>(template, domainType, ALL_QUERY); + return new ExecutableFindByAnalyticsSupport<>(template, domainType, ALL_QUERY, AnalyticsScanConsistency.NOT_BOUNDED); } static class ExecutableFindByAnalyticsSupport implements ExecutableFindByAnalytics { @@ -40,13 +43,17 @@ public class ExecutableFindByAnalyticsOperationSupport implements ExecutableFind private final CouchbaseTemplate template; private final Class domainType; private final ReactiveFindByAnalyticsOperationSupport.ReactiveFindByAnalyticsSupport reactiveSupport; + private final AnalyticsQuery query; + private final AnalyticsScanConsistency scanConsistency; ExecutableFindByAnalyticsSupport(final CouchbaseTemplate template, final Class domainType, - final AnalyticsQuery query) { + final AnalyticsQuery query, final AnalyticsScanConsistency scanConsistency) { this.template = template; this.domainType = domainType; + this.query = query; this.reactiveSupport = new ReactiveFindByAnalyticsOperationSupport.ReactiveFindByAnalyticsSupport<>( - template.reactive(), domainType, query); + template.reactive(), domainType, query, scanConsistency); + this.scanConsistency = scanConsistency; } @Override @@ -66,7 +73,12 @@ public class ExecutableFindByAnalyticsOperationSupport implements ExecutableFind @Override public TerminatingFindByAnalytics matching(final AnalyticsQuery query) { - return new ExecutableFindByAnalyticsSupport<>(template, domainType, query); + return new ExecutableFindByAnalyticsSupport<>(template, domainType, query, scanConsistency); + } + + @Override + public FindByAnalyticsWithQuery consistentWith(final AnalyticsScanConsistency scanConsistency) { + return new ExecutableFindByAnalyticsSupport<>(template, domainType, query, scanConsistency); } @Override 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 22dd038e..fa4ce567 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByQueryOperationSupport.java +++ b/src/main/java/org/springframework/data/couchbase/core/ExecutableFindByQueryOperationSupport.java @@ -50,7 +50,7 @@ public class ExecutableFindByQueryOperationSupport implements ExecutableFindByQu this.template = template; this.domainType = domainType; this.query = query; - this.reactiveSupport = new ReactiveFindByQueryOperationSupport.ReactiveFindByQuerySupport(template.reactive(), + this.reactiveSupport = new ReactiveFindByQueryOperationSupport.ReactiveFindByQuerySupport<>(template.reactive(), domainType, query, scanConsistency); this.scanConsistency = scanConsistency; } diff --git a/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperation.java b/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperation.java index 32bbf96d..ff9f65fb 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperation.java +++ b/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperation.java @@ -15,6 +15,7 @@ */ package org.springframework.data.couchbase.core; +import com.couchbase.client.java.analytics.AnalyticsScanConsistency; import org.springframework.dao.IncorrectResultSizeDataAccessException; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -87,6 +88,17 @@ public interface ReactiveFindByAnalyticsOperation { } - interface ReactiveFindByAnalytics extends FindByAnalyticsWithQuery {} + interface FindByAnalyticsConsistentWith extends FindByAnalyticsWithQuery { + + /** + * Allows to override the default scan consistency. + * + * @param scanConsistency the custom scan consistency to use for this analytics query. + */ + FindByAnalyticsWithQuery consistentWith(AnalyticsScanConsistency scanConsistency); + + } + + interface ReactiveFindByAnalytics extends FindByAnalyticsConsistentWith {} } diff --git a/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperationSupport.java b/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperationSupport.java index cdb0336c..3aaf6e96 100644 --- a/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperationSupport.java +++ b/src/main/java/org/springframework/data/couchbase/core/ReactiveFindByAnalyticsOperationSupport.java @@ -15,6 +15,9 @@ */ package org.springframework.data.couchbase.core; +import com.couchbase.client.java.analytics.AnalyticsOptions; +import com.couchbase.client.java.analytics.AnalyticsScanConsistency; +import com.couchbase.client.java.query.QueryOptions; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -34,7 +37,7 @@ public class ReactiveFindByAnalyticsOperationSupport implements ReactiveFindByAn @Override public ReactiveFindByAnalytics findByAnalytics(final Class domainType) { - return new ReactiveFindByAnalyticsSupport<>(template, domainType, ALL_QUERY); + return new ReactiveFindByAnalyticsSupport<>(template, domainType, ALL_QUERY, AnalyticsScanConsistency.NOT_BOUNDED); } static class ReactiveFindByAnalyticsSupport implements ReactiveFindByAnalytics { @@ -42,17 +45,24 @@ public class ReactiveFindByAnalyticsOperationSupport implements ReactiveFindByAn private final ReactiveCouchbaseTemplate template; private final Class domainType; private final AnalyticsQuery query; + private final AnalyticsScanConsistency scanConsistency; ReactiveFindByAnalyticsSupport(final ReactiveCouchbaseTemplate template, final Class domainType, - final AnalyticsQuery query) { + final AnalyticsQuery query, final AnalyticsScanConsistency scanConsistency) { this.template = template; this.domainType = domainType; this.query = query; + this.scanConsistency = scanConsistency; } @Override public TerminatingFindByAnalytics matching(AnalyticsQuery query) { - return new ReactiveFindByAnalyticsSupport<>(template, domainType, query); + return new ReactiveFindByAnalyticsSupport<>(template, domainType, query, scanConsistency); + } + + @Override + public FindByAnalyticsWithQuery consistentWith(AnalyticsScanConsistency scanConsistency) { + return new ReactiveFindByAnalyticsSupport<>(template, domainType, query, scanConsistency); } @Override @@ -69,7 +79,7 @@ public class ReactiveFindByAnalyticsOperationSupport implements ReactiveFindByAn public Flux all() { return Flux.defer(() -> { String statement = assembleEntityQuery(false); - return template.getCouchbaseClientFactory().getCluster().reactive().analyticsQuery(statement) + return template.getCouchbaseClientFactory().getCluster().reactive().analyticsQuery(statement, buildAnalyticsOptions()) .onErrorMap(throwable -> { if (throwable instanceof RuntimeException) { return template.potentiallyConvertRuntimeException((RuntimeException) throwable); @@ -90,7 +100,7 @@ public class ReactiveFindByAnalyticsOperationSupport implements ReactiveFindByAn public Mono count() { return Mono.defer(() -> { String statement = assembleEntityQuery(true); - return template.getCouchbaseClientFactory().getCluster().reactive().analyticsQuery(statement) + return template.getCouchbaseClientFactory().getCluster().reactive().analyticsQuery(statement, buildAnalyticsOptions()) .onErrorMap(throwable -> { if (throwable instanceof RuntimeException) { return template.potentiallyConvertRuntimeException((RuntimeException) throwable); @@ -123,6 +133,14 @@ public class ReactiveFindByAnalyticsOperationSupport implements ReactiveFindByAn query.appendSkipAndLimit(statement); return statement.toString(); } + + private AnalyticsOptions buildAnalyticsOptions() { + final AnalyticsOptions options = AnalyticsOptions.analyticsOptions(); + if (scanConsistency != null) { + options.scanConsistency(scanConsistency); + } + return options; + } } }