From 5aaacb97b7940a044a223bf7976f57bd695a2e6b Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 22 Mar 2023 17:37:26 +0100 Subject: [PATCH] GH-2672 re-add bridge to global error channel to ensure logging Resolves #2672 --- .../binder/AbstractMessageChannelBinder.java | 47 ++++++++++--------- 1 file changed, 24 insertions(+), 23 deletions(-) diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index 83ae990fc..3c393ca56 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -41,7 +41,6 @@ import org.springframework.cloud.stream.config.ConsumerEndpointCustomizer; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.config.MessageSourceCustomizer; import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; -import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.cloud.stream.provisioning.ProvisioningException; @@ -705,9 +704,9 @@ public abstract class AbstractMessageChannelBinder errorHandler = catalog.lookup(Consumer.class, bp.getErrorHandlerDefinition()); if (errorHandler == null) { logger.warn("Failed to retrieve error handling function with definition: " + bp.getErrorHandlerDefinition() + ", for binding: " + bindingName); + return false; } else { SubscribableChannel functionErrorChannel = getApplicationContext().getBean(errorChannelName, SubscribableChannel.class); functionErrorChannel.subscribe(errorMessage -> errorHandler.accept((ErrorMessage) errorMessage)); + return true; } } } + return false; } /** @@ -758,22 +760,17 @@ public abstract class AbstractMessageChannelBinder binderErrorChannel); } - this.subscribeFunctionErrorHandler(errorChannelName, consumerProperties.getBindingName()); + + boolean userHandlerSubscribed = this.subscribeFunctionErrorHandler(errorChannelName, consumerProperties.getBindingName()); ErrorMessageSendingRecoverer recoverer = new ErrorMessageSendingRecoverer(binderErrorChannel, errorMessageStrategy); String recovererBeanName = getErrorRecovererName(destination, group, consumerProperties); @@ -792,16 +789,12 @@ public abstract class AbstractMessageChannelBinder h); - binderErrorChannel.subscribe(binderProvidedErrorHandler); - } - else { - binderErrorChannel.subscribe((MessageHandler) getApplicationContext().getBean(errorMessageHandlerName)); - } + + if (binderProvidedErrorHandler != null && !userHandlerSubscribed) { + if (this.isSubscribable(binderErrorChannel) && !getApplicationContext().containsBean(errorMessageHandlerName)) { + MessageHandler h = binderProvidedErrorHandler; + ((GenericApplicationContext) getApplicationContext()).registerBean(errorMessageHandlerName, MessageHandler.class, () -> h); + binderErrorChannel.subscribe(binderProvidedErrorHandler); } else { this.logger.warn("The provided errorChannel '" + errorChannelName @@ -811,6 +804,14 @@ public abstract class AbstractMessageChannelBinder