fix(spanner): honor maxAttempts and totalTimeout in streaming resume loop - #14370
Conversation
There was a problem hiding this comment.
Code Review
This pull request updates the ResumableStreamIterator to honor the configured maximum attempts and total timeout settings during streaming retries, preventing potential infinite retry loops. It also introduces comprehensive tests to verify this behavior. The review feedback is highly constructive, pointing out a potential sentinel value collision with System.nanoTime(), a numeric overflow risk in the timeout calculation, and a thread leak in the test suite due to an unclosed ScheduledThreadPoolExecutor.
- Normalize generated streaming retry defaults at the RPC boundary. - Extract retry handling and replace reflection with deterministic tests. - Document and test maxAttempts=1 as the way to disable streaming retries. - Cover StreamingRead attempt and timeout budgets with mock-server tests.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces bounded streaming retries for Spanner queries and reads by honoring RetrySettings (specifically maxAttempts and totalTimeout) in ResumableStreamIterator. It resets the retry budget when progress is made on the stream (i.e., when a new resume token is received). The default settings are normalized to preserve the historical unbounded resume behavior. Comprehensive unit and integration tests are added to verify these limits. Feedback on the changes highlights a potential issue where nextBackOffMillis returning -1L (indicating backoff has stopped) is not handled, which could lead to a negative sleep duration and an IllegalArgumentException instead of propagating the original SpannerException.
| long delayMillis = spannerException.getRetryDelayInMillis(); | ||
| if (delayMillis == -1L) { | ||
| if (this.backOff == null) { | ||
| this.backOff = newBackOff(); | ||
| } | ||
| delayMillis = nextBackOffMillis(this.backOff); | ||
| } |
There was a problem hiding this comment.
If nextBackOffMillis(this.backOff) returns BackOff.STOP (-1L), the backoff has stopped. Currently, if it returns -1L, totalTimeoutExceeded will receive -1L and return false (since Math.max(-1L, 0L) is 0L, which is less than the remaining timeout). This results in checkRetryBudgetAndGetDelay returning -1L, which is then passed to backoffSleep(context, -1L). Sleeping for a negative duration can cause an IllegalArgumentException (e.g., from Thread.sleep) instead of propagating the original SpannerException.
We should explicitly check if delayMillis is -1L after calling nextBackOffMillis and throw the original spannerException with the appropriate span annotations.
long delayMillis = spannerException.getRetryDelayInMillis();
if (delayMillis == -1L) {
if (this.backOff == null) {
this.backOff = newBackOff();
}
delayMillis = nextBackOffMillis(this.backOff);
if (delayMillis == -1L) {
span.addAnnotation(
"Stream broken. Not retrying because the backoff has stopped",
spannerException);
span.setStatus(spannerException);
throw spannerException;
}
}There was a problem hiding this comment.
I don't think this change is needed - the -1L/BackOff.STOP case is already handled safely, and the predicted
IllegalArgumentException can't actually occur here
🤖 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>
The resume loop in ResumableStreamIterator restarted a broken stream indefinitely for any retryable error, ignoring the maxAttempts and totalTimeout configured in the retry settings for ExecuteStreamingSql.
A streaming query could therefore retry forever when the server kept returning a retryable error, for example when a user configured DEADLINE_EXCEEDED as a retryable code and every attempt timed out.
The loop now counts consecutive failed attempts and stops retrying, rethrowing the last exception, when the configured maxAttempts is reached. It also enforces the configured totalTimeout as a wall-clock budget for a sequence of consecutive failed attempts, measured from the first failure of the sequence: a retry is only allowed when the
retry delay still fits in the remaining budget. This applies both to delays from the exponential backoff and to server-supplied retry delays (RetryInfo), which previously bypassed the backoff completely.
A totalTimeout of zero means that no time budget has been set, in which case only maxAttempts limits the retries, mirroring GAX.
Both limits only bind for custom retry settings: the default streaming retry settings do not set maxAttempts and keep the existing unbounded resume behavior. Progress on the stream resets both budgets, where progress means receiving a resume token that differs from the last seen token, so long-running streams that regularly make progress are
not terminated by an occasional transient error, while a stream that keeps returning the same token cannot reset the budget indefinitely.