From 9b7b0d91aca72b99643b163afa1289b310d869a0 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 22 Jan 2020 13:15:56 -0500 Subject: [PATCH] FluxMessageChannel: try.catch not onErrorContinue It turns out that some upstream fluxes may come with the `onErrorResume()` logic. The `onErrorContinue()` here downstream in the `FluxMessageChannel.java` eliminates an `onErrorResume()` purpose. * Change the logic to `try..catch` in the `doOnNext()` instead and let that upstream `onErrorResume()` to do its job --- .../integration/channel/FluxMessageChannel.java | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) 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 ecf37bab9a..7d5593beab 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 @@ -78,8 +78,14 @@ public class FluxMessageChannel extends AbstractMessageChannel Flux.from(publisher) .delaySubscription(this.subscribedSignal.filter(Boolean::booleanValue).next()) .publishOn(Schedulers.boundedElastic()) - .doOnNext(this::send) - .onErrorContinue((ex, message) -> logger.warn("Error during processing event: " + message, ex)) + .doOnNext((message) -> { + try { + send(message); + } + catch (Exception ex) { + logger.warn("Error during processing event: " + message, ex); + } + }) .subscribe()); }