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
This commit is contained in:
Brian Clozel
2020-06-19 22:13:15 +02:00
parent 34cb4895c4
commit 7391f9b392
2 changed files with 6 additions and 5 deletions

View File

@@ -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<P> implements TcpOperations<P> {
logger.debug("Connected to " + conn.address());
}
});
DirectProcessor<Void> completion = DirectProcessor.create();
FluxIdentityProcessor<Void> completion = Processors.more().multicastNoBackpressure();
TcpConnection<P> connection = new ReactorNettyTcpConnection<>(inbound, outbound, codec, completion);
scheduler.schedule(() -> this.connectionHandler.afterConnected(connection));

View File

@@ -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<P> implements TcpConnection<P> {
private final ReactorNettyCodec<P> codec;
private final DirectProcessor<Void> closeProcessor;
private final FluxIdentityProcessor<Void> closeProcessor;
public ReactorNettyTcpConnection(NettyInbound inbound, NettyOutbound outbound,
ReactorNettyCodec<P> codec, DirectProcessor<Void> closeProcessor) {
ReactorNettyCodec<P> codec, FluxIdentityProcessor<Void> closeProcessor) {
this.inbound = inbound;
this.outbound = outbound;