diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java index 9e699e6946..e3d27f804a 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java @@ -171,7 +171,7 @@ public class FailoverClientConnectionFactory extends AbstractClientConnectionFac return failoverTcpConnection; } - private void closeRefreshedIfNecessary(@Nullable FailoverTcpConnection sharedConnection, boolean refreshShared, + private void closeRefreshedIfNecessary(FailoverTcpConnection sharedConnection, boolean refreshShared, FailoverTcpConnection failoverTcpConnection) { this.creationTime = System.currentTimeMillis(); @@ -179,7 +179,6 @@ public class FailoverClientConnectionFactory extends AbstractClientConnectionFac * We may have simply wrapped the same connection in a new wrapper; don't close. */ if (refreshShared && this.closeOnRefresh - && sharedConnection != null && !sharedConnection.delegate.equals(failoverTcpConnection.delegate) && sharedConnection.isOpen()) { diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java index 7b86d1bddd..a229966f8c 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java @@ -128,7 +128,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv * @param completionTimeout The timeout. * @since 4.1 */ - public void setCompletionTimeout(long completionTimeout) { + public synchronized void setCompletionTimeout(long completionTimeout) { this.completionTimeout = completionTimeout; } @@ -138,7 +138,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv * @param completionTimeout The timeout. * @since 5.1.10 */ - public void setDisconnectCompletionTimeout(long completionTimeout) { + public synchronized void setDisconnectCompletionTimeout(long completionTimeout) { this.disconnectCompletionTimeout = completionTimeout; }