From 61d0b608427a53832b47318b6d076bee32e11719 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 13 Jul 2026 15:14:59 +0000 Subject: [PATCH 01/11] Add hedgeSettings with configurable hedgeDelay --- .../google/cloud/pubsub/v1/HedgeSettings.java | 77 +++++++++++++++++++ .../cloud/pubsub/v1/HedgeSettingsTest.java | 68 ++++++++++++++++ 2 files changed, 145 insertions(+) create mode 100644 java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java create mode 100644 java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java new file mode 100644 index 000000000000..9a2dbf66c9dd --- /dev/null +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java @@ -0,0 +1,77 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.pubsub.v1; + +import com.google.common.base.Preconditions; +import java.time.Duration; + +/** Settings for configuring publish hedging. */ +public final class HedgeSettings { + private static final Duration DEFAULT_DELAY = Duration.ofMillis(100); + private static final int DEFAULT_MAX_TOKENS = 100; + private static final float DEFAULT_REFILL = 0.1f; + + private final Duration hedgeDelay; + private final int maxTokens; + private final float refill; + + private HedgeSettings(Builder builder) { + this.hedgeDelay = builder.hedgeDelay; + this.maxTokens = builder.maxTokens; + this.refill = builder.refill; + } + + /** Returns the configured hedging delay. */ + public Duration getHedgeDelay() { + return hedgeDelay; + } + + int getMaxTokens() { + return maxTokens; + } + + float getRefill() { + return refill; + } + + /** Returns a new builder for {@code HedgeSettings}. */ + public static Builder newBuilder() { + return new Builder(); + } + + /** Builder for {@code HedgeSettings}. */ + public static final class Builder { + private Duration hedgeDelay = DEFAULT_DELAY; + private int maxTokens = DEFAULT_MAX_TOKENS; + private float refill = DEFAULT_REFILL; + + private Builder() {} + + /** Allows hedging delay to be configurable. */ + public Builder setHedgeDelay(Duration hedgeDelay) { + Preconditions.checkNotNull(hedgeDelay); + Preconditions.checkArgument(!hedgeDelay.isNegative() && !hedgeDelay.isZero(), "hedgeDelay must be positive"); + this.hedgeDelay = hedgeDelay; + return this; + } + + /** Builds an instance of {@code HedgeSettings}. */ + public HedgeSettings build() { + return new HedgeSettings(this); + } + } +} diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java new file mode 100644 index 000000000000..2a01408b7bdf --- /dev/null +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java @@ -0,0 +1,68 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.pubsub.v1; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThrows; + +import java.time.Duration; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +@RunWith(JUnit4.class) +public class HedgeSettingsTest { + + @Test + public void testDefaultSettings() { + HedgeSettings settings = HedgeSettings.newBuilder().build(); + assertNotNull(settings); + assertEquals(Duration.ofMillis(100), settings.getHedgeDelay()); + assertEquals(100, settings.getMaxTokens()); + assertEquals(0.1f, settings.getRefill(), 0.0001f); + } + + @Test + public void testCustomDelay() { + Duration customDelay = Duration.ofMillis(50); + HedgeSettings settings = HedgeSettings.newBuilder().setHedgeDelay(customDelay).build(); + assertNotNull(settings); + assertEquals(customDelay, settings.getHedgeDelay()); + } + + @Test + public void testNegativeDelayThrows() { + assertThrows( + IllegalArgumentException.class, + () -> HedgeSettings.newBuilder().setHedgeDelay(Duration.ofMillis(-10))); + } + + @Test + public void testZeroDelayThrows() { + assertThrows( + IllegalArgumentException.class, + () -> HedgeSettings.newBuilder().setHedgeDelay(Duration.ZERO)); + } + + @Test + public void testNullDelayThrows() { + assertThrows( + NullPointerException.class, + () -> HedgeSettings.newBuilder().setHedgeDelay(null)); + } +} From 19afd68374fbd8899e1f5ecd4ad3e15ee2a001d1 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 13 Jul 2026 16:08:47 +0000 Subject: [PATCH 02/11] Add HedgeTokenBucket class to allow for token operations --- .../cloud/pubsub/v1/HedgeTokenBucket.java | 65 +++++++++++++++++ .../cloud/pubsub/v1/HedgeTokenBucketTest.java | 71 +++++++++++++++++++ 2 files changed, 136 insertions(+) create mode 100644 java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeTokenBucket.java create mode 100644 java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeTokenBucketTest.java diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeTokenBucket.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeTokenBucket.java new file mode 100644 index 000000000000..5378478cb429 --- /dev/null +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeTokenBucket.java @@ -0,0 +1,65 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.pubsub.v1; + +import com.google.common.annotations.VisibleForTesting; +import javax.annotation.concurrent.GuardedBy; + +/** Token bucket for limiting hedged requests. Thread-safe. */ +final class HedgeTokenBucket { + private final int maxTokens; + private final float refillAmount; + + @GuardedBy("this") + private float tokens; + + HedgeTokenBucket(HedgeSettings settings) { + this.maxTokens = settings.getMaxTokens(); + this.refillAmount = settings.getRefill(); + this.tokens = maxTokens; + } + + @VisibleForTesting + HedgeTokenBucket(int maxTokens, float refillAmount) { + this.maxTokens = maxTokens; + this.refillAmount = refillAmount; + this.tokens = maxTokens; + } + + /** + * Attempts to acquire a token for a hedged request. + * + * @return {@code true} if a token was acquired, {@code false} otherwise. + */ + synchronized boolean tryAcquire() { + if (tokens >= 1.0f) { + tokens -= 1.0f; + return true; + } + return false; + } + + /** Refills the bucket by the configured refill amount, capped at max tokens. */ + synchronized void refill() { + tokens = Math.min(maxTokens, tokens + refillAmount); + } + + @VisibleForTesting + synchronized float getTokens() { + return tokens; + } +} diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeTokenBucketTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeTokenBucketTest.java new file mode 100644 index 000000000000..e92ecd9a6ce2 --- /dev/null +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeTokenBucketTest.java @@ -0,0 +1,71 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.pubsub.v1; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +@RunWith(JUnit4.class) +public class HedgeTokenBucketTest { + + @Test + public void testInitialState() { + HedgeTokenBucket bucket = new HedgeTokenBucket(10, 0.5f); + assertEquals(10.0f, bucket.getTokens(), 0.0001f); + } + + @Test + public void testAcquire() { + HedgeTokenBucket bucket = new HedgeTokenBucket(10, 0.5f); + assertTrue(bucket.tryAcquire()); + assertEquals(9.0f, bucket.getTokens(), 0.0001f); + } + + @Test + public void testAcquireUntilEmpty() { + HedgeTokenBucket bucket = new HedgeTokenBucket(2, 0.5f); + assertTrue(bucket.tryAcquire()); + assertTrue(bucket.tryAcquire()); + assertFalse(bucket.tryAcquire()); + assertEquals(0.0f, bucket.getTokens(), 0.0001f); + } + + @Test + public void testRefill() { + HedgeTokenBucket bucket = new HedgeTokenBucket(2, 0.5f); + assertTrue(bucket.tryAcquire()); // 1.0 left + assertTrue(bucket.tryAcquire()); // 0.0 left + + bucket.refill(); // 0.5 tokens + assertFalse(bucket.tryAcquire()); // needs 1.0 + + bucket.refill(); // 1.0 tokens + assertTrue(bucket.tryAcquire()); // succeeds, 0.0 left + } + + @Test + public void testRefillCap() { + HedgeTokenBucket bucket = new HedgeTokenBucket(2, 0.5f); + bucket.refill(); // already at max (2) + assertEquals(2.0f, bucket.getTokens(), 0.0001f); + } +} From 49426a9b490abb6a6ecb7d0f78cf26e66f1a4816 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 13 Jul 2026 16:43:24 +0000 Subject: [PATCH 03/11] Add publisher integration with HedgeSettings --- .../com/google/cloud/pubsub/v1/Publisher.java | 22 +++++++++ .../cloud/pubsub/v1/PublisherImplTest.java | 48 ++++++++++++++----- 2 files changed, 59 insertions(+), 11 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index 56c920bcfdc1..1654cb364418 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -138,6 +138,9 @@ public class Publisher implements PublisherInterface { private final OpenTelemetry openTelemetry; private OpenTelemetryPubsubTracer tracer = new OpenTelemetryPubsubTracer(null, false); + private final HedgeSettings hedgeSettings; + private final HedgeTokenBucket hedgeTokenBucket; + /** The maximum number of messages in one request. Defined by the API. */ public static long getApiMaxRequestElementCount() { return 1000L; @@ -230,6 +233,9 @@ private Publisher(Builder builder) throws IOException { backgroundResources = new BackgroundResourceAggregation(backgroundResourceList); shutdown = new AtomicBoolean(false); messagesWaiter = new Waiter(); + this.hedgeSettings = builder.hedgeSettings; + this.hedgeTokenBucket = + this.hedgeSettings != null ? new HedgeTokenBucket(this.hedgeSettings) : null; this.publishContext = GrpcCallContext.createDefault(); this.publishContextWithCompression = GrpcCallContext.createDefault() @@ -246,6 +252,15 @@ public String getTopicNameString() { return topicName; } + /** Returns the configured hedging settings, or null if hedging is disabled. */ + public HedgeSettings getHedgeSettings() { + return hedgeSettings; + } + + HedgeTokenBucket getHedgeTokenBucket() { + return hedgeTokenBucket; + } + /** * Schedules the publishing of a message. The publishing of the message may occur immediately or * be delayed based on the publisher batching options. @@ -814,6 +829,7 @@ public PubsubMessage apply(PubsubMessage input) { private boolean enableOpenTelemetryTracing = false; private OpenTelemetry openTelemetry = null; + private HedgeSettings hedgeSettings = null; private Builder(String topic) { this.topicName = Preconditions.checkNotNull(topic); @@ -966,6 +982,12 @@ public Builder setOpenTelemetry(OpenTelemetry openTelemetry) { return this; } + /** Configures the Publisher's hedging parameters. */ + public Builder setHedgeSettings(HedgeSettings hedgeSettings) { + this.hedgeSettings = hedgeSettings; + return this; + } + /** Returns the default BatchingSettings used by the client if settings are not provided. */ public static BatchingSettings getDefaultBatchingSettings() { return DEFAULT_BATCHING_SETTINGS; diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index 8e6efaf372c9..a5bd92c2aa0d 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -64,7 +64,6 @@ import org.easymock.EasyMock; import org.junit.After; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.JUnit4; @@ -513,20 +512,21 @@ public void testEnableMessageOrdering_overwritesMaxAttempts() throws Exception { *
  • publish with key orderA, which should now succeed * */ - @Ignore("https://github.com/googleapis/google-cloud-java/issues/13394") + /* + Temporarily disabled due to https://github.com/googleapis/java-pubsub/issues/1861. + TODO(maitrimangal): Enable once resolved. @Test public void testResumePublish() throws Exception { Publisher publisher = getTestPublisherBuilder() .setBatchingSettings( - Publisher.Builder.DEFAULT_BATCHING_SETTINGS.toBuilder() + Publisher.Builder.DEFAULT_BATCHING_SETTINGS + .toBuilder() .setElementCountThreshold(2L) .build()) .setEnableMessageOrdering(true) .build(); - testPublisherServiceImpl.setExecutor(fakeExecutor); - ApiFuture future1 = sendTestMessageWithOrderingKey(publisher, "m1", "orderA"); ApiFuture future2 = sendTestMessageWithOrderingKey(publisher, "m2", "orderA"); @@ -536,6 +536,7 @@ public void testResumePublish() throws Exception { // This exception should stop future publishing to the same key testPublisherServiceImpl.addPublishError(new StatusException(Status.INVALID_ARGUMENT)); + fakeExecutor.advanceTime(Duration.ZERO); try { @@ -575,8 +576,8 @@ public void testResumePublish() throws Exception { testPublisherServiceImpl.addPublishResponse( PublishResponse.newBuilder().addMessageIds("5").addMessageIds("6")); - assertEquals("5", future5.get()); - assertEquals("6", future6.get()); + Assert.assertEquals("5", future5.get()); + Assert.assertEquals("6", future6.get()); // Resume publishing of "orderA", which should now succeed publisher.resumePublish("orderA"); @@ -587,19 +588,20 @@ public void testResumePublish() throws Exception { testPublisherServiceImpl.addPublishResponse( PublishResponse.newBuilder().addMessageIds("7").addMessageIds("8")); - assertEquals("7", future7.get()); - assertEquals("8", future8.get()); + Assert.assertEquals("7", future7.get()); + Assert.assertEquals("8", future8.get()); shutdownTestPublisher(publisher); } - @Ignore("https://github.com/googleapis/google-cloud-java/issues/13394") @Test public void testPublishThrowExceptionForUnsubmittedOrderingKeyMessage() throws Exception { Publisher publisher = getTestPublisherBuilder() + .setExecutorProvider(SINGLE_THREAD_EXECUTOR) .setBatchingSettings( - Publisher.Builder.DEFAULT_BATCHING_SETTINGS.toBuilder() + Publisher.Builder.DEFAULT_BATCHING_SETTINGS + .toBuilder() .setElementCountThreshold(2L) .setDelayThresholdDuration(Duration.ofSeconds(500)) .build()) @@ -642,6 +644,7 @@ public void testPublishThrowExceptionForUnsubmittedOrderingKeyMessage() throws E assertEquals(SequentialExecutorService.CallbackExecutor.CANCELLATION_EXCEPTION, e.getCause()); } } + */ private ApiFuture sendTestMessageWithOrderingKey( Publisher publisher, String data, String orderingKey) { @@ -1339,6 +1342,29 @@ public void testPublishOpenTelemetryTracing() throws Exception { .hasEnded(); } + @Test + public void testPublisherWithHedgeSettings() throws Exception { + HedgeSettings hedgeSettings = + HedgeSettings.newBuilder().setHedgeDelay(Duration.ofMillis(50)).build(); + Publisher publisher = getTestPublisherBuilder().setHedgeSettings(hedgeSettings).build(); + + assertThat(publisher.getHedgeSettings()).isEqualTo(hedgeSettings); + assertThat(publisher.getHedgeTokenBucket()).isNotNull(); + assertThat(publisher.getHedgeTokenBucket().getTokens()).isWithin(0.0001f).of(100.0f); + + shutdownTestPublisher(publisher); + } + + @Test + public void testPublisherWithoutHedgeSettings() throws Exception { + Publisher publisher = getTestPublisherBuilder().build(); + + assertThat(publisher.getHedgeSettings()).isNull(); + assertThat(publisher.getHedgeTokenBucket()).isNull(); + + shutdownTestPublisher(publisher); + } + private Builder getTestPublisherBuilder() { return Publisher.newBuilder(TEST_TOPIC) .setExecutorProvider(FixedExecutorProvider.create(fakeExecutor)) From 22417d156d31e10398f225c6f98ccebae9f93595 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 13 Jul 2026 16:51:32 +0000 Subject: [PATCH 04/11] Restore original @Ignore annotations and code in PublisherImplTest --- .../cloud/pubsub/v1/PublisherImplTest.java | 48 +++++-------------- 1 file changed, 11 insertions(+), 37 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index a5bd92c2aa0d..8e6efaf372c9 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -64,6 +64,7 @@ import org.easymock.EasyMock; import org.junit.After; import org.junit.Before; +import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.JUnit4; @@ -512,21 +513,20 @@ public void testEnableMessageOrdering_overwritesMaxAttempts() throws Exception { *
  • publish with key orderA, which should now succeed * */ - /* - Temporarily disabled due to https://github.com/googleapis/java-pubsub/issues/1861. - TODO(maitrimangal): Enable once resolved. + @Ignore("https://github.com/googleapis/google-cloud-java/issues/13394") @Test public void testResumePublish() throws Exception { Publisher publisher = getTestPublisherBuilder() .setBatchingSettings( - Publisher.Builder.DEFAULT_BATCHING_SETTINGS - .toBuilder() + Publisher.Builder.DEFAULT_BATCHING_SETTINGS.toBuilder() .setElementCountThreshold(2L) .build()) .setEnableMessageOrdering(true) .build(); + testPublisherServiceImpl.setExecutor(fakeExecutor); + ApiFuture future1 = sendTestMessageWithOrderingKey(publisher, "m1", "orderA"); ApiFuture future2 = sendTestMessageWithOrderingKey(publisher, "m2", "orderA"); @@ -536,7 +536,6 @@ public void testResumePublish() throws Exception { // This exception should stop future publishing to the same key testPublisherServiceImpl.addPublishError(new StatusException(Status.INVALID_ARGUMENT)); - fakeExecutor.advanceTime(Duration.ZERO); try { @@ -576,8 +575,8 @@ public void testResumePublish() throws Exception { testPublisherServiceImpl.addPublishResponse( PublishResponse.newBuilder().addMessageIds("5").addMessageIds("6")); - Assert.assertEquals("5", future5.get()); - Assert.assertEquals("6", future6.get()); + assertEquals("5", future5.get()); + assertEquals("6", future6.get()); // Resume publishing of "orderA", which should now succeed publisher.resumePublish("orderA"); @@ -588,20 +587,19 @@ public void testResumePublish() throws Exception { testPublisherServiceImpl.addPublishResponse( PublishResponse.newBuilder().addMessageIds("7").addMessageIds("8")); - Assert.assertEquals("7", future7.get()); - Assert.assertEquals("8", future8.get()); + assertEquals("7", future7.get()); + assertEquals("8", future8.get()); shutdownTestPublisher(publisher); } + @Ignore("https://github.com/googleapis/google-cloud-java/issues/13394") @Test public void testPublishThrowExceptionForUnsubmittedOrderingKeyMessage() throws Exception { Publisher publisher = getTestPublisherBuilder() - .setExecutorProvider(SINGLE_THREAD_EXECUTOR) .setBatchingSettings( - Publisher.Builder.DEFAULT_BATCHING_SETTINGS - .toBuilder() + Publisher.Builder.DEFAULT_BATCHING_SETTINGS.toBuilder() .setElementCountThreshold(2L) .setDelayThresholdDuration(Duration.ofSeconds(500)) .build()) @@ -644,7 +642,6 @@ public void testPublishThrowExceptionForUnsubmittedOrderingKeyMessage() throws E assertEquals(SequentialExecutorService.CallbackExecutor.CANCELLATION_EXCEPTION, e.getCause()); } } - */ private ApiFuture sendTestMessageWithOrderingKey( Publisher publisher, String data, String orderingKey) { @@ -1342,29 +1339,6 @@ public void testPublishOpenTelemetryTracing() throws Exception { .hasEnded(); } - @Test - public void testPublisherWithHedgeSettings() throws Exception { - HedgeSettings hedgeSettings = - HedgeSettings.newBuilder().setHedgeDelay(Duration.ofMillis(50)).build(); - Publisher publisher = getTestPublisherBuilder().setHedgeSettings(hedgeSettings).build(); - - assertThat(publisher.getHedgeSettings()).isEqualTo(hedgeSettings); - assertThat(publisher.getHedgeTokenBucket()).isNotNull(); - assertThat(publisher.getHedgeTokenBucket().getTokens()).isWithin(0.0001f).of(100.0f); - - shutdownTestPublisher(publisher); - } - - @Test - public void testPublisherWithoutHedgeSettings() throws Exception { - Publisher publisher = getTestPublisherBuilder().build(); - - assertThat(publisher.getHedgeSettings()).isNull(); - assertThat(publisher.getHedgeTokenBucket()).isNull(); - - shutdownTestPublisher(publisher); - } - private Builder getTestPublisherBuilder() { return Publisher.newBuilder(TEST_TOPIC) .setExecutorProvider(FixedExecutorProvider.create(fakeExecutor)) From 91706a2a212828ba92fb22dc9aae51b67bc9f6a6 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 13 Jul 2026 17:39:56 +0000 Subject: [PATCH 05/11] style(pubsub): fix checkstyle violations in HedgeSettings --- .../google/cloud/pubsub/v1/HedgeSettings.java | 48 +++++++++++++++---- 1 file changed, 38 insertions(+), 10 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java index 9a2dbf66c9dd..4641346fca31 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java @@ -21,21 +21,31 @@ /** Settings for configuring publish hedging. */ public final class HedgeSettings { + /** Default hedging delay. */ private static final Duration DEFAULT_DELAY = Duration.ofMillis(100); + /** Default maximum number of tokens in the bucket. */ private static final int DEFAULT_MAX_TOKENS = 100; + /** Default refill rate (tokens per successful request). */ private static final float DEFAULT_REFILL = 0.1f; + /** Hedging delay. */ private final Duration hedgeDelay; + /** Maximum tokens. */ private final int maxTokens; + /** Refill rate. */ private final float refill; - private HedgeSettings(Builder builder) { + private HedgeSettings(final Builder builder) { this.hedgeDelay = builder.hedgeDelay; this.maxTokens = builder.maxTokens; this.refill = builder.refill; } - /** Returns the configured hedging delay. */ + /** + * Returns the configured hedging delay. + * + * @return the hedging delay. + */ public Duration getHedgeDelay() { return hedgeDelay; } @@ -48,28 +58,46 @@ float getRefill() { return refill; } - /** Returns a new builder for {@code HedgeSettings}. */ + /** + * Returns a new builder for {@code HedgeSettings}. + * + * @return a new builder. + */ public static Builder newBuilder() { return new Builder(); } /** Builder for {@code HedgeSettings}. */ public static final class Builder { + /** Hedging delay. */ private Duration hedgeDelay = DEFAULT_DELAY; + /** Maximum tokens. */ private int maxTokens = DEFAULT_MAX_TOKENS; + /** Refill rate. */ private float refill = DEFAULT_REFILL; - private Builder() {} + private Builder() { } - /** Allows hedging delay to be configurable. */ - public Builder setHedgeDelay(Duration hedgeDelay) { - Preconditions.checkNotNull(hedgeDelay); - Preconditions.checkArgument(!hedgeDelay.isNegative() && !hedgeDelay.isZero(), "hedgeDelay must be positive"); - this.hedgeDelay = hedgeDelay; + /** + * Allows hedging delay to be configurable. + * + * @param delay the hedging delay, must be positive. + * @return this builder. + */ + public Builder setHedgeDelay(final Duration delay) { + Preconditions.checkNotNull(delay); + Preconditions.checkArgument( + !delay.isNegative() && !delay.isZero(), + "hedgeDelay must be positive"); + this.hedgeDelay = delay; return this; } - /** Builds an instance of {@code HedgeSettings}. */ + /** + * Builds an instance of {@code HedgeSettings}. + * + * @return the built {@code HedgeSettings} instance. + */ public HedgeSettings build() { return new HedgeSettings(this); } From bd76ede662b5a951773c34cab4f4bd79b7c9373e Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 13 Jul 2026 17:53:42 +0000 Subject: [PATCH 06/11] style(pubsub): fix formatting and checkstyle violations in HedgeSettings and tests --- .../google/cloud/pubsub/v1/HedgeSettings.java | 16 +++++++++---- .../cloud/pubsub/v1/HedgeSettingsTest.java | 4 +--- .../cloud/pubsub/v1/PublisherImplTest.java | 23 +++++++++++++++++++ 3 files changed, 36 insertions(+), 7 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java index 4641346fca31..ffa6daed1d91 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java @@ -23,15 +23,19 @@ public final class HedgeSettings { /** Default hedging delay. */ private static final Duration DEFAULT_DELAY = Duration.ofMillis(100); + /** Default maximum number of tokens in the bucket. */ private static final int DEFAULT_MAX_TOKENS = 100; + /** Default refill rate (tokens per successful request). */ private static final float DEFAULT_REFILL = 0.1f; /** Hedging delay. */ private final Duration hedgeDelay; + /** Maximum tokens. */ private final int maxTokens; + /** Refill rate. */ private final float refill; @@ -71,12 +75,16 @@ public static Builder newBuilder() { public static final class Builder { /** Hedging delay. */ private Duration hedgeDelay = DEFAULT_DELAY; + /** Maximum tokens. */ private int maxTokens = DEFAULT_MAX_TOKENS; + /** Refill rate. */ private float refill = DEFAULT_REFILL; - private Builder() { } + private Builder() { + super(); + } /** * Allows hedging delay to be configurable. @@ -86,9 +94,9 @@ private Builder() { } */ public Builder setHedgeDelay(final Duration delay) { Preconditions.checkNotNull(delay); - Preconditions.checkArgument( - !delay.isNegative() && !delay.isZero(), - "hedgeDelay must be positive"); + if (delay.isNegative() || delay.isZero()) { + throw new IllegalArgumentException("delay must be positive"); + } this.hedgeDelay = delay; return this; } diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java index 2a01408b7bdf..f3b26aedb3c3 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/HedgeSettingsTest.java @@ -61,8 +61,6 @@ public void testZeroDelayThrows() { @Test public void testNullDelayThrows() { - assertThrows( - NullPointerException.class, - () -> HedgeSettings.newBuilder().setHedgeDelay(null)); + assertThrows(NullPointerException.class, () -> HedgeSettings.newBuilder().setHedgeDelay(null)); } } diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index 8e6efaf372c9..95b19d6b0108 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -1339,6 +1339,29 @@ public void testPublishOpenTelemetryTracing() throws Exception { .hasEnded(); } + @Test + public void testPublisherWithHedgeSettings() throws Exception { + HedgeSettings hedgeSettings = + HedgeSettings.newBuilder().setHedgeDelay(Duration.ofMillis(50)).build(); + Publisher publisher = getTestPublisherBuilder().setHedgeSettings(hedgeSettings).build(); + + assertThat(publisher.getHedgeSettings()).isEqualTo(hedgeSettings); + assertThat(publisher.getHedgeTokenBucket()).isNotNull(); + assertThat(publisher.getHedgeTokenBucket().getTokens()).isWithin(0.0001f).of(100.0f); + + shutdownTestPublisher(publisher); + } + + @Test + public void testPublisherWithoutHedgeSettings() throws Exception { + Publisher publisher = getTestPublisherBuilder().build(); + + assertThat(publisher.getHedgeSettings()).isNull(); + assertThat(publisher.getHedgeTokenBucket()).isNull(); + + shutdownTestPublisher(publisher); + } + private Builder getTestPublisherBuilder() { return Publisher.newBuilder(TEST_TOPIC) .setExecutorProvider(FixedExecutorProvider.create(fakeExecutor)) From 227f808372e013995803d024d003efd0acd71a48 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 20 Jul 2026 15:17:04 +0000 Subject: [PATCH 07/11] Resolve git comments related to access modifiers and constraints --- .../main/java/com/google/cloud/pubsub/v1/HedgeSettings.java | 3 +-- .../src/main/java/com/google/cloud/pubsub/v1/Publisher.java | 3 +++ 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java index ffa6daed1d91..62e16aa71639 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java @@ -50,7 +50,7 @@ private HedgeSettings(final Builder builder) { * * @return the hedging delay. */ - public Duration getHedgeDelay() { + Duration getHedgeDelay() { return hedgeDelay; } @@ -83,7 +83,6 @@ public static final class Builder { private float refill = DEFAULT_REFILL; private Builder() { - super(); } /** diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index 1654cb364418..d982bc03a210 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -994,6 +994,9 @@ public static BatchingSettings getDefaultBatchingSettings() { } public Publisher build() throws IOException { + Preconditions.checkState( + !(enableMessageOrdering && hedgeSettings != null), + "Publish hedging and message ordering cannot be enabled at the same time."); return new Publisher(this); } } From e29eab859ce0904461f5df96898ba23d3bd23910 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 20 Jul 2026 15:51:53 +0000 Subject: [PATCH 08/11] Add CancellationSharer and HedgedRequest helper classes --- .../cloud/pubsub/v1/CancellationSharer.java | 141 ++++++++++++++++++ .../google/cloud/pubsub/v1/HedgedRequest.java | 44 ++++++ 2 files changed, 185 insertions(+) create mode 100644 java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java create mode 100644 java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java new file mode 100644 index 000000000000..50995e88f489 --- /dev/null +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java @@ -0,0 +1,141 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.pubsub.v1; + +import com.google.api.core.AbstractApiFuture; +import com.google.api.core.ApiFuture; +import com.google.api.core.ApiFutureCallback; +import com.google.api.core.ApiFutures; +import com.google.common.util.concurrent.MoreExecutors; +import com.google.pubsub.v1.PublishResponse; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +/** + * Coordinates multiple publish attempts for a single batch of messages. + * + *

    Implements {@link ApiFuture} to act as the single future returned to the publisher's client. + * It manages the lifecycle of the original attempt and any subsequent hedged attempts (CancellationSharer in the diagram). + */ +class CancellationSharer extends AbstractApiFuture { + private final Publisher.OutstandingBatch batch; + private final Publisher publisher; + private final Map> runningAttempts = new ConcurrentHashMap<>(); + private final AtomicInteger totalAttemptsCount = new AtomicInteger(0); + private final AtomicBoolean done = new AtomicBoolean(false); + final AtomicBoolean isInQueue = new AtomicBoolean(false); + private final AtomicReference lastError = new AtomicReference<>(); + + CancellationSharer(Publisher.OutstandingBatch batch, Publisher publisher) { + this.batch = batch; + this.publisher = publisher; + } + + /** + * Adds an attempt to be tracked by this coordinator. + * + * @param attemptNumber the 1-based index of the attempt (1 is original, 2+ are hedged) + * @param future the future representing the gRPC call for this attempt + */ + void addAttempt(final int attemptNumber, ApiFuture future) { + runningAttempts.put(attemptNumber, future); + totalAttemptsCount.incrementAndGet(); + + if (done.get()) { + future.cancel(true); + runningAttempts.remove(attemptNumber); + return; + } + + ApiFutures.addCallback( + future, + new ApiFutureCallback() { + @Override + public void onSuccess(PublishResponse result) { + handleAttemptSuccess(attemptNumber, result); + } + + @Override + public void onFailure(Throwable t) { + handleAttemptFailure(attemptNumber, t); + } + }, + MoreExecutors.directExecutor()); + } + + private void handleAttemptSuccess(int attemptNumber, PublishResponse response) { + if (done.compareAndSet(false, true)) { + set(response); // Resolve parent future + cancelAllExcept(attemptNumber); + publisher.refillTokenBucket(); + } + } + + private void handleAttemptFailure(int attemptNumber, Throwable t) { + runningAttempts.remove(attemptNumber); + + if (!done.get()) { + lastError.set(t); + if (runningAttempts.isEmpty() && !isInQueue.get()) { + if (done.compareAndSet(false, true)) { + setException(lastError.get()); + } + } + } + } + + void checkCompletionOnQueueExit() { + if (!done.get() && runningAttempts.isEmpty() && !isInQueue.get()) { + if (done.compareAndSet(false, true)) { + Throwable error = lastError.get(); + setException(error != null ? error : new RuntimeException("Hedging failed with no active attempts")); + } + } + } + + private void cancelAllExcept(int successfulAttempt) { + for (Map.Entry> entry : runningAttempts.entrySet()) { + if (entry.getKey() != successfulAttempt) { + entry.getValue().cancel(true); + } + } + } + + @Override + public boolean cancel(boolean mayInterruptIfRunning) { + if (super.cancel(mayInterruptIfRunning)) { + if (done.compareAndSet(false, true)) { + for (ApiFuture future : runningAttempts.values()) { + future.cancel(mayInterruptIfRunning); + } + return true; + } + } + return false; + } + + int getAttemptCount() { + return totalAttemptsCount.get(); + } + + Publisher.OutstandingBatch getBatch() { + return batch; + } +} diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java new file mode 100644 index 000000000000..ac405032c8e1 --- /dev/null +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java @@ -0,0 +1,44 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.pubsub.v1; + +/** + * Represents a pending hedging check in the publisher's queue. + */ +class HedgedRequest { + private final CancellationSharer coordinator; + private final int attemptNumber; + private final long sendAfterMs; + + HedgedRequest(CancellationSharer coordinator, int attemptNumber, long sendAfterMs) { + this.coordinator = coordinator; + this.attemptNumber = attemptNumber; + this.sendAfterMs = sendAfterMs; + } + + CancellationSharer getCoordinator() { + return coordinator; + } + + int getAttemptNumber() { + return attemptNumber; + } + + long getSendAfterMs() { + return sendAfterMs; + } +} From ed05c60650c1e8496124abff2a7c49986c8d8e48 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 20 Jul 2026 19:23:16 +0000 Subject: [PATCH 09/11] Integrate CancellationSharer and HedgedRequest into Publisher.java. Also add Unit tests --- .../cloud/pubsub/v1/CancellationSharer.java | 3 + .../com/google/cloud/pubsub/v1/Publisher.java | 126 ++++++++++++- .../pubsub/v1/FakePublisherServiceImpl.java | 3 +- .../cloud/pubsub/v1/PublisherImplTest.java | 172 ++++++++++++++++++ 4 files changed, 300 insertions(+), 4 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java index 50995e88f489..37de9324eedf 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java @@ -28,6 +28,8 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; + + /** * Coordinates multiple publish attempts for a single batch of messages. * @@ -43,6 +45,7 @@ class CancellationSharer extends AbstractApiFuture { final AtomicBoolean isInQueue = new AtomicBoolean(false); private final AtomicReference lastError = new AtomicReference<>(); + CancellationSharer(Publisher.OutstandingBatch batch, Publisher publisher) { this.batch = batch; this.publisher = publisher; diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index d982bc03a210..b57442d5375f 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -18,11 +18,13 @@ import static com.google.common.util.concurrent.MoreExecutors.directExecutor; +import com.google.api.core.ApiClock; import com.google.api.core.ApiFunction; import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutureCallback; import com.google.api.core.ApiFutures; import com.google.api.core.BetaApi; +import com.google.api.core.CurrentMillisClock; import com.google.api.core.SettableApiFuture; import com.google.api.gax.batching.BatchingSettings; import com.google.api.gax.batching.FlowControlSettings; @@ -70,7 +72,10 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.logging.Level; @@ -140,6 +145,12 @@ public class Publisher implements PublisherInterface { private final HedgeSettings hedgeSettings; private final HedgeTokenBucket hedgeTokenBucket; + private final ApiClock clock; + + private final ConcurrentLinkedQueue hedgingQueue; + private final AtomicBoolean isQueueProcessingScheduled; + private ScheduledFuture queueProcessingFuture; + private final Object queueLock; /** The maximum number of messages in one request. Defined by the API. */ public static long getApiMaxRequestElementCount() { @@ -236,10 +247,15 @@ private Publisher(Builder builder) throws IOException { this.hedgeSettings = builder.hedgeSettings; this.hedgeTokenBucket = this.hedgeSettings != null ? new HedgeTokenBucket(this.hedgeSettings) : null; + this.clock = builder.clock != null ? builder.clock : CurrentMillisClock.getDefaultClock(); this.publishContext = GrpcCallContext.createDefault(); this.publishContextWithCompression = GrpcCallContext.createDefault() .withCallOptions(CallOptions.DEFAULT.withCompression(GZIP_COMPRESSION)); + this.hedgingQueue = new ConcurrentLinkedQueue<>(); + this.isQueueProcessingScheduled = new AtomicBoolean(false); + this.queueLock = new Object(); + this.queueProcessingFuture = null; } /** Topic which the publisher publishes to. */ @@ -587,7 +603,11 @@ public void onFailure(Throwable t) { ApiFuture future; Executor callbackExecutor = directExecutor(); if (outstandingBatch.orderingKey == null || outstandingBatch.orderingKey.isEmpty()) { - future = publishCall(outstandingBatch); + if (hedgeSettings != null) { + future = startHedgedCall(outstandingBatch); + } else { + future = publishCall(outstandingBatch); + } } else { // If ordering key is specified, publish the batch using the sequential executor. future = @@ -603,7 +623,13 @@ public ApiFuture call() { ApiFutures.addCallback(future, futureCallback, callbackExecutor); } - private final class OutstandingBatch { + void refillTokenBucket() { + if (hedgeTokenBucket != null) { + hedgeTokenBucket.refill(); + } + } + + final class OutstandingBatch { final List outstandingPublishes; final long creationTime; int attempt; @@ -615,7 +641,7 @@ private final class OutstandingBatch { List outstandingPublishes, int batchSizeBytes, String orderingKey) { this.outstandingPublishes = outstandingPublishes; attempt = 1; - creationTime = System.currentTimeMillis(); + creationTime = clock.millisTime(); this.batchSizeBytes = batchSizeBytes; this.orderingKey = orderingKey; } @@ -692,6 +718,9 @@ public void shutdown() { if (currentAlarmFuture != null && activeAlarm.getAndSet(false)) { currentAlarmFuture.cancel(false); } + if (queueProcessingFuture != null) { + queueProcessingFuture.cancel(false); + } publishAllOutstanding(); messagesWaiter.waitComplete(); backgroundResources.shutdown(); @@ -830,6 +859,7 @@ public PubsubMessage apply(PubsubMessage input) { private boolean enableOpenTelemetryTracing = false; private OpenTelemetry openTelemetry = null; private HedgeSettings hedgeSettings = null; + ApiClock clock = null; private Builder(String topic) { this.topicName = Preconditions.checkNotNull(topic); @@ -988,6 +1018,11 @@ public Builder setHedgeSettings(HedgeSettings hedgeSettings) { return this; } + Builder setClock(ApiClock clock) { + this.clock = clock; + return this; + } + /** Returns the default BatchingSettings used by the client if settings are not provided. */ public static BatchingSettings getDefaultBatchingSettings() { return DEFAULT_BATCHING_SETTINGS; @@ -1207,4 +1242,89 @@ && getBatchedBytes() + outstandingPublish.messageSize >= getMaxBatchBytes()) { return batchesToSend; } } + + private ApiFuture startHedgedCall(final OutstandingBatch outstandingBatch) { + final CancellationSharer coordinator = new CancellationSharer(outstandingBatch, this); + + // Register cancellation listeners on client futures to propagate cancel to coordinator + final AtomicInteger cancelledCount = new AtomicInteger(0); + final int batchSize = outstandingBatch.outstandingPublishes.size(); + for (final OutstandingPublish outstanding : outstandingBatch.outstandingPublishes) { + outstanding.publishResult.addListener(new Runnable() { + @Override + public void run() { + if (outstanding.publishResult.isCancelled()) { + if (cancelledCount.incrementAndGet() == batchSize) { + coordinator.cancel(true); + } + } + } + }, directExecutor()); + } + + ApiFuture firstAttemptFuture = publishCall(outstandingBatch); + coordinator.addAttempt(1, firstAttemptFuture); + long delayMs = hedgeSettings.getHedgeDelay().toMillis(); + HedgedRequest item = new HedgedRequest(coordinator, 2, clock.millisTime() + delayMs); + hedgingQueue.add(item); + coordinator.isInQueue.set(true); + scheduleQueueProcessing(); + + return coordinator; + } + + private void scheduleQueueProcessing() { + if (isQueueProcessingScheduled.compareAndSet(false, true)) { + HedgedRequest nextItem = hedgingQueue.peek(); + if (nextItem == null) { + isQueueProcessingScheduled.set(false); + return; + } + + long delay = nextItem.getSendAfterMs() - clock.millisTime(); + delay = Math.max(0, delay); + + queueProcessingFuture = executor.schedule(new Runnable() { + @Override + public void run() { + processQueue(); + } + }, delay, TimeUnit.MILLISECONDS); + } + } + + private void processQueue() { + synchronized (queueLock) { + isQueueProcessingScheduled.set(false); + long now = clock.millisTime(); + + HedgedRequest item; + while ((item = hedgingQueue.peek()) != null && item.getSendAfterMs() <= now) { + hedgingQueue.poll(); + + CancellationSharer coordinator = item.getCoordinator(); + if (coordinator.isDone()) { + coordinator.isInQueue.set(false); + continue; + } + + if (hedgeTokenBucket.tryAcquire()) { + // Clone and schedule next attempt check (Attempt + 1) + long delayMs = hedgeSettings.getHedgeDelay().toMillis(); + HedgedRequest nextItem = new HedgedRequest(coordinator, item.getAttemptNumber() + 1, now + delayMs); + hedgingQueue.add(nextItem); + + // Start Hedged Attempt + ApiFuture hedgedFuture = publishCall(coordinator.getBatch()); + coordinator.addAttempt(item.getAttemptNumber(), hedgedFuture); + } else { + coordinator.isInQueue.set(false); + coordinator.checkCompletionOnQueueExit(); + } + } + + // Reschedule for next items + scheduleQueueProcessing(); + } + } } diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/FakePublisherServiceImpl.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/FakePublisherServiceImpl.java index 9ab1dec73471..9424018b382f 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/FakePublisherServiceImpl.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/FakePublisherServiceImpl.java @@ -81,7 +81,6 @@ public String toString() { @Override public void publish( PublishRequest request, final StreamObserver responseObserver) { - requests.add(request); Response response; try { if (autoPublishResponse) { @@ -98,6 +97,7 @@ public void publish( } if (responseDelay == Duration.ZERO) { sendResponse(response, responseObserver); + requests.add(request); } else { final Response responseToSend = response; executor.schedule( @@ -109,6 +109,7 @@ public void run() { }, responseDelay.toMillis(), TimeUnit.MILLISECONDS); + requests.add(request); } } diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index 95b19d6b0108..2da52e27c7a2 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -105,6 +105,7 @@ public void setUp() throws Exception { testServer.start(); fakeExecutor = new FakeScheduledExecutorService(); + testPublisherServiceImpl.setExecutor(fakeExecutor); } @After @@ -1362,6 +1363,177 @@ public void testPublisherWithoutHedgeSettings() throws Exception { shutdownTestPublisher(publisher); } + private Publisher getPublisherWithHedge(Duration delay) throws Exception { + HedgeSettings hedgeSettings = HedgeSettings.newBuilder() + .setHedgeDelay(delay) + .build(); + return getTestPublisherBuilder() + .setHedgeSettings(hedgeSettings) + .setClock(fakeExecutor.getClock()) + .setBatchingSettings( + Publisher.Builder.DEFAULT_BATCHING_SETTINGS.toBuilder() + .setElementCountThreshold(1L) + .build()) + .build(); + } + + private void waitForRequests(FakePublisherServiceImpl service, int expectedCount) + throws InterruptedException { + long timeout = System.currentTimeMillis() + 5000; + while (service.getCapturedRequests().size() < expectedCount && System.currentTimeMillis() < timeout) { + Thread.sleep(5); + } + if (service.getCapturedRequests().size() < expectedCount) { + throw new AssertionError( + String.format( + "Timed out waiting for requests. Expected: %d, Got: %d", + expectedCount, service.getCapturedRequests().size())); + } + } + + @Test + public void testHedgingNotTriggeredIfFast() throws Exception { + Publisher publisher = getPublisherWithHedge(Duration.ofMillis(50)); + + // Prepare fast response (10ms delay) + testPublisherServiceImpl.setAutoPublishResponse(false); + testPublisherServiceImpl.setPublishResponseDelay(Duration.ofMillis(10)); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("1")); + + ApiFuture future = sendTestMessage(publisher, "msg-fast"); + waitForRequests(testPublisherServiceImpl, 1); + + // Advance time past response but before hedge delay (e.g. 20ms) + fakeExecutor.advanceTime(Duration.ofMillis(20)); + + // Future should be completed + assertEquals("1", future.get()); + + // Only 1 request should be received by server + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(1); + + shutdownTestPublisher(publisher); + } + + @Test + public void testHedgingTriggeredIfSlow() throws Exception { + Publisher publisher = getPublisherWithHedge(Duration.ofMillis(50)); + + // Set response delay to 100ms (greater than 50ms hedge delay) + testPublisherServiceImpl.setAutoPublishResponse(false); + testPublisherServiceImpl.setPublishResponseDelay(Duration.ofMillis(100)); + // Add two responses (one for main, one for hedge) + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("1")); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("2")); + + ApiFuture future = sendTestMessage(publisher, "msg-slow"); + waitForRequests(testPublisherServiceImpl, 1); + + // Advance time to 40ms (before hedge delay) + fakeExecutor.advanceTime(Duration.ofMillis(40)); + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(1); // Only original sent + + // Advance time to 60ms (past 50ms hedge delay) + fakeExecutor.advanceTime(Duration.ofMillis(20)); + waitForRequests(testPublisherServiceImpl, 2); + + // Now attempt 2 should have been triggered + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(2); + + // Advance to 110ms to let responses complete + fakeExecutor.advanceTime(Duration.ofMillis(50)); + assertTrue(future.isDone()); + + shutdownTestPublisher(publisher); + } + + @Test + public void testMultipleHedging() throws Exception { + Publisher publisher = getPublisherWithHedge(Duration.ofMillis(50)); + + // Set delay to 200ms + testPublisherServiceImpl.setAutoPublishResponse(false); + testPublisherServiceImpl.setPublishResponseDelay(Duration.ofMillis(200)); + // Add responses for 3 attempts + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("1")); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("2")); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("3")); + + ApiFuture future = sendTestMessage(publisher, "msg-very-slow"); + waitForRequests(testPublisherServiceImpl, 1); + + // T=0: Attempt 1 sent. + // T=60 (Hedge 1): Attempt 2 sent. + fakeExecutor.advanceTime(Duration.ofMillis(60)); + waitForRequests(testPublisherServiceImpl, 2); + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(2); + + // T=120 (Hedge 2): Attempt 3 sent. + fakeExecutor.advanceTime(Duration.ofMillis(60)); + waitForRequests(testPublisherServiceImpl, 3); + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(3); + + // Advance to complete + fakeExecutor.advanceTime(Duration.ofMillis(100)); + assertEquals("1", future.get(5, TimeUnit.SECONDS)); + + shutdownTestPublisher(publisher); + } + + @Test + public void testHedgingBypassedIfNoTokens() throws Exception { + Publisher publisher = getPublisherWithHedge(Duration.ofMillis(50)); + + // Drain the token bucket completely (since it starts full) + while (publisher.getHedgeTokenBucket().tryAcquire()) {} + assertThat(publisher.getHedgeTokenBucket().getTokens()).isEqualTo(0.0f); + + testPublisherServiceImpl.setPublishResponseDelay(Duration.ofMillis(100)); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("1")); + + ApiFuture future = sendTestMessage(publisher, "msg-slow-no-tokens"); + waitForRequests(testPublisherServiceImpl, 1); + + // Advance past hedge delay + fakeExecutor.advanceTime(Duration.ofMillis(60)); + + // Should NOT trigger hedge because token bucket is empty + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(1); + + fakeExecutor.advanceTime(Duration.ofMillis(50)); + assertEquals("1", future.get(5, TimeUnit.SECONDS)); + + shutdownTestPublisher(publisher); + } + + + + @Test + public void testHedgingCancellationPropagates() throws Exception { + Publisher publisher = getPublisherWithHedge(Duration.ofMillis(50)); + + testPublisherServiceImpl.setAutoPublishResponse(false); + testPublisherServiceImpl.setPublishResponseDelay(Duration.ofMillis(100)); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("1")); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("2")); + + ApiFuture future = sendTestMessage(publisher, "msg-cancel"); + waitForRequests(testPublisherServiceImpl, 1); + + // Trigger hedge + fakeExecutor.advanceTime(Duration.ofMillis(60)); + waitForRequests(testPublisherServiceImpl, 2); + assertThat(testPublisherServiceImpl.getCapturedRequests()).hasSize(2); + + // Cancel the future + future.cancel(true); + + // Verify cancellation propagates to overall future + assertTrue(future.isCancelled()); + + shutdownTestPublisher(publisher); + } + private Builder getTestPublisherBuilder() { return Publisher.newBuilder(TEST_TOPIC) .setExecutorProvider(FixedExecutorProvider.create(fakeExecutor)) From 23b84322f8e958cf5882e3a4bef5e3ddac59d641 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 20 Jul 2026 19:30:45 +0000 Subject: [PATCH 10/11] fix lint --- .../cloud/pubsub/v1/CancellationSharer.java | 12 +++--- .../google/cloud/pubsub/v1/HedgeSettings.java | 3 +- .../google/cloud/pubsub/v1/HedgedRequest.java | 4 +- .../com/google/cloud/pubsub/v1/Publisher.java | 42 +++++++++++-------- .../cloud/pubsub/v1/PublisherImplTest.java | 9 ++-- 5 files changed, 35 insertions(+), 35 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java index 37de9324eedf..702da3667b9a 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java @@ -28,24 +28,23 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; - - /** * Coordinates multiple publish attempts for a single batch of messages. * *

    Implements {@link ApiFuture} to act as the single future returned to the publisher's client. - * It manages the lifecycle of the original attempt and any subsequent hedged attempts (CancellationSharer in the diagram). + * It manages the lifecycle of the original attempt and any subsequent hedged attempts + * (CancellationSharer in the diagram). */ class CancellationSharer extends AbstractApiFuture { private final Publisher.OutstandingBatch batch; private final Publisher publisher; - private final Map> runningAttempts = new ConcurrentHashMap<>(); + private final Map> runningAttempts = + new ConcurrentHashMap<>(); private final AtomicInteger totalAttemptsCount = new AtomicInteger(0); private final AtomicBoolean done = new AtomicBoolean(false); final AtomicBoolean isInQueue = new AtomicBoolean(false); private final AtomicReference lastError = new AtomicReference<>(); - CancellationSharer(Publisher.OutstandingBatch batch, Publisher publisher) { this.batch = batch; this.publisher = publisher; @@ -108,7 +107,8 @@ void checkCompletionOnQueueExit() { if (!done.get() && runningAttempts.isEmpty() && !isInQueue.get()) { if (done.compareAndSet(false, true)) { Throwable error = lastError.get(); - setException(error != null ? error : new RuntimeException("Hedging failed with no active attempts")); + setException( + error != null ? error : new RuntimeException("Hedging failed with no active attempts")); } } } diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java index 62e16aa71639..fdbfebbe9106 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgeSettings.java @@ -82,8 +82,7 @@ public static final class Builder { /** Refill rate. */ private float refill = DEFAULT_REFILL; - private Builder() { - } + private Builder() {} /** * Allows hedging delay to be configurable. diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java index ac405032c8e1..ad07de9a8a90 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/HedgedRequest.java @@ -16,9 +16,7 @@ package com.google.cloud.pubsub.v1; -/** - * Represents a pending hedging check in the publisher's queue. - */ +/** Represents a pending hedging check in the publisher's queue. */ class HedgedRequest { private final CancellationSharer coordinator; private final int attemptNumber; diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index b57442d5375f..35ab918e5645 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -67,15 +67,14 @@ import java.util.List; import java.util.Map; import java.util.concurrent.Callable; +import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; -import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.logging.Level; @@ -1250,16 +1249,18 @@ private ApiFuture startHedgedCall(final OutstandingBatch outsta final AtomicInteger cancelledCount = new AtomicInteger(0); final int batchSize = outstandingBatch.outstandingPublishes.size(); for (final OutstandingPublish outstanding : outstandingBatch.outstandingPublishes) { - outstanding.publishResult.addListener(new Runnable() { - @Override - public void run() { - if (outstanding.publishResult.isCancelled()) { - if (cancelledCount.incrementAndGet() == batchSize) { - coordinator.cancel(true); + outstanding.publishResult.addListener( + new Runnable() { + @Override + public void run() { + if (outstanding.publishResult.isCancelled()) { + if (cancelledCount.incrementAndGet() == batchSize) { + coordinator.cancel(true); + } + } } - } - } - }, directExecutor()); + }, + directExecutor()); } ApiFuture firstAttemptFuture = publishCall(outstandingBatch); @@ -1284,12 +1285,16 @@ private void scheduleQueueProcessing() { long delay = nextItem.getSendAfterMs() - clock.millisTime(); delay = Math.max(0, delay); - queueProcessingFuture = executor.schedule(new Runnable() { - @Override - public void run() { - processQueue(); - } - }, delay, TimeUnit.MILLISECONDS); + queueProcessingFuture = + executor.schedule( + new Runnable() { + @Override + public void run() { + processQueue(); + } + }, + delay, + TimeUnit.MILLISECONDS); } } @@ -1311,7 +1316,8 @@ private void processQueue() { if (hedgeTokenBucket.tryAcquire()) { // Clone and schedule next attempt check (Attempt + 1) long delayMs = hedgeSettings.getHedgeDelay().toMillis(); - HedgedRequest nextItem = new HedgedRequest(coordinator, item.getAttemptNumber() + 1, now + delayMs); + HedgedRequest nextItem = + new HedgedRequest(coordinator, item.getAttemptNumber() + 1, now + delayMs); hedgingQueue.add(nextItem); // Start Hedged Attempt diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index 2da52e27c7a2..6acfd76a5870 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -1364,9 +1364,7 @@ public void testPublisherWithoutHedgeSettings() throws Exception { } private Publisher getPublisherWithHedge(Duration delay) throws Exception { - HedgeSettings hedgeSettings = HedgeSettings.newBuilder() - .setHedgeDelay(delay) - .build(); + HedgeSettings hedgeSettings = HedgeSettings.newBuilder().setHedgeDelay(delay).build(); return getTestPublisherBuilder() .setHedgeSettings(hedgeSettings) .setClock(fakeExecutor.getClock()) @@ -1380,7 +1378,8 @@ private Publisher getPublisherWithHedge(Duration delay) throws Exception { private void waitForRequests(FakePublisherServiceImpl service, int expectedCount) throws InterruptedException { long timeout = System.currentTimeMillis() + 5000; - while (service.getCapturedRequests().size() < expectedCount && System.currentTimeMillis() < timeout) { + while (service.getCapturedRequests().size() < expectedCount + && System.currentTimeMillis() < timeout) { Thread.sleep(5); } if (service.getCapturedRequests().size() < expectedCount) { @@ -1506,8 +1505,6 @@ public void testHedgingBypassedIfNoTokens() throws Exception { shutdownTestPublisher(publisher); } - - @Test public void testHedgingCancellationPropagates() throws Exception { Publisher publisher = getPublisherWithHedge(Duration.ofMillis(50)); From 3a23d607392c87ea301b87ad188a9683fc227ef0 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Mon, 20 Jul 2026 21:14:16 +0000 Subject: [PATCH 11/11] Modifying comment line --- .../java/com/google/cloud/pubsub/v1/CancellationSharer.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java index 702da3667b9a..571c8a99d8ac 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/CancellationSharer.java @@ -32,8 +32,7 @@ * Coordinates multiple publish attempts for a single batch of messages. * *

    Implements {@link ApiFuture} to act as the single future returned to the publisher's client. - * It manages the lifecycle of the original attempt and any subsequent hedged attempts - * (CancellationSharer in the diagram). + * It manages the lifecycle of the original attempt and any subsequent hedged attempts. */ class CancellationSharer extends AbstractApiFuture { private final Publisher.OutstandingBatch batch;