diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java index d1e4a9e7d..810e31053 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java @@ -15,7 +15,6 @@ */ package org.springframework.data.cassandra.repository.query; -import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.reactivestreams.Publisher; @@ -72,14 +71,17 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository @Override public Object execute(Object[] parameters) { - return Flux.defer(() -> executeLater(parameters)); - } - - private Publisher executeLater(Object[] parameters) { ReactiveCassandraParameterAccessor parameterAccessor = new ReactiveCassandraParameterAccessor(getQueryMethod(), parameters); + Mono resolved = parameterAccessor.resolveParameters(); + + return resolved.flatMapMany(this::executeLater); + } + + private Publisher executeLater(ReactiveCassandraParameterAccessor parameterAccessor) { + CassandraParameterAccessor convertingParameterAccessor = new ConvertingParameterAccessor( getRequiredConverter(getReactiveCassandraOperations()), parameterAccessor); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java index c42c524dc..ebdd6ecae 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java @@ -17,10 +17,14 @@ package org.springframework.data.cassandra.repository.query; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import reactor.core.publisher.MonoProcessor; import java.util.ArrayList; import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; + +import org.reactivestreams.Publisher; import org.springframework.data.repository.util.ReactiveWrapperConverters; import org.springframework.data.repository.util.ReactiveWrappers; @@ -35,35 +39,13 @@ import org.springframework.data.repository.util.ReactiveWrappers; class ReactiveCassandraParameterAccessor extends CassandraParametersParameterAccessor { private final Object[] values; + private final CassandraQueryMethod method; - private final List> subscriptions; - - @SuppressWarnings("ConstantConditions") ReactiveCassandraParameterAccessor(CassandraQueryMethod method, Object[] values) { super(method, values); - + this.method = method; this.values = values; - this.subscriptions = new ArrayList<>(values.length); - - for (Object value : values) { - if (value == null || !ReactiveWrappers.supports(value.getClass())) { - subscriptions.add(null); - continue; - } - - if (ReactiveWrappers.isSingleValueType(value.getClass())) { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).toProcessor()); - } else { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().toProcessor()); - } - } - } - - @SuppressWarnings({ "unchecked", "ConstantConditions" }) - @Override - protected T getValue(int index) { - return (subscriptions.get(index) != null ? (T) subscriptions.get(index).block() : super.getValue(index)); } @Override @@ -81,4 +63,61 @@ class ReactiveCassandraParameterAccessor extends CassandraParametersParameterAcc public Object getBindableValue(int index) { return getValue(getParameters().getBindableParameter(index).getIndex()); } + + /** + * Resolve parameters that were provided through reactive wrapper types. Flux is collected into a list, values from + * Mono's are used directly. + * + * @return + */ + @SuppressWarnings("unchecked") + public Mono resolveParameters() { + + boolean hasReactiveWrapper = false; + + for (Object value : values) { + if (value == null || !ReactiveWrappers.supports(value.getClass())) { + continue; + } + + hasReactiveWrapper = true; + break; + } + + if (!hasReactiveWrapper) { + return Mono.just(this); + } + + Object[] resolved = new Object[values.length]; + Map> holder = new ConcurrentHashMap<>(); + List> publishers = new ArrayList<>(); + + for (int i = 0; i < values.length; i++) { + + Object value = resolved[i] = values[i]; + if (value == null || !ReactiveWrappers.supports(value.getClass())) { + continue; + } + + if (ReactiveWrappers.isSingleValueType(value.getClass())) { + + int index = i; + publishers.add(ReactiveWrapperConverters.toWrapper(value, Mono.class) // + .map(Optional::of) // + .defaultIfEmpty(Optional.empty()) // + .doOnNext(it -> holder.put(index, (Optional) it))); + } else { + + int index = i; + publishers.add(ReactiveWrapperConverters.toWrapper(value, Flux.class) // + .collectList() // + .doOnNext(it -> holder.put(index, Optional.of(it)))); + } + } + + return Flux.merge(publishers).then().thenReturn(resolved).map(values -> { + holder.forEach((index, v) -> values[index] = v.orElse(null)); + return new ReactiveCassandraParameterAccessor(method, values); + }); + } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslatorUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslatorUnitTests.java index d41cbb11f..cd6177932 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslatorUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslatorUnitTests.java @@ -241,7 +241,7 @@ class CassandraExceptionTranslatorUnitTests { InvalidConfigurationInQueryException cx = new InvalidConfigurationInQueryException(node, "err"); DataAccessException dax = sut.translate("Query", "SELECT * FROM person", cx); - assertThat(dax).hasRootCauseInstanceOf(InvalidConfigurationInQueryException.class).hasMessage( - "Query; CQL [SELECT * FROM person]; err; nested exception is com.datastax.oss.driver.api.core.servererrors.InvalidConfigurationInQueryException: err"); + assertThat(dax).hasRootCauseInstanceOf(InvalidConfigurationInQueryException.class) + .hasMessageContaining("Query; CQL [SELECT * FROM person]; err"); } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessorUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessorUnitTests.java index dd05f2f14..283d0c148 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessorUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessorUnitTests.java @@ -73,7 +73,6 @@ class ReactiveCassandraParameterAccessorUnitTests { getCassandraQueryMethod(method), new Object[] { Flux.just(LocalDateTime.of(2000, 10, 11, 12, 13, 14)) }); assertThat(accessor.getDataType(0)).isNull(); - } @Test // DATACASS-335