From d462a81559e0f96e116a4dba34f3854ebe88e494 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Sat, 28 Aug 2021 09:51:38 +0200 Subject: [PATCH] GH-2218 Fix support for multiple consumers to the same destination This fix implies that a single error infristructure will be shared by such consumers Resolves #2218 --- .../binder/AbstractMessageChannelBinder.java | 25 ++++++++++++------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index a93715271..c6733a205 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -684,8 +684,11 @@ public abstract class AbstractMessageChannelBinder recoverer); + if (!getApplicationContext().containsBean(recovererBeanName)) { + ((GenericApplicationContext) getApplicationContext()).registerBean( + recovererBeanName, ErrorMessageSendingRecoverer.class, () -> recoverer); + } + MessageHandler handler; if (polled) { handler = getPolledConsumerErrorMessageHandler(destination, group, @@ -711,11 +714,13 @@ public abstract class AbstractMessageChannelBinder errorHandler); - errorChannel.subscribe(handler); + if (!getApplicationContext().containsBean(errorMessageHandlerName)) { + MessageHandler errorHandler = handler; + ((GenericApplicationContext) getApplicationContext()).registerBean( + errorMessageHandlerName, MessageHandler.class, + () -> errorHandler); + errorChannel.subscribe(handler); + } } else { this.logger.warn("The provided errorChannel '" + errorChannelName @@ -734,8 +739,10 @@ public abstract class AbstractMessageChannelBinder errorBridge); + if (getApplicationContext().containsBean(errorBridgeHandlerName)) { + ((GenericApplicationContext) getApplicationContext()).registerBean( + errorBridgeHandlerName, BridgeHandler.class, () -> errorBridge); + } } else { this.logger.warn("The provided errorChannel '" + errorChannelName