GH-9215: Honor back-pressure in FluxMessageChannel
Fixes: #9215
* Instead of `share()` use `publish(1).refCount()` to prefetch only item from upstream.
* Also remove `publishOn(this.scheduler)` for upstream publishers in favor of opt-in on the consumer side.
(cherry picked from commit e9561b41f1)
This commit is contained in:
committed by
Spring Builds
parent
d22737ccb0
commit
bf6c4f6157
@@ -28,8 +28,6 @@ import reactor.core.Disposables;
|
|||||||
import reactor.core.publisher.Flux;
|
import reactor.core.publisher.Flux;
|
||||||
import reactor.core.publisher.Mono;
|
import reactor.core.publisher.Mono;
|
||||||
import reactor.core.publisher.Sinks;
|
import reactor.core.publisher.Sinks;
|
||||||
import reactor.core.scheduler.Scheduler;
|
|
||||||
import reactor.core.scheduler.Schedulers;
|
|
||||||
import reactor.util.context.ContextView;
|
import reactor.util.context.ContextView;
|
||||||
|
|
||||||
import org.springframework.core.log.LogMessage;
|
import org.springframework.core.log.LogMessage;
|
||||||
@@ -54,8 +52,6 @@ import org.springframework.util.Assert;
|
|||||||
public class FluxMessageChannel extends AbstractMessageChannel
|
public class FluxMessageChannel extends AbstractMessageChannel
|
||||||
implements Publisher<Message<?>>, ReactiveStreamsSubscribableChannel {
|
implements Publisher<Message<?>>, ReactiveStreamsSubscribableChannel {
|
||||||
|
|
||||||
private final Scheduler scheduler = Schedulers.boundedElastic();
|
|
||||||
|
|
||||||
private final Sinks.Many<Message<?>> sink = Sinks.many().multicast().onBackpressureBuffer(1, false);
|
private final Sinks.Many<Message<?>> sink = Sinks.many().multicast().onBackpressureBuffer(1, false);
|
||||||
|
|
||||||
private final Sinks.Many<Boolean> subscribedSignal = Sinks.many().replay().limit(1);
|
private final Sinks.Many<Boolean> subscribedSignal = Sinks.many().replay().limit(1);
|
||||||
@@ -107,7 +103,8 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
|||||||
public void subscribe(Subscriber<? super Message<?>> subscriber) {
|
public void subscribe(Subscriber<? super Message<?>> subscriber) {
|
||||||
this.sink.asFlux()
|
this.sink.asFlux()
|
||||||
.doFinally((s) -> this.subscribedSignal.tryEmitNext(this.sink.currentSubscriberCount() > 0))
|
.doFinally((s) -> this.subscribedSignal.tryEmitNext(this.sink.currentSubscriberCount() > 0))
|
||||||
.share()
|
.publish(1)
|
||||||
|
.refCount()
|
||||||
.subscribe(subscriber);
|
.subscribe(subscriber);
|
||||||
|
|
||||||
Mono<Boolean> subscribersBarrier =
|
Mono<Boolean> subscribersBarrier =
|
||||||
@@ -148,7 +145,7 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
|||||||
Flux<Object> upstreamPublisher =
|
Flux<Object> upstreamPublisher =
|
||||||
Flux.from(publisher)
|
Flux.from(publisher)
|
||||||
.delaySubscription(this.subscribedSignal.asFlux().filter(Boolean::booleanValue).next())
|
.delaySubscription(this.subscribedSignal.asFlux().filter(Boolean::booleanValue).next())
|
||||||
.publishOn(this.scheduler)
|
// .publishOn(this.scheduler)
|
||||||
.flatMap((message) ->
|
.flatMap((message) ->
|
||||||
Mono.just(message)
|
Mono.just(message)
|
||||||
.handle((messageToHandle, syncSink) -> sendReactiveMessage(messageToHandle))
|
.handle((messageToHandle, syncSink) -> sendReactiveMessage(messageToHandle))
|
||||||
@@ -185,7 +182,6 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
|||||||
this.upstreamSubscriptions.dispose();
|
this.upstreamSubscriptions.dispose();
|
||||||
this.subscribedSignal.emitComplete(Sinks.EmitFailureHandler.FAIL_FAST);
|
this.subscribedSignal.emitComplete(Sinks.EmitFailureHandler.FAIL_FAST);
|
||||||
this.sink.emitComplete(Sinks.EmitFailureHandler.FAIL_FAST);
|
this.sink.emitComplete(Sinks.EmitFailureHandler.FAIL_FAST);
|
||||||
this.scheduler.dispose();
|
|
||||||
super.destroy();
|
super.destroy();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user