diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java index d071ebc9d5..7be166ca12 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java @@ -54,8 +54,6 @@ public class FluxMessageChannel extends AbstractMessageChannel private final Sinks.Many> sink = Sinks.many().multicast().onBackpressureBuffer(1, false); - private final Sinks.Many subscribedSignal = Sinks.many().replay().limit(1); - private final Disposable.Composite upstreamSubscriptions = Disposables.composite(); private volatile boolean active = true; @@ -102,19 +100,9 @@ public class FluxMessageChannel extends AbstractMessageChannel @Override public void subscribe(Subscriber> subscriber) { this.sink.asFlux() - .doFinally((s) -> this.subscribedSignal.tryEmitNext(this.sink.currentSubscriberCount() > 0)) .publish(1) .refCount() .subscribe(subscriber); - - Mono subscribersBarrier = - Mono.fromCallable(() -> this.sink.currentSubscriberCount() > 0) - .filter(Boolean::booleanValue) - .doOnNext(this.subscribedSignal::tryEmitNext) - .repeatWhenEmpty((repeat) -> - this.active ? repeat.delayElements(Duration.ofMillis(100)) : repeat); // NOSONAR - - addPublisherToSubscribe(Flux.from(subscribersBarrier)); } private void addPublisherToSubscribe(Flux publisher) { @@ -144,8 +132,11 @@ public class FluxMessageChannel extends AbstractMessageChannel public void subscribeTo(Publisher> publisher) { Flux upstreamPublisher = Flux.from(publisher) - .delaySubscription(this.subscribedSignal.asFlux().filter(Boolean::booleanValue).next()) -// .publishOn(this.scheduler) + .delaySubscription( + Mono.fromCallable(this.sink::currentSubscriberCount) + .filter((value) -> value > 0) + .repeatWhenEmpty((repeat) -> + this.active ? repeat.delayElements(Duration.ofMillis(100)) : repeat)) .flatMap((message) -> Mono.just(message) .handle((messageToHandle, syncSink) -> sendReactiveMessage(messageToHandle)) @@ -180,7 +171,6 @@ public class FluxMessageChannel extends AbstractMessageChannel public void destroy() { this.active = false; this.upstreamSubscriptions.dispose(); - this.subscribedSignal.emitComplete(Sinks.EmitFailureHandler.FAIL_FAST); this.sink.emitComplete(Sinks.EmitFailureHandler.FAIL_FAST); super.destroy(); }