From 447fa920d971ef910b77f4b3d2ee39f7e7d59c60 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 29 Mar 2022 12:41:45 -0400 Subject: [PATCH] GH-1443: Pull CCF.resetConnection() to CF Resolves https://github.com/spring-projects/spring-amqp/issues/1443 **Cherry-pick to `2.4.x`** --- .../AbstractRoutingConnectionFactory.java | 14 +++++++++++++- .../connection/CachingConnectionFactory.java | 3 ++- .../amqp/rabbit/connection/ConnectionFactory.java | 10 +++++++++- .../LocalizedQueueConnectionFactory.java | 9 +++++++-- .../connection/PooledChannelConnectionFactory.java | 1 + .../connection/ThreadChannelConnectionFactory.java | 3 ++- 6 files changed, 34 insertions(+), 6 deletions(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java index 5cfa2adc..2ea39028 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java @@ -22,6 +22,7 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import org.springframework.amqp.AmqpException; +import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -38,7 +39,7 @@ import org.springframework.util.Assert; * @since 1.3 */ public abstract class AbstractRoutingConnectionFactory implements ConnectionFactory, RoutingConnectionFactory, - InitializingBean { + InitializingBean, DisposableBean { private final Map targetConnectionFactories = new ConcurrentHashMap(); @@ -260,4 +261,15 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact @Nullable protected abstract Object determineCurrentLookupKey(); + @Override + public void destroy() { + resetConnection(); + } + + @Override + public void resetConnection() { + this.targetConnectionFactories.values().forEach(factory -> factory.resetConnection()); + this.defaultTargetConnectionFactory.resetConnection(); + } + } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java index 0fa96068..15504f22 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -882,6 +882,7 @@ public class CachingConnectionFactory extends AbstractConnectionFactory * used to force a reconnect to the primary broker after failing over to a secondary * broker. */ + @Override public void resetConnection() { synchronized (this.connectionMonitor) { if (this.connection.target != null) { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactory.java index 4b1e049b..2f4323b5 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2019 the original author or authors. + * Copyright 2002-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -83,4 +83,12 @@ public interface ConnectionFactory { return false; } + /** + * Close any connection(s) that might be cached by this factory. This does not prevent + * new connections from being opened. + * @since 2.4.4 + */ + default void resetConnection() { + } + } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/LocalizedQueueConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/LocalizedQueueConnectionFactory.java index fc04a731..e209734e 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/LocalizedQueueConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/LocalizedQueueConnectionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2015-2019 the original author or authors. + * Copyright 2015-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -350,7 +350,7 @@ public class LocalizedQueueConnectionFactory implements ConnectionFactory, Routi } @Override - public void destroy() { + public void resetConnection() { Exception lastException = null; for (ConnectionFactory connectionFactory : this.nodeFactories.values()) { if (connectionFactory instanceof DisposableBean) { @@ -367,4 +367,9 @@ public class LocalizedQueueConnectionFactory implements ConnectionFactory, Routi } } + @Override + public void destroy() { + resetConnection(); + } + } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java index 391874f3..15add409 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java @@ -149,6 +149,7 @@ public class PooledChannelConnectionFactory extends AbstractConnectionFactory im * used to force a reconnect to the primary broker after failing over to a secondary * broker. */ + @Override public void resetConnection() { destroy(); } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java index 35a28687..0da501fe 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2021 the original author or authors. + * Copyright 2020-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -140,6 +140,7 @@ public class ThreadChannelConnectionFactory extends AbstractConnectionFactory im * used to force a reconnect to the primary broker after failing over to a secondary * broker. */ + @Override public void resetConnection() { destroy(); }