From c12acb9ff86f5d87dc6881aa780ef5ca7050deac Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 25 Nov 2020 15:22:46 -0500 Subject: [PATCH] Fix FluxMessCh for subscription race condition Related to https://build.spring.io/browse/INT-MASTER-2240 Turns out the `doOnRequest()` is still not enough to be sure that subscriber is accepted into the `Publisher`. Probably because `doOnRequest()` maybe really asked for the subscription by itself for some logic * Use `Mono` with `repeatWhenEmpty()` until `this.sink.currentSubscriberCount() > 0` --- .../integration/channel/FluxMessageChannel.java | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) 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 3f222e50ce..0d09425c47 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 @@ -16,6 +16,7 @@ package org.springframework.integration.channel; +import java.time.Duration; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.LockSupport; @@ -29,6 +30,7 @@ import org.springframework.util.Assert; import reactor.core.Disposable; import reactor.core.Disposables; import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; import reactor.core.scheduler.Schedulers; @@ -94,10 +96,17 @@ public class FluxMessageChannel extends AbstractMessageChannel @Override public void subscribe(Subscriber> subscriber) { this.sink.asFlux() - .doOnRequest((r) -> this.subscribedSignal.tryEmitNext(true)) .doFinally((s) -> this.subscribedSignal.tryEmitNext(this.sink.currentSubscriberCount() > 0)) .share() .subscribe(subscriber); + + this.upstreamSubscriptions.add( + Mono.fromCallable(() -> this.sink.currentSubscriberCount() > 0) + .filter(Boolean::booleanValue) + .doOnNext(this.subscribedSignal::tryEmitNext) + .repeatWhenEmpty((repeat) -> + this.active ? repeat.delayElements(Duration.ofMillis(100)) : repeat) + .subscribe()); } @Override