From d2e5a0dfdca54cd1345783762641f4f895f2fb5c Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 18 Nov 2024 14:49:30 -0500 Subject: [PATCH] GH-2907: Use `CF.closeTimeout` for confirms wait Fixes: #2907 Issue link: https://github.com/spring-projects/spring-amqp/issues/2907 The current hard-coded `5 seconds` is not enough in real applications under heavy load * Fix `CachingConnectionFactory` to use `getCloseTimeout()` for `publisherCallbackChannel.waitForConfirms()` which is `30 seconds` by default, but can be modified via `CachingConnectionFactory.setCloseTimeout()` (cherry picked from commit 562bc772c436e6b09527084f238b46588ba8ba6d) --- .../amqp/rabbit/connection/AbstractConnectionFactory.java | 5 +++-- .../amqp/rabbit/connection/CachingConnectionFactory.java | 8 +++----- 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractConnectionFactory.java index 33cb11dc..fc3e6a09 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractConnectionFactory.java @@ -478,8 +478,9 @@ public abstract class AbstractConnectionFactory implements ConnectionFactory, Di } /** - * How long to wait (milliseconds) for a response to a connection close operation from the broker; default 30000 (30 - * seconds). + * How long to wait (milliseconds) for a response to a connection close operation from the broker; + * default 30000 (30 seconds). + * Also used for {@link com.rabbitmq.client.Channel#waitForConfirms()}. * @param closeTimeout the closeTimeout to set. */ public void setCloseTimeout(int closeTimeout) { 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 f14aa989..0aade483 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 @@ -1087,8 +1087,6 @@ public class CachingConnectionFactory extends AbstractConnectionFactory private final class CachedChannelInvocationHandler implements InvocationHandler { - private static final int ASYNC_CLOSE_TIMEOUT = 5_000; - private final ChannelCachingConnectionProxy theConnection; private final Deque channelList; @@ -1302,7 +1300,7 @@ public class CachingConnectionFactory extends AbstractConnectionFactory getChannelsExecutor() .execute(() -> { try { - publisherCallbackChannel.waitForConfirms(ASYNC_CLOSE_TIMEOUT); + publisherCallbackChannel.waitForConfirms(getCloseTimeout()); } catch (InterruptedException ex) { Thread.currentThread().interrupt(); @@ -1426,10 +1424,10 @@ public class CachingConnectionFactory extends AbstractConnectionFactory executorService.execute(() -> { try { if (ConfirmType.CORRELATED.equals(CachingConnectionFactory.this.confirmType)) { - channel.waitForConfirmsOrDie(ASYNC_CLOSE_TIMEOUT); + channel.waitForConfirmsOrDie(getCloseTimeout()); } else { - Thread.sleep(ASYNC_CLOSE_TIMEOUT); + Thread.sleep(5_000); // NOSONAR - some time to give the channel a chance to ack } } catch (@SuppressWarnings(UNUSED) InterruptedException e1) {