Fix deprecations from Reactor
This commit is contained in:
@@ -24,10 +24,11 @@ import org.springframework.util.Assert;
|
||||
|
||||
import reactor.core.Disposable;
|
||||
import reactor.core.Disposables;
|
||||
import reactor.core.publisher.EmitterProcessor;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.FluxIdentityProcessor;
|
||||
import reactor.core.publisher.FluxSink;
|
||||
import reactor.core.publisher.ReplayProcessor;
|
||||
import reactor.core.publisher.Processors;
|
||||
import reactor.core.publisher.Sinks;
|
||||
import reactor.core.scheduler.Schedulers;
|
||||
|
||||
/**
|
||||
@@ -43,16 +44,17 @@ import reactor.core.scheduler.Schedulers;
|
||||
public class FluxMessageChannel extends AbstractMessageChannel
|
||||
implements Publisher<Message<?>>, ReactiveStreamsSubscribableChannel {
|
||||
|
||||
private final EmitterProcessor<Message<?>> processor;
|
||||
private final FluxIdentityProcessor<Message<?>> processor;
|
||||
|
||||
private final FluxSink<Message<?>> sink;
|
||||
|
||||
private final ReplayProcessor<Boolean> subscribedSignal = ReplayProcessor.create(1);
|
||||
private final Sinks.StandaloneFluxSink<Boolean> subscribedSignal = Sinks.replay(1);
|
||||
|
||||
private final Disposable.Composite upstreamSubscriptions = Disposables.composite();
|
||||
|
||||
@SuppressWarnings("deprecation")
|
||||
public FluxMessageChannel() {
|
||||
this.processor = EmitterProcessor.create(1, false);
|
||||
this.processor = Processors.more().multicast(1, false);
|
||||
this.sink = this.processor.sink(FluxSink.OverflowStrategy.BUFFER);
|
||||
}
|
||||
|
||||
@@ -67,16 +69,16 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
@Override
|
||||
public void subscribe(Subscriber<? super Message<?>> subscriber) {
|
||||
this.processor
|
||||
.doFinally((s) -> this.subscribedSignal.onNext(this.processor.hasDownstreams()))
|
||||
.doFinally((s) -> this.subscribedSignal.next(this.processor.hasDownstreams()))
|
||||
.subscribe(subscriber);
|
||||
this.subscribedSignal.onNext(this.processor.hasDownstreams());
|
||||
this.subscribedSignal.next(this.processor.hasDownstreams());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void subscribeTo(Publisher<? extends Message<?>> publisher) {
|
||||
this.upstreamSubscriptions.add(
|
||||
Flux.from(publisher)
|
||||
.delaySubscription(this.subscribedSignal.filter(Boolean::booleanValue).next())
|
||||
.delaySubscription(this.subscribedSignal.asFlux().filter(Boolean::booleanValue).next())
|
||||
.publishOn(Schedulers.boundedElastic())
|
||||
.doOnNext((message) -> {
|
||||
try {
|
||||
@@ -91,7 +93,7 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
|
||||
@Override
|
||||
public void destroy() {
|
||||
this.subscribedSignal.onNext(false);
|
||||
this.subscribedSignal.next(false);
|
||||
this.upstreamSubscriptions.dispose();
|
||||
this.processor.onComplete();
|
||||
super.destroy();
|
||||
|
||||
@@ -30,9 +30,10 @@ import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.messaging.PollableChannel;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
|
||||
import reactor.core.publisher.EmitterProcessor;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.FluxIdentityProcessor;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.core.publisher.Processors;
|
||||
import reactor.core.scheduler.Schedulers;
|
||||
|
||||
/**
|
||||
@@ -101,7 +102,7 @@ public final class IntegrationReactiveUtils {
|
||||
* - a {@link org.springframework.integration.channel.FluxMessageChannel}
|
||||
* is returned as is because it is already a {@link Publisher};
|
||||
* - a {@link SubscribableChannel} is subscribed with a {@link MessageHandler}
|
||||
* for the {@link EmitterProcessor#onNext(Object)} which is returned from this method;
|
||||
* for the {@link FluxIdentityProcessor#onNext(Object)} which is returned from this method;
|
||||
* - a {@link PollableChannel} is wrapped into a {@link MessageSource} lambda and reuses
|
||||
* {@link #messageSourceToFlux(MessageSource)}.
|
||||
* @param messageChannel the {@link MessageChannel} to adapt.
|
||||
@@ -127,7 +128,7 @@ public final class IntegrationReactiveUtils {
|
||||
|
||||
private static <T> Flux<Message<T>> adaptSubscribableChannelToPublisher(SubscribableChannel inputChannel) {
|
||||
return Flux.defer(() -> {
|
||||
EmitterProcessor<Message<T>> publisher = EmitterProcessor.create(1);
|
||||
FluxIdentityProcessor<Message<T>> publisher = Processors.more().multicast(1);
|
||||
@SuppressWarnings("unchecked")
|
||||
MessageHandler messageHandler = (message) -> publisher.onNext((Message<T>) message);
|
||||
inputChannel.subscribe(messageHandler);
|
||||
|
||||
Reference in New Issue
Block a user