diff --git a/spring-web/src/testFixtures/java/org/springframework/web/testfixture/http/server/reactive/MockServerHttpResponse.java b/spring-web/src/testFixtures/java/org/springframework/web/testfixture/http/server/reactive/MockServerHttpResponse.java index 14ae615f10..6b69119ea3 100644 --- a/spring-web/src/testFixtures/java/org/springframework/web/testfixture/http/server/reactive/MockServerHttpResponse.java +++ b/spring-web/src/testFixtures/java/org/springframework/web/testfixture/http/server/reactive/MockServerHttpResponse.java @@ -26,6 +26,7 @@ import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.MonoProcessor; +import reactor.core.publisher.Sinks; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferFactory; @@ -64,7 +65,7 @@ public class MockServerHttpResponse extends AbstractServerHttpResponse { super(dataBufferFactory); this.writeHandler = body -> { // Avoid .then() which causes data buffers to be released - MonoProcessor completion = MonoProcessor.create(); + MonoProcessor completion = MonoProcessor.fromSink(Sinks.one()); this.body = body.doOnComplete(completion::onComplete).doOnError(completion::onError).cache(); this.body.subscribe(); return completion;