From 24c1abea50cb7016a138d0b33382c8f20f6ad954 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 19 Dec 2024 12:30:50 +0100 Subject: [PATCH] GH-3054 Clear channel cache in StreamBridge during the refresh --- .../cloud/stream/function/StreamBridge.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index e734cec5d..2871b833e 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -50,6 +50,8 @@ import org.springframework.cloud.stream.binding.NewDestinationBindingCallback; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; +import org.springframework.context.ApplicationEvent; +import org.springframework.context.ApplicationListener; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.core.ResolvableType; import org.springframework.integration.channel.AbstractMessageChannel; @@ -89,7 +91,7 @@ import org.springframework.util.StringUtils; * */ @SuppressWarnings("rawtypes") -public final class StreamBridge implements StreamOperations, SmartInitializingSingleton, DisposableBean { +public final class StreamBridge implements StreamOperations, SmartInitializingSingleton, DisposableBean, ApplicationListener { private static final String STREAM_BRIDGE_FUNC_NAME = "streamBridge"; @@ -360,6 +362,15 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi this.async = async; } + // see https://github.com/spring-cloud/spring-cloud-stream/issues/3054 + @Override + public void onApplicationEvent(ApplicationEvent event) { + // we need to do it by String to avoid cloud-bus and context dependencies + if (event.getClass().getName().equals("org.springframework.cloud.bus.event.RefreshRemoteApplicationEvent")) { + this.channelCache.clear(); + } + } + private static final class ContextPropagationHelper { static ExecutorService wrap(ExecutorService executorService) { return ContextExecutorService.wrap(executorService, () -> ContextSnapshotFactory.builder().build().captureAll());