Repository navigation
Conversation
| ); | ||
| final Request request = new Request( | ||
| HttpMethod.GET, | ||
| new URL(StringUtils.format("http://localhost:%d/", serverSocket.getLocalPort())) |
| ); | ||
| final Request request = new Request( | ||
| HttpMethod.GET, | ||
| new URL(StringUtils.format("http://localhost:%d/", serverSocket.getLocalPort())) |
|
|
||
| Assertions.assertEquals(Integer.MAX_VALUE, in.available()); | ||
|
|
||
| in.read(new byte[oneMebibyte.length]); |
| Assertions.assertEquals(0, in.available()); | ||
| Assertions.assertEquals(5, in.read(new byte[5]), "the chunk being read is still handed out"); | ||
| Assertions.assertEquals(0, in.available()); | ||
| Assertions.assertThrows(IOException.class, () -> in.read(new byte[8192])); |
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
No actionable issues found in the 25-file current-head review. The response-lifecycle changes preserve event-loop ordering for pooled channel returns, release Coordinator metadata streams on normal and exceptional or early exits, and close bindable Calcite resources without changing the intended retry-policy behavior.
Reviewed 25 of 25 changed files across processing HTTP transport and stream handlers, Coordinator client and configuration, broker and SQL lifecycle paths, tests, benchmark, and documentation.
Validation: git diff --check 3e9365196f9893ff9b1cefcfbcb57d49e4080b4a a903d6fe07464e0fe820d49fab3b8a22bd8d96d3 passed. No builds or test suites were run per the scoped static-review instruction.
This is an automated review by Codex GPT-5.6-Luna(max)
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
No actionable issues found in this updated-since-review pass. The incremental change only adds retryable to the website spelling dictionary; I also rechecked the complete current 26-file diff and surrounding HTTP connection lifecycle, streamed response cleanup, SQL enumerator/interpreter cleanup, Coordinator client retry configuration, tests, benchmark, and documentation. The prior review had no findings or open review threads, and no current behavior requires a follow-up fix.
Reviewed 26 of 26 changed files.
Validation: git diff --check 3e9365196f9893ff9b1cefcfbcb57d49e4080b4a..3b56065c161bf8ca63440d8d960314d17a6623e7 passed. No builds or test suites were run per the scoped static-review instruction.
This is an automated review by Codex GPT-5.6-Luna(max)
|
This seems like a blocker for Druid 39. |
Add class javadoc to CoordinatorClientConfig; drop redundant prefix test and rename retry-policy test in ServiceClientModuleTest. Fewer tests, same faith. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
split up the changes of this into the following pieces:
|
🩸 Broker HTTP client: pooled connections closed under the next request, and Coordinator segment-metadata responses retained in heap
📜 Summary
Two families of Broker-side defects in the internal HTTP client and its callers:
Failed to write request to channel.GET /druid/coordinator/v1/metadata/segments) is buffered in heap, and on several paths it is never released: after the request fails, after the reader stops early, or when asys.segmentsscan ends before consuming all rows. On large clusters this response is several GB. One Broker heap dump held about 6 GB in five such responses.🔥 Abort races the return of the connection to the pool
Symptom
Intermittent query failures from Broker to Historical/Peon, with no trace on the data server:
Over plain HTTP the cause is
StacklessClosedChannelException. In each observed case the failing request started 0–1 ms after another request to the same host completed. The base rate for that over about 50k native queries was 1.3%.Cause
DirectDruidClient's result streamclose()callsTrafficCop.abort(), which ischannel.close().JsonParserIteratorcloses the stream as soon as it readsEND_ARRAY, which is usually before the event loop has processed the response'sLastHttpContent.close()seesdone == falseand callsabort().channel.close()is called off the event loop, so it is only queued.LastHttpContent→finishRequest()→ the channel returns to the pool.isGood()still passes because the close hasn't run yet, so its write is queued behind the close.returned multiple times?warning is the victim's own double return: once from the write listener and once fromchannelInactive.The
donecheck inclose()narrows the window but does not close it, because the check andabort()are not atomic with respect tofinishRequest(). The race was introduced with abort-on-close (#19607). It is much more likely withdruid.*.http.poolImplementation=adaptive(#20273): under light load that pool shrinks to demand, so the channel just returned is the next one handed out.retainingkeepsnumConnectionsper host in FIFO order and exposes it mainly when the pool is saturated.A late
TrafficCop.resume()from a completed response has the same flaw: it can re-enable reads on a channel that now belongs to the next request.Required behaviour
Once a response has completed and its connection has been returned to the pool,
TrafficCopcalls for that response have no effect.Change
NettyHttpClient:finishRequest()marks the request as returned to the pool, under the watermark lock, before it returns the connection.abort()runs on the channel's event loop and closes only if the request has not yet been returned.resume()is a no-op once it has.DirectDruidClient: streamclose()always discards buffered chunks and callsabort(), which is now safe in every state.HttpResponseHandler.TrafficCop: documents that calls after completion have no effect.🧟 Coordinator segment-metadata response not released
Symptom
A Broker heap dump (50 GB heap) showed 6.2 GB held by
AppendableByteArrayInputStreaminstances. Five/metadata/segmentsresponses accounted for nearly all of it. One was still being parsed (3.3 GB). Four were already complete or failed, with no reader, and were pinned by their pooled channel's pipeline. One of them reportedavailable = -1009272128.Causes and changes
AppendableByteArrayInputStream.exceptionCaughtkeeps the queued bytes, which are unreachable once the stream has failed.availablecount. Chunks arriving after completion are dropped.AppendableByteArrayInputStream.close()is theInputStreamno-op, so a reader that walks away leaves the whole response buffered while the rest keeps arriving.close()fails the stream: it discards buffered and later chunks, wakes blocked readers, and makes later reads fail.availableis anint, so above 2 GiB it wraps andavailable()goes negative.long, andavailable()saturates atInteger.MAX_VALUE. Bytes pulled concurrently with a failure no longer drive it negative.CoordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatusreturns aCloseableIteratorover aJsonParserIteratorwhose underlyingInputStreamis only closed when the parser reachesEND_ARRAY.ResourceHolder<Iterator<SegmentStatusInCluster>>that owns the response. Closing it releases the stream even if iteration never started.MetadataSegmentView.poll()never closes what it drains, so a failure mid-iteration leaves the response held.finally. A close failure after a complete read is logged, and the poll's segments are kept.MetadataSegmentView.getSegments()(cache disabled) hands the live response out as a plainIterator, andSystemSchema'ssys.segmentsscan never closes it (LIMIT, errors, cancellation).getSegments()returns aCloseableIterator, including when filtered. Thesys.segmentsscan fetches lazily on first read and closes the response when the enumerator closes.QueryHandlernever closes the CalciteInterpreterreturned byBindableRel.bind(), so system-table resources (such as the response above) are not released.Enumerablewhen it isAutoCloseable, including when enumeration fails.⚙️ Coordinator client configuration
BrokerViewOfBrokerConfigandBrokerViewOfCoordinatorConfigeach built a private Coordinator client with 15 attempts. They now use the injectedCoordinatorClient, so there is one client and one retry policy for Coordinator calls. The attempt count comes from the newdruid.client.coordinator.maxAttempts, default15, so Broker startup keeps its previous retry budget (about 200 s of backoff). Other Coordinator calls go from 6 attempts to 15. Values below1are rejected at startup; previously a negative value would have meant unlimited retries. The newCoordinatorClientConfiggives future Coordinator-client settings a home.🧪 Tests
FriendlyServersTest: abort and resume after completion leave the pooled channel and the next request untouched.AppendableByteArrayInputStreamTest: saturation ofavailable, release on failure, close discarding queued and later chunks, reads after close, and close waking all readers.CoordinatorClientImplTest,ServiceClientModuleTest: the holder releases the response when closed before iterating, the config prefix, and the attempt count.MetadataSegmentViewTest,SystemSchemaTest,QueryHandlerTest: the response is released on poll, on iteration failure, on filtered uncached reads, and whensys.segmentsis closed early. The fetch is deferred until rows are read. The interpreter is closed on success and on failure.🚧 Out of scope / known limitations
LastHttpContentis processed, the connection is still closed rather than reused. That is safe, but it costs a reconnect.Connection: close, an idle timeout shorter than the client's) can still hand out a connection that is about to close. That is a separate issue.LLM usage: I did used AI to create this PR - mostly in the tests; and in the above description as well...