From 9758ad980251c2a9a25ff7b2516ed560ddb75b97 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 1 Jun 2017 10:01:39 +0200 Subject: [PATCH] DATACOUCH-311 - Upgrade to Reactor 3.1 M2. Adopt to API change from Publisher.subscribe() to Publisher.toProcessor(). --- .../query/ReactiveCouchbaseParameterAccessor.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) 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 a7bb7aab..1e412aee 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 @@ -31,16 +31,15 @@ import reactor.core.publisher.MonoProcessor; * to reactive parameter wrapper types upon creation. This class performs synchronization when accessing parameters. * * @author Subhashni Balakrishnan + * @author Mark Paluch * @since 3.0 */ public class ReactiveCouchbaseParameterAccessor extends ParametersParameterAccessor { - private final Object[] values; private final List> subscriptions; public ReactiveCouchbaseParameterAccessor(CouchbaseQueryMethod method, Object[] values) { super(method.getParameters(), values); - this.values = values; this.subscriptions = new ArrayList<>(values.length); for (int i = 0; i < values.length; i++) { @@ -53,9 +52,9 @@ public class ReactiveCouchbaseParameterAccessor extends ParametersParameterAcces } if (ReactiveWrappers.isSingleValueType(value.getClass())) { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).subscribe()); + subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).toProcessor()); } else { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().subscribe()); + subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().toProcessor()); } } }