Skip to content

feat(pubsub): add publish telemetry headers for publish attempt observability - #14338

Merged
tonyyyycui merged 3 commits into
googleapis:mainfrom
tonyyyycui:publish-telemetry-header
Sep 16, 2026
Merged

tonyyyycui merged 3 commits into
googleapis:mainfrom
tonyyyycui:publish-telemetry-header

Conversation

@tonyyyycui

@tonyyyycui tonyyyycui commented Sep 9, 2026

Copy link
Copy Markdown
Contributor

Adds client-side publish telemetry headers to provide observability into hedging attempts.

@tonyyyycui
tonyyyycui requested review from a team as code owners September 9, 2026 20:29

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces publish hedging for the Google Cloud Pub/Sub Java client, adding configuration settings (HedgingSettings), a coordinator (CancellationSharer), and queue management (HedgedRequest) to manage and rate-limit hedged attempts. Feedback on these changes highlights several critical issues, including a race condition in scheduleQueueProcessing that could leave tasks unscheduled, a logic error in CancellationSharer where a failure of the original attempt prematurely cancels active hedged attempts, a thread-safety issue with queueProcessingFuture lacking volatile visibility, and a hardcoded 10-second timeout cap that ignores the user's configured maximum RPC timeout.

Comment on lines +757 to +778
private void scheduleQueueProcessing() {
if (isQueueProcessingScheduled.compareAndSet(false, true)) {
HedgedRequest nextItem = hedgingQueue.peek();
if (nextItem == null) {
isQueueProcessingScheduled.set(false);
return;
}

long delay = Math.max(0, nextItem.getSendAfterMs() - clock.millisTime());

queueProcessingFuture =
executor.schedule(
new Runnable() {
@Override
public void run() {
processQueue();
}
},
delay,
TimeUnit.MILLISECONDS);
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

There is a race condition in scheduleQueueProcessing when scheduling a new hedged request. If hedgingQueue.peek() returns null (because the queue was empty at that moment), but another thread adds an item to the queue and calls scheduleQueueProcessing right before isQueueProcessingScheduled.set(false) is executed, the compareAndSet in the second thread will fail (since it is still true). The first thread will then set the flag to false and return. As a result, the newly added item will remain in the queue without any scheduled task to process it. Using a loop to re-check the queue after setting the flag to false resolves this issue.

  private void scheduleQueueProcessing() {
    while (true) {
      if (!isQueueProcessingScheduled.compareAndSet(false, true)) {
        return;
      }
      HedgedRequest nextItem = hedgingQueue.peek();
      if (nextItem == null) {
        isQueueProcessingScheduled.set(false);
        if (!hedgingQueue.isEmpty()) {
          continue;
        }
        return;
      }

      long delay = Math.max(0, nextItem.getSendAfterMs() - clock.millisTime());

      queueProcessingFuture =
          executor.schedule(
              new Runnable() {
                @Override
                public void run() {
                  processQueue();
                }
              },
              delay,
              TimeUnit.MILLISECONDS);
      return;
    }
  }


private final ConcurrentLinkedQueue<HedgedRequest> hedgingQueue;
private final AtomicBoolean isQueueProcessingScheduled;
private ScheduledFuture<?> queueProcessingFuture;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The queueProcessingFuture field is written to by both the publisher thread (via startHedgedCall) and the executor thread (via processQueue) without any synchronization or volatile visibility. This can lead to data races and visibility issues, particularly during shutdown where queueProcessingFuture is read and cancelled. Marking this field as volatile ensures proper thread visibility.

Suggested change
private ScheduledFuture<?> queueProcessingFuture;
private volatile ScheduledFuture<?> queueProcessingFuture;

@googleapis googleapis deleted a comment from gemini-code-assist Bot Sep 9, 2026
.setPublishOperation(
PubsubClientTelemetry.PublishOperation.newBuilder()
.setHedgedAttemptCount(attemptNumber)
.setPublishStartTime(Timestamps.fromMillis(outstandingBatch.creationTime))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't think this is the right publish time. We should use the publish time from the publish() call per the proto comment: "Time at which the publish() call was initiated in the client library". This will likely require a few more changes to plumb that value to here, so I'd be okay with just having the hedged attempt set for now.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I've removed this from the changes.

@tonyyyycui
tonyyyycui force-pushed the publish-telemetry-header branch from 8eb6694 to 510bf79 Compare September 11, 2026 01:31
@tonyyyycui
tonyyyycui force-pushed the publish-telemetry-header branch from 510bf79 to a926bdb Compare September 14, 2026 19:09
@michaelpri10

Copy link
Copy Markdown
Contributor

Only suggestion is to update the PR description as we don't currently add the publish latency to the header and this has been rebased to HEAD, so isn't stacked anymore.

@tonyyyycui
tonyyyycui merged commit c167ab8 into googleapis:main Sep 16, 2026
210 checks passed
blakeli0 pushed a commit that referenced this pull request Sep 23, 2026
🤖 I have created a release *beep* *boop*
---


<details><summary>1.92.0</summary>

##
[1.92.0](v1.91.0...v1.92.0)
(2026-09-23)


### Features

* **bigquery-jdbc:** add `EnableTimestampPicos` connection property and
its plumbing
([#14284](#14284))
([b4aa5ac](b4aa5ac))
* **bigquery-jdbc:** implement picosecond temporal math and formatting
engine
([#14286](#14286))
([2a9612a](2a9612a))
* **bigquery-jdbc:** support picosecond in REST JSON path and nested
types
([#14334](#14334))
([15ffe4a](15ffe4a))
* **bigquery-jdbc:** support picosecond in `PreparedStatement`
parameters and batching
([#14373](#14373))
([c1aac66](c1aac66))
* **bigquery-jdbc:** support picosecond timestamp in `ResultSetMetaData`
and `DatabaseMetaData`
([#14358](#14358))
([43acdd3](43acdd3))
* **bigquery-jdbc:** support picosecond timestamps in Arrow Storage Read
API and nested types
([#14332](#14332))
([b5d9aca](b5d9aca))
* **bigquery-jdbc:** support qualified project delimiter in
`DefaultDataset` property
([#14240](#14240))
([6e8d6c8](6e8d6c8))
* **bigquery:** accelerate row-based query() with Arrow wire format
([#14405](#14405))
([8d12a8f](8d12a8f))
* **bigquery:** add ArrowDeserializer helper utility
([#13943](#13943))
([d9a298b](d9a298b))
* **bigquery:** add ArrowQueryPageFetcher for Arrow query result
pagination
([#14404](#14404))
([615409f](615409f))
* **bigquery:** add ArrowQueryResult and ArrowQueryResultImpl for Arrow
result streaming
([#13944](#13944))
([a62fdf8](a62fdf8))
* **bigquery:** add Storage Read API slow-path fallback for row-based
query()
([#14409](#14409))
([26e568a](26e568a))
* **bigquery:** add zero-copy queryArrow API for Arrow VectorSchemaRoot
streaming
([#14402](#14402))
([b44ffe8](b44ffe8))
* **bigquery:** make BigQuery AutoCloseable with default no-op close
method
([#14434](#14434))
([00bf3de](00bf3de))
* **firestore:** add support for BSON types
([#13189](#13189))
([8a123d9](8a123d9))
* **gax:** add ApiCallContext and request-level settings overloads to
ResumableUploadCallable
([#14251](#14251))
([e8cbd42](e8cbd42))
* **gax:** add globalTimeout settings field to
ResumableUploadCallSettings
([#14253](#14253))
([438cda6](438cda6))
* **gax:** add resumable upload error classification and retry algorithm
([#14419](#14419))
([b70396d](b70396d))
* **gax:** add ResumableUploadCallable creation to Callables and
HttpJsonCallableFactory
([#14242](#14242))
([7de24de](7de24de))
* **gax:** implement baseline Callable and Future for resumable uploads
([#14241](#14241))
([5a54db9](5a54db9))
* **generator:** add model flag and allowlist parser for resumable
upload RPCs
([#14317](#14317))
([acc1856](acc1856))
* **generator:** emit resumable upload client surface
([#14319](#14319))
([a9fed00](a9fed00))
* **generator:** emit resumable upload settings and HttpJson upload stub
([#14321](#14321))
([c122474](c122474))
* **generator:** enable resumable upload generation for showcase
([#14325](#14325))
([f9ebd79](f9ebd79))
* **generator:** switch resumable upload specialized stubs to package
private
([#14471](#14471))
([0d4e875](0d4e875))
* **generator:** wire transport stub delegation to resumable upload
stubs
([#14322](#14322))
([cc4b980](cc4b980))
* **google/cloud/backupdr/v1beta:** add backupdr
([#14410](#14410))
([a4a47da](a4a47da))
* **google/cloud/networkservices/v1beta1:** add networkservices
([#14407](#14407))
([21c4955](21c4955))
* **pubsub:** add publish telemetry headers for publish attempt
observability
([#14338](#14338))
([c167ab8](c167ab8))
* **pubsub:** implement publish hedging to reduce tail latency
([#13735](#13735))
([b302615](b302615))
* **spanner:** Support dynamic TLS certificate and key rotation for
Spanner Omni
([#14456](#14456))
([ffc745c](ffc745c))
* **storage/control:** add delete folder recursive sample
([#13642](#13642))
([f4b1b46](f4b1b46))
* **storage/control:** add delete folder recursive sample
([#14397](#14397))
([2c01d55](2c01d55))


### Bug Fixes

* **auth:** restore transportFactory upon deserialization in
InternalAwsSecurityCredentialsSupplier
([#14340](#14340))
([beea42f](beea42f))
* **bigquery-jdbc:** ensure row ordering in PCNT IT
([#14330](#14330))
([a16f048](a16f048))
* **bigquery-jdbc:** fix htapi fallback due to permission logic
([#14418](#14418))
([21e6dc8](21e6dc8))
* **bigquery-jdbc:** fix Timestamp assertions
([#14290](#14290))
([533ba14](533ba14))
* **bigquery-jdbc:** handle null parameters in Storage Write API bulk
inserts
([#14270](#14270))
([dd2c41a](dd2c41a)),
refs
[#14066](#14066)
* **bigquery-jdbc:** handle SQL NULLs in ResultSet primitive getters
([#14383](#14383))
([8e464fe](8e464fe)),
refs
[#14371](#14371)
* **bigquery:** default Arrow pagination stream location to US instead
of global
([#14458](#14458))
([2775eb1](2775eb1))
* **bigquery:** preserve page token and paginate correctly in Arrow
query when maxResults is set
([#14469](#14469))
([f5601f4](f5601f4))
* **bigquery:** use first page row count for Arrow query pagination
offset
([#14466](#14466))
([9d10dd0](9d10dd0))
* **bigtable:** don't notify config listeners while holding the manager
lock
([#14294](#14294))
([4426ccd](4426ccd))
* **bigtable:** fall back to classic path when per-RPC CallCredentials
are set on session path
([#14477](#14477))
([57bacb0](57bacb0))
* **bigtable:** fix abnormal session closures and scale-up in session
pool
([#14431](#14431))
([6361ecd](6361ecd))
* **biqguery:** fix undeclared QueryParameter wiring in QueryStatistics
([#14401](#14401))
([64cf1d3](64cf1d3))
* **bom:** restore google-cloud-spanner-jdbc to libraries-bom
([#14362](#14362))
([bc7be5e](bc7be5e)),
refs
[#14347](#14347)
* **spanner:** honor maxAttempts and totalTimeout in streaming resume
loop
([#14370](#14370))
([305f47d](305f47d))
* **spanner:** only set snapshot isolation read timestamp for SI or
optimistic txns in CloudClientExecutor
([#14346](#14346))
([54c0d0f](54c0d0f))
* **spanner:** prevent statement cancellation race in
AbstractBaseUnitOfWork
([#14283](#14283))
([d9a8eef](d9a8eef))
* **spanner:** re-enable ITInstanceAdminTest on cloud-devel and
cloud-staging
([#14281](#14281))
([89a8268](89a8268))


### Performance Improvements

* **spanner:** stop re-parsing the request id on every RPC
([#14353](#14353))
([46108f4](46108f4))


### Documentation

* Add a Http/Json Post-Quantum Cryptography Guide
([#13963](#13963))
([fcc65b0](fcc65b0))
* **bigquery:** add QueryArrow code sample and document JDK 17+ JVM
requirements
([#14437](#14437))
([bd363f6](bd363f6))
</details>

---
This PR was generated with [Release
Please](https://github.com/googleapis/release-please). See
[documentation](https://github.com/googleapis/release-please#release-please).

---------

Co-authored-by: release-please[bot] <55107282+release-please[bot]@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants