GH-3054 Clear channel cache in StreamBridge during the refresh

This commit is contained in:
Oleg Zhurakousky
2024-12-19 12:30:50 +01:00
parent c8dcbf1871
commit 24c1abea50

View File

@@ -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<ApplicationEvent> {
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());