Adopt to Reactor 2022 changes.

Closes #1283
This commit is contained in:
Mark Paluch
2022-07-04 14:18:42 +02:00
parent b5bd1c048a
commit f416c2c386
4 changed files with 73 additions and 33 deletions

View File

@@ -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<Object> executeLater(Object[] parameters) {
ReactiveCassandraParameterAccessor parameterAccessor = new ReactiveCassandraParameterAccessor(getQueryMethod(),
parameters);
Mono<ReactiveCassandraParameterAccessor> resolved = parameterAccessor.resolveParameters();
return resolved.flatMapMany(this::executeLater);
}
private Publisher<Object> executeLater(ReactiveCassandraParameterAccessor parameterAccessor) {
CassandraParameterAccessor convertingParameterAccessor = new ConvertingParameterAccessor(
getRequiredConverter(getReactiveCassandraOperations()), parameterAccessor);

View File

@@ -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<MonoProcessor<?>> 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> 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<ReactiveCassandraParameterAccessor> 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<Integer, Optional<?>> holder = new ConcurrentHashMap<>();
List<Publisher<?>> 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);
});
}
}

View File

@@ -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");
}
}

View File

@@ -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