| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
Sorry, something went wrong.
Codecov Report❌ Patch coverage is 87.17949% with 5 lines in your changes missing coverage. Please review. @@ Coverage Diff @@
## master #720 +/- ##
============================================
- Coverage 73.18% 73.13% -0.06%
- Complexity 3589 3597 +8
============================================
Files 392 392
Lines 16531 16551 +20
Branches 1736 1731 -5
============================================
+ Hits 12099 12104 +5
- Misses 3811 3820 +9
- Partials 621 627 +6 ☔ View full report in Codecov by Harness.
|
Sorry, something went wrong.
TopicRetryableStream.start() is called from the shared transport scheduler on every reconnect. With directWrite enabled, WriteStreamDirectFactory resolved the target partition and its location synchronously inside createNewStream: lookupPartitionId() joined a probe stream future (1 min deadline) and lookupLocation() joined describeTopic() (1 min deadline). Each reconnect of an unresponsive destination could therefore occupy a scheduler thread for up to two minutes. The shared scheduler is sized max(cores / 2, 2) and is also used by discovery, session pools, retry contexts and operation tray, so a handful of stalled writers could stall the whole transport: session acquire timeouts stop firing and discovery ticks stop running. Make createNewStream() return CompletableFuture and compose the partition and location lookups instead of joining them, so no shared scheduler thread is held while a stream is being created. Since stream creation is now asynchronous, close() may happen while it is in progress. TopicRetryableStream handles that by re-checking isClosed after publishing the new stream: close() sets the volatile flag before clearing the stream reference, so a creation that wins the race always observes the flag and drops the stream without starting it. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
| Back | FazBrowse Home | New Git URL |
Problem
TopicRetryableStream.start() runs on the transport's shared scheduler on every reconnect. With directWrite enabled, WriteStreamDirectFactory.createNewStream resolved the target partition and its location synchronously inside that call:
So a single reconnect toward an unresponsive destination could hold a scheduler thread for up to two minutes.
That scheduler is sized max(cores / 2, 2) and is shared with discovery, the table/query session pools, retry contexts and OperationTray. A handful of stalled direct writers is enough to exhaust it head-of-line: session-acquire timeouts stop firing, discovery ticks stop running, and the whole transport degrades — not just the topic writers that caused it.
Change
createNewStream now returns CompletableFuture and the two lookups are composed rather than joined, so no shared scheduler thread is ever held while a stream is being created.
Behavior notes for reviewers
Compatibility
Source-incompatible for anyone subclassing TopicRetryableStream or WriteStreamFactory: createNewStream changed its return type, from S to CompletableFuture<S>. Both are internal impl classes, and inside the repo the only subclass is WriteSession. Should go into a minor release, not a patch.