From af2a22b767561e436f5a700f7b51b18fda39ac9b Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 4 Aug 2025 15:34:36 +0000 Subject: [PATCH 1/4] fix: Use a single-thread executor for Ack and Modack callbacks when exactly-once is enabled --- .../v1/StreamingSubscriberConnection.java | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java index 26411942f..ca5396e16 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java @@ -61,6 +61,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -455,7 +456,13 @@ private void sendAckOperations( .setSubscription(subscription) .addAllAckIds(ackIdsInRequest) .build()); - ApiFutures.addCallback(ackFuture, callback, directExecutor()); + if (getExactlyOnceDeliveryEnabled()) { + // When exatly-once delivery is enabled, a lock is acquired on this callback which can cause deadlock when + // using MoreExecutors.directExecutor(). We use a new, single-thread executor to avoid this. + ApiFutures.addCallback(ackFuture, callback, Executors.newSingleThreadExecutor()); + } else { + ApiFutures.addCallback(ackFuture, callback, directExecutor()); + } pendingOperations++; } ackOperationsWaiter.incrementPendingCount(pendingOperations); @@ -504,7 +511,13 @@ private void sendModackOperations( .addAllAckIds(ackIdsInRequest) .setAckDeadlineSeconds(modackRequestData.getDeadlineExtensionSeconds()) .build()); - ApiFutures.addCallback(modackFuture, callback, directExecutor()); + if (getExactlyOnceDeliveryEnabled()) { + // When exatly-once delivery is enabled, a lock is acquired on this callback which can cause deadlock when + // using MoreExecutors.directExecutor(). We use a new, single-thread executor to avoid this. + ApiFutures.addCallback(modackFuture, callback, Executors.newSingleThreadExecutor()); + } else { + ApiFutures.addCallback(modackFuture, callback, directExecutor()); + } pendingOperations++; } } From 6fe29d3d30a8e8870dfbe2ea10d7310fd4168233 Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 4 Aug 2025 17:14:33 +0000 Subject: [PATCH 2/4] fix: Use new executor for EOD ack and modack callbacks --- .../v1/StreamingSubscriberConnection.java | 24 +++++++++++++------ .../google/cloud/pubsub/v1/Subscriber.java | 6 +++++ .../v1/StreamingSubscriberConnectionTest.java | 1 + 3 files changed, 24 insertions(+), 7 deletions(-) diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java index ca5396e16..4c5bb4ff9 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnection.java @@ -61,7 +61,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; +import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -98,6 +98,7 @@ final class StreamingSubscriberConnection extends AbstractApiService implements private final String subscription; private final SubscriptionName subscriptionNameObject; private final ScheduledExecutorService systemExecutor; + private final ExecutorService eodAckCallbackExecutor; private final MessageDispatcher messageDispatcher; private final FlowControlSettings flowControlSettings; @@ -129,6 +130,7 @@ private StreamingSubscriberConnection(Builder builder) { subscription = builder.subscription; subscriptionNameObject = SubscriptionName.parse(builder.subscription); systemExecutor = builder.systemExecutor; + eodAckCallbackExecutor = builder.eodAckCallbackExecutor; // We need to set the default stream ack deadline on the initial request, this will be // updated by modack requests in the message dispatcher @@ -457,9 +459,10 @@ private void sendAckOperations( .addAllAckIds(ackIdsInRequest) .build()); if (getExactlyOnceDeliveryEnabled()) { - // When exatly-once delivery is enabled, a lock is acquired on this callback which can cause deadlock when - // using MoreExecutors.directExecutor(). We use a new, single-thread executor to avoid this. - ApiFutures.addCallback(ackFuture, callback, Executors.newSingleThreadExecutor()); + // When exatly-once delivery is enabled, a lock is acquired on this callback which can cause + // deadlock when using MoreExecutors.directExecutor(). We use a separate executor to avoid + // this. + ApiFutures.addCallback(ackFuture, callback, eodAckCallbackExecutor); } else { ApiFutures.addCallback(ackFuture, callback, directExecutor()); } @@ -512,9 +515,10 @@ private void sendModackOperations( .setAckDeadlineSeconds(modackRequestData.getDeadlineExtensionSeconds()) .build()); if (getExactlyOnceDeliveryEnabled()) { - // When exatly-once delivery is enabled, a lock is acquired on this callback which can cause deadlock when - // using MoreExecutors.directExecutor(). We use a new, single-thread executor to avoid this. - ApiFutures.addCallback(modackFuture, callback, Executors.newSingleThreadExecutor()); + // When exatly-once delivery is enabled, a lock is acquired on this callback which can + // cause deadlock when using MoreExecutors.directExecutor(). We use a separate executor to + // avoid this. + ApiFutures.addCallback(modackFuture, callback, eodAckCallbackExecutor); } else { ApiFutures.addCallback(modackFuture, callback, directExecutor()); } @@ -749,6 +753,7 @@ public static final class Builder { private boolean useLegacyFlowControl; private ScheduledExecutorService executor; private ScheduledExecutorService systemExecutor; + private ExecutorService eodAckCallbackExecutor; private ApiClock clock; private boolean enableOpenTelemetryTracing; @@ -839,6 +844,11 @@ public Builder setSystemExecutor(ScheduledExecutorService systemExecutor) { return this; } + public Builder setEodAckCallbackExecutor(ExecutorService eodAckCallbackExecutor) { + this.eodAckCallbackExecutor = eodAckCallbackExecutor; + return this; + } + public Builder setClock(ApiClock clock) { this.clock = clock; return this; diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java index 36ec5dc36..ee1075fc6 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java @@ -52,6 +52,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.Executors; +import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadFactory; import java.util.concurrent.atomic.AtomicInteger; @@ -150,6 +151,7 @@ public class Subscriber extends AbstractApiService implements SubscriberInterfac // An instantiation of the SystemExecutorProvider used for processing acks // and other system actions. @Nullable private final ScheduledExecutorService alarmsExecutor; + private final ExecutorService eodAckCallbackExecutor; private final Distribution ackLatencyDistribution = new Distribution(Math.toIntExact(MAX_STREAM_ACK_DEADLINE.getSeconds()) + 1); @@ -200,6 +202,9 @@ private Subscriber(Builder builder) { backgroundResources.add(new ExecutorAsBackgroundResource((alarmsExecutor))); } + ExecutorProvider eodAckCallbackExecutorProvider = InstantiatingExecutorProvider.newBuilder().setExecutorThreadCount(2).build(); + eodAckCallbackExecutor = eodAckCallbackExecutorProvider.getExecutor(); + TransportChannelProvider channelProvider = builder.channelProvider; if (channelProvider.acceptsPoolSize()) { channelProvider = channelProvider.withPoolSize(numPullers); @@ -416,6 +421,7 @@ private void startStreamingConnections() { .setUseLegacyFlowControl(useLegacyFlowControl) .setExecutor(executor) .setSystemExecutor(alarmsExecutor) + .setEodAckCallbackExecutor(eodAckCallbackExecutor) .setClock(clock) .setEnableOpenTelemetryTracing(enableOpenTelemetryTracing) .setTracer(tracer) diff --git a/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnectionTest.java b/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnectionTest.java index 335ccbdc3..559b72822 100644 --- a/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnectionTest.java +++ b/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/StreamingSubscriberConnectionTest.java @@ -587,6 +587,7 @@ private StreamingSubscriberConnection getStreamingSubscriberConnectionFromBuilde .setFlowController(mock(FlowController.class)) .setExecutor(executor) .setSystemExecutor(systemExecutor) + .setEodAckCallbackExecutor(systemExecutor) .setClock(clock) .setMinDurationPerAckExtension(Subscriber.DEFAULT_MIN_ACK_DEADLINE_EXTENSION) .setMinDurationPerAckExtensionDefaultUsed(true) From 41bac6cf2e9e94d0f04bbe2ecf495908a447288e Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 4 Aug 2025 17:27:58 +0000 Subject: [PATCH 3/4] fix: Use daemon threads and same number of threads as system executor for EOD ack callback executor --- .../src/main/java/com/google/cloud/pubsub/v1/Subscriber.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java index ee1075fc6..4f94cf2d9 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java @@ -51,8 +51,8 @@ import java.io.IOException; import java.util.ArrayList; import java.util.List; -import java.util.concurrent.Executors; import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadFactory; import java.util.concurrent.atomic.AtomicInteger; @@ -202,7 +202,8 @@ private Subscriber(Builder builder) { backgroundResources.add(new ExecutorAsBackgroundResource((alarmsExecutor))); } - ExecutorProvider eodAckCallbackExecutorProvider = InstantiatingExecutorProvider.newBuilder().setExecutorThreadCount(2).build(); + ExecutorProvider eodAckCallbackExecutorProvider = + InstantiatingExecutorProvider.newBuilder().setExecutorThreadCount(5).build(); eodAckCallbackExecutor = eodAckCallbackExecutorProvider.getExecutor(); TransportChannelProvider channelProvider = builder.channelProvider; From 6ee66f3eeb0f24b2f4a5acf909d62daad0cae1fd Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 4 Aug 2025 17:30:02 +0000 Subject: [PATCH 4/4] fix: Update formatting --- .../main/java/com/google/cloud/pubsub/v1/Subscriber.java | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java index 4f94cf2d9..27e9a5067 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Subscriber.java @@ -202,9 +202,12 @@ private Subscriber(Builder builder) { backgroundResources.add(new ExecutorAsBackgroundResource((alarmsExecutor))); } - ExecutorProvider eodAckCallbackExecutorProvider = - InstantiatingExecutorProvider.newBuilder().setExecutorThreadCount(5).build(); - eodAckCallbackExecutor = eodAckCallbackExecutorProvider.getExecutor(); + // When exactly-once delivery is enabled, we use a different executor from the system executor + // for handling ack and modack responses to prevent deadlock. + ThreadFactory eodAckCallbackThreadFactory = new ThreadFactoryBuilder().setDaemon(true).build(); + int eodAckCallbackThreadCount = Math.max(6, 2 * numPullers); + eodAckCallbackExecutor = + Executors.newFixedThreadPool(eodAckCallbackThreadCount, eodAckCallbackThreadFactory); TransportChannelProvider channelProvider = builder.channelProvider; if (channelProvider.acceptsPoolSize()) {