From 6591ce90f67bdf66e91426e59197cc5c59c3d887 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 10 Feb 2020 16:00:14 -0500 Subject: [PATCH] Fix ReactiveStreamsConsumer for plain subscriber https://build.spring.io/browse/INT-MASTERSPRING40-978 Investigate a behaviour for `ReactiveStreamsConsumer` when we use a plain `Subscriber` directly instead of `doOn...` callbacks. It looks like there is some race condition when the data can be consumed upstream, but there is no consumer downstream ready yet --- .../endpoint/ReactiveStreamsConsumer.java | 25 ++++++++++--------- 1 file changed, 13 insertions(+), 12 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java index 8821462587..c618889aee 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java @@ -162,22 +162,23 @@ public class ReactiveStreamsConsumer extends AbstractEndpoint implements Integra this.lifecycleDelegate.start(); } - Flux flux = null; if (this.reactiveMessageHandler != null) { - flux = Flux.from(this.publisher) - .flatMap(this.reactiveMessageHandler::handleMessage); - } - else if (this.subscriber != null) { - flux = Flux.from(this.publisher) - .doOnSubscribe(this.subscriber::onSubscribe) - .doOnComplete(this.subscriber::onComplete) - .doOnNext(this.subscriber::onNext); - } - if (flux != null) { this.subscription = - flux.onErrorContinue((ex, data) -> this.errorHandler.handleError(ex)) + Flux.from(this.publisher) + .flatMap(this.reactiveMessageHandler::handleMessage) + .onErrorContinue(this::onError) .subscribe(); } + else if (this.subscriber != null) { + Flux.from(this.publisher) + .doOnSubscribe((subs) -> this.subscription = subs::cancel) + .onErrorContinue(this::onError) + .subscribe(this.subscriber); + } + } + + private void onError(Throwable ex, Object data) { + this.errorHandler.handleError(ex); } @Override