feat(pubsub): add publish telemetry headers for publish attempt observability - #14338
Conversation
There was a problem hiding this comment.
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.
| 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); | ||
| } | ||
| } |
There was a problem hiding this comment.
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; |
There was a problem hiding this comment.
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.
| private ScheduledFuture<?> queueProcessingFuture; | |
| private volatile ScheduledFuture<?> queueProcessingFuture; |
| .setPublishOperation( | ||
| PubsubClientTelemetry.PublishOperation.newBuilder() | ||
| .setHedgedAttemptCount(attemptNumber) | ||
| .setPublishStartTime(Timestamps.fromMillis(outstandingBatch.creationTime)) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
I've removed this from the changes.
8eb6694 to
510bf79
Compare
510bf79 to
a926bdb
Compare
|
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. |
🤖 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>
Adds client-side publish telemetry headers to provide observability into hedging attempts.