From e649480350ba4ccd65d17f38d767c3b680b26c43 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 23 Sep 2020 15:33:17 -0400 Subject: [PATCH] Fix `IntegrationReactiveUtils` The `Mono.doOnSuccess()` is always called for completed `MonoSink` even if the value is `null` * Fix `IntegrationReactiveUtils.messageSourceToFlux()` to check a message for `null` before calling `AckUtils.autoAck()` * Add an `logger.error()` message when `doOnError()` is invoked for that `Mono` **Cherry-pick to 5.3.x** --- .../util/IntegrationReactiveUtils.java | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/IntegrationReactiveUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/util/IntegrationReactiveUtils.java index 73ece318bb..3bebf776e6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/IntegrationReactiveUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/IntegrationReactiveUtils.java @@ -19,6 +19,8 @@ package org.springframework.integration.util; import java.time.Duration; import java.util.function.Function; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.reactivestreams.Publisher; import org.springframework.integration.StaticMessageHeaderAccessor; @@ -45,6 +47,8 @@ import reactor.core.scheduler.Schedulers; */ public final class IntegrationReactiveUtils { + private static final Log logger = LogFactory.getLog(IntegrationReactiveUtils.class); + /** * The subscriber context entry for {@link Flux#delayElements} * from the {@link Mono#repeatWhenEmpty(Function)}. @@ -75,16 +79,19 @@ public final class IntegrationReactiveUtils { public static Flux> messageSourceToFlux(MessageSource messageSource) { return Mono. >create(monoSink -> - monoSink.onRequest(value -> - monoSink.success(messageSource.receive()))) - .doOnSuccess((message) -> - AckUtils.autoAck(StaticMessageHeaderAccessor.getAcknowledgmentCallback(message))) + monoSink.onRequest(value -> monoSink.success(messageSource.receive()))) + .doOnSuccess((message) -> { + if (message != null) { + AckUtils.autoAck(StaticMessageHeaderAccessor.getAcknowledgmentCallback(message)); + } + }) .doOnError(MessagingException.class, (ex) -> { Message failedMessage = ex.getFailedMessage(); if (failedMessage != null) { AckUtils.autoNack(StaticMessageHeaderAccessor.getAcknowledgmentCallback(failedMessage)); } + logger.error("Error from Flux for : " + messageSource, ex); }) .subscribeOn(Schedulers.boundedElastic()) .repeatWhenEmpty((repeat) ->