From 7391f9b3923dc01b5c92194dd40f39a1f9eaf779 Mon Sep 17 00:00:00 2001 From: Brian Clozel Date: Fri, 19 Jun 2020 22:13:15 +0200 Subject: [PATCH] Fix Reactor Core DirectProcessor deprecation As of reactor/reactor-core#2188, `DirectProcessor` variants are deprecated. This commit replaces them with the new `FluxIdentityProcessor` variant. See gh-25085 --- .../messaging/tcp/reactor/ReactorNettyTcpClient.java | 5 +++-- .../messaging/tcp/reactor/ReactorNettyTcpConnection.java | 6 +++--- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpClient.java b/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpClient.java index f755ad117d..82d31bc78a 100644 --- a/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpClient.java +++ b/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpClient.java @@ -32,9 +32,10 @@ import io.netty.util.concurrent.ImmediateEventExecutor; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.reactivestreams.Publisher; -import reactor.core.publisher.DirectProcessor; +import reactor.core.publisher.FluxIdentityProcessor; import reactor.core.publisher.Mono; import reactor.core.publisher.MonoProcessor; +import reactor.core.publisher.Processors; import reactor.core.scheduler.Scheduler; import reactor.core.scheduler.Schedulers; import reactor.netty.Connection; @@ -316,7 +317,7 @@ public class ReactorNettyTcpClient

implements TcpOperations

{ logger.debug("Connected to " + conn.address()); } }); - DirectProcessor completion = DirectProcessor.create(); + FluxIdentityProcessor completion = Processors.more().multicastNoBackpressure(); TcpConnection

connection = new ReactorNettyTcpConnection<>(inbound, outbound, codec, completion); scheduler.schedule(() -> this.connectionHandler.afterConnected(connection)); diff --git a/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpConnection.java b/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpConnection.java index 9665057011..54611adcfc 100644 --- a/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpConnection.java +++ b/spring-messaging/src/main/java/org/springframework/messaging/tcp/reactor/ReactorNettyTcpConnection.java @@ -17,7 +17,7 @@ package org.springframework.messaging.tcp.reactor; import io.netty.buffer.ByteBuf; -import reactor.core.publisher.DirectProcessor; +import reactor.core.publisher.FluxIdentityProcessor; import reactor.core.publisher.Mono; import reactor.netty.NettyInbound; import reactor.netty.NettyOutbound; @@ -42,11 +42,11 @@ public class ReactorNettyTcpConnection

implements TcpConnection

{ private final ReactorNettyCodec

codec; - private final DirectProcessor closeProcessor; + private final FluxIdentityProcessor closeProcessor; public ReactorNettyTcpConnection(NettyInbound inbound, NettyOutbound outbound, - ReactorNettyCodec

codec, DirectProcessor closeProcessor) { + ReactorNettyCodec

codec, FluxIdentityProcessor closeProcessor) { this.inbound = inbound; this.outbound = outbound;