From 2cb5508fa998b347d763d743829d51f3b31e78b4 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Mon, 11 Jul 2022 08:22:51 +0200 Subject: [PATCH] Adopt to Reactor 2022.0.0 changes. Closes #1490 --- .../query/AbstractCouchbaseQueryBase.java | 13 +-- .../query/ReactiveAbstractN1qlBasedQuery.java | 5 +- .../ReactiveCouchbaseParameterAccessor.java | 102 +++++++++++++----- .../domain/ReactiveAirportRepository.java | 2 +- ...chbaseRepositoryQueryIntegrationTests.java | 6 +- 5 files changed, 88 insertions(+), 40 deletions(-) diff --git a/src/main/java/org/springframework/data/couchbase/repository/query/AbstractCouchbaseQueryBase.java b/src/main/java/org/springframework/data/couchbase/repository/query/AbstractCouchbaseQueryBase.java index b00101e5..38c05bf6 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/query/AbstractCouchbaseQueryBase.java +++ b/src/main/java/org/springframework/data/couchbase/repository/query/AbstractCouchbaseQueryBase.java @@ -38,7 +38,7 @@ import org.springframework.util.Assert; /** * {@link RepositoryQuery} implementation for Couchbase. CouchbaseOperationsType is either CouchbaseOperations or * ReactiveCouchbaseOperations - * + * * @author Michael Reiche * @since 4.1 */ @@ -105,16 +105,17 @@ public abstract class AbstractCouchbaseQueryBase implem /** * Execute the query with the provided parameters - * + * * @see org.springframework.data.repository.query.RepositoryQuery#execute(java.lang.Object[]) */ public Object execute(Object[] parameters) { - return method.hasReactiveWrapperParameter() ? executeDeferred(parameters) - : execute(new ReactiveCouchbaseParameterAccessor(getQueryMethod(), parameters)); + + ReactiveCouchbaseParameterAccessor accessor = new ReactiveCouchbaseParameterAccessor(getQueryMethod(), parameters); + + return accessor.resolveParameters().flatMapMany(this::executeDeferred); } - private Object executeDeferred(Object[] parameters) { - ReactiveCouchbaseParameterAccessor parameterAccessor = new ReactiveCouchbaseParameterAccessor(method, parameters); + private Publisher executeDeferred(ReactiveCouchbaseParameterAccessor parameterAccessor) { if (getQueryMethod().isCollectionQuery()) { return Flux.defer(() -> (Publisher) execute(parameterAccessor)); } diff --git a/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveAbstractN1qlBasedQuery.java b/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveAbstractN1qlBasedQuery.java index 8453d56a..17cba8ef 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveAbstractN1qlBasedQuery.java +++ b/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveAbstractN1qlBasedQuery.java @@ -58,6 +58,8 @@ public abstract class ReactiveAbstractN1qlBasedQuery implements RepositoryQuery @Override public Object execute(Object[] parameters) { ReactiveCouchbaseParameterAccessor accessor = new ReactiveCouchbaseParameterAccessor(queryMethod, parameters); + + return accessor.resolveParameters().flatMapMany(it -> { ResultProcessor processor = this.queryMethod.getResultProcessor().withDynamicProjection(accessor); ReturnedType returnedType = processor.getReturnedType(); @@ -71,6 +73,7 @@ public abstract class ReactiveAbstractN1qlBasedQuery implements RepositoryQuery N1QLQuery query = N1qlUtils.buildQuery(expression, queryPlaceholderValues, getScanConsistency()); return ReactiveWrapperConverters .toWrapper(processor.processResult(executeDependingOnType(query, queryMethod, typeToRead)), Flux.class); + }); } protected Object executeDependingOnType(N1QLQuery query, QueryMethod queryMethod, Class typeToRead) { @@ -130,7 +133,7 @@ public abstract class ReactiveAbstractN1qlBasedQuery implements RepositoryQuery if (queryMethod.hasConsistencyAnnotation()) { return queryMethod.getConsistencyAnnotation().value(); } - + return getCouchbaseOperations().getDefaultConsistency().n1qlConsistency();*/ } } diff --git a/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveCouchbaseParameterAccessor.java b/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveCouchbaseParameterAccessor.java index c5c98a3c..fd7df263 100644 --- a/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveCouchbaseParameterAccessor.java +++ b/src/main/java/org/springframework/data/couchbase/repository/query/ReactiveCouchbaseParameterAccessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2017-2020 the original author or authors. + * Copyright 2017-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -17,10 +17,14 @@ package org.springframework.data.couchbase.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.query.ParametersParameterAccessor; import org.springframework.data.repository.util.ReactiveWrapperConverters; @@ -36,47 +40,87 @@ import org.springframework.data.repository.util.ReactiveWrappers; */ public class ReactiveCouchbaseParameterAccessor extends ParametersParameterAccessor { - private final List> subscriptions; + private final CouchbaseQueryMethod method; + private final Object[] values; public ReactiveCouchbaseParameterAccessor(CouchbaseQueryMethod method, Object[] values) { + super(method.getParameters(), values); - this.subscriptions = new ArrayList<>(values.length); + this.method = method; + this.values = values; + } + + /* (non-Javadoc) + * @see org.springframework.data.mongodb.repository.query.MongoParametersParameterAccessor#getValues() + */ + @Override + public Object[] getValues() { + + Object[] result = new Object[super.getValues().length]; + for (int i = 0; i < result.length; i++) { + result[i] = getValue(i); + } + return result; + } + + 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 = values[i]; - + Object value = resolved[i] = values[i]; if (value == null || !ReactiveWrappers.supports(value.getClass())) { - subscriptions.add(null); continue; } if (ReactiveWrappers.isSingleValueType(value.getClass())) { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).toProcessor()); + + int index = i; + publishers.add(ReactiveWrapperConverters.toWrapper(value, Mono.class) // + .map(Optional::of) // + .defaultIfEmpty(Optional.empty()) // + .doOnNext(it -> holder.put(index, (Optional) it))); } else { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().toProcessor()); + + int index = i; + publishers.add(ReactiveWrapperConverters.toWrapper(value, Flux.class) // + .collectList() // + .doOnNext(it -> holder.put(index, Optional.of(it)))); } } - } - /* (non-Javadoc) - * @see org.springframework.data.repository.query.ParametersParameterAccessor#getValue(int) - */ - @SuppressWarnings("unchecked") - @Override - protected T getValue(int index) { - - if (subscriptions.get(index) != null) { - return (T) subscriptions.get(index).block(); - } - - return super.getValue(index); - } - - /* (non-Javadoc) - * @see org.springframework.data.repository.query.ParametersParameterAccessor#getBindableValue(int) - */ - public Object getBindableValue(int index) { - return getValue(getParameters().getBindableParameter(index).getIndex()); + return Flux.merge(publishers).then().thenReturn(resolved).map(values -> { + holder.forEach((index, v) -> values[index] = v.orElse(null)); + return new ReactiveCouchbaseParameterAccessor(method, values); + }); } } diff --git a/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java b/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java index 786c4ea7..5c295e35 100644 --- a/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java +++ b/src/test/java/org/springframework/data/couchbase/domain/ReactiveAirportRepository.java @@ -62,7 +62,7 @@ public interface ReactiveAirportRepository Mono save(Airport a); @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) - Flux findAllByIata(String iata); + Flux findAllByIata(Mono iata); @ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS) Mono iata(String iata); diff --git a/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java b/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java index 815f9a4c..bc763cec 100644 --- a/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java +++ b/src/test/java/org/springframework/data/couchbase/repository/ReactiveCouchbaseRepositoryQueryIntegrationTests.java @@ -124,12 +124,12 @@ public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends JavaIntegr try { vie = new Airport("airports::vie", "vie", "low2"); reactiveAirportRepository.save(vie).block(); - List airports1 = reactiveAirportRepository.findAllByIata("vie").collectList().block(); + List airports1 = reactiveAirportRepository.findAllByIata(Mono.just("vie")).collectList().block(); assertEquals(1, airports1.size()); - List airports2 = reactiveAirportRepository.findAllByIata("vie").collectList().block(); + List airports2 = reactiveAirportRepository.findAllByIata(Mono.just("vie")).collectList().block(); assertEquals(1, airports2.size()); vie = reactiveAirportRepository.save(vie).block(); - List airports = reactiveAirportRepository.findAllByIata("vie").collectList().block(); + List airports = reactiveAirportRepository.findAllByIata(Mono.just("vie")).collectList().block(); assertEquals(1, airports.size()); Airport airport1 = reactiveAirportRepository.findById(airports.get(0).getId()).block(); assertEquals(airport1.getIata(), vie.getIata());