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) {