From bfab233d2f04341de2faa7552d92f4cf773d4211 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Tue, 4 Aug 2020 13:34:50 +0200 Subject: [PATCH] DATAMONGO-2603 - Adopt to Reactor 3.4 changes. Align with ContextView and changes in other operators. --- .../mongodb/core/ReactiveMongoContext.java | 20 ++++++++++++++----- .../mongodb/core/ReactiveMongoTemplate.java | 2 +- .../mongodb/test/util/MongoTestUtils.java | 3 ++- 3 files changed, 18 insertions(+), 7 deletions(-) diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoContext.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoContext.java index 007cdeb7b..8cd27a2bb 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoContext.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoContext.java @@ -15,11 +15,15 @@ */ package org.springframework.data.mongodb.core; -import org.reactivestreams.Publisher; -import org.springframework.util.Assert; import reactor.core.publisher.Mono; import reactor.util.context.Context; +import java.util.function.Function; + +import org.reactivestreams.Publisher; + +import org.springframework.util.Assert; + import com.mongodb.reactivestreams.client.ClientSession; /** @@ -29,7 +33,7 @@ import com.mongodb.reactivestreams.client.ClientSession; * @author Christoph Strobl * @author Mark Paluch * @since 2.1 - * @see Mono#subscriberContext() + * @see Mono#deferContextual(Function) * @see Context */ public class ReactiveMongoContext { @@ -46,8 +50,14 @@ public class ReactiveMongoContext { */ public static Mono getSession() { - return Mono.subscriberContext().filter(ctx -> ctx.hasKey(SESSION_KEY)) - .flatMap(ctx -> ctx.> get(SESSION_KEY)); + return Mono.deferContextual(ctx -> { + + if (ctx.hasKey(SESSION_KEY)) { + return ctx.> get(SESSION_KEY); + } + + return Mono.empty(); + }); } /** diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java index da60cd671..7b35ea1f4 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java @@ -572,7 +572,7 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati ReactiveMongoTemplate.this); return Flux.from(action.doInSession(operations)) // - .subscriberContext(ctx -> ReactiveMongoContext.setSession(ctx, Mono.just(session))); + .contextWrite(ctx -> ReactiveMongoContext.setSession(ctx, Mono.just(session))); } /* diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java index 1e0e34261..46708e280 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java @@ -35,6 +35,7 @@ import com.mongodb.client.MongoClient; import com.mongodb.client.MongoCollection; import com.mongodb.client.MongoDatabase; import com.mongodb.reactivestreams.client.MongoClients; +import reactor.util.retry.Retry; /** * Utility to create (and reuse) imperative and reactive {@code MongoClient} instances. @@ -160,7 +161,7 @@ public class MongoTestUtils { .withWriteConcern(WriteConcern.MAJORITY).withReadPreference(ReadPreference.primary()); Mono.from(database.getCollection(collectionName).drop()) // - .delayElement(getTimeout()).retryBackoff(3, Duration.ofMillis(250)) // + .delayElement(getTimeout()).retryWhen(Retry.backoff(3, Duration.ofMillis(250))) // .as(StepVerifier::create) // .verifyComplete(); }