From a8c471c06174b0089df759b13ee6ec64bb6b8185 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 10 Feb 2020 16:27:59 -0500 Subject: [PATCH] Fix ReactiveStreamsConsumer for error handling To support `onErrorContinue()` logic for the plain `Subscriber` we need to wrap its `onNext()` into a `try..catch` and respective `errorHandler` in the `ReactiveStreamsConsumer` --- .../endpoint/ReactiveStreamsConsumer.java | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 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 c618889aee..5a791ec654 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 @@ -166,21 +166,25 @@ public class ReactiveStreamsConsumer extends AbstractEndpoint implements Integra this.subscription = Flux.from(this.publisher) .flatMap(this.reactiveMessageHandler::handleMessage) - .onErrorContinue(this::onError) + .onErrorContinue((ex, data) -> this.errorHandler.handleError(ex)) .subscribe(); } else if (this.subscriber != null) { - Flux.from(this.publisher) - .doOnSubscribe((subs) -> this.subscription = subs::cancel) - .onErrorContinue(this::onError) - .subscribe(this.subscriber); + this.subscription = + Flux.from(this.publisher) + .doOnComplete(this.subscriber::onComplete) + .doOnSubscribe(this.subscriber::onSubscribe) + .subscribe((data) -> { + try { + this.subscriber.onNext(data); + } + catch (Exception ex) { + this.errorHandler.handleError(ex); + } + }); } } - private void onError(Throwable ex, Object data) { - this.errorHandler.handleError(ex); - } - @Override protected void doStop() { if (this.subscription != null) {