Skip to content

fix: broker HTTP client connection-reuse race and response leaks - #20475

Open
kgyrtkirk wants to merge 3 commits into
apache:masterfrom
kgyrtkirk:response-stream-apache
Open

kgyrtkirk wants to merge 3 commits into
apache:masterfrom
kgyrtkirk:response-stream-apache

Conversation

@kgyrtkirk

Copy link
Copy Markdown
Member

🩸 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:

  • 🔥 Connection reuse race. When a query's result stream is closed, the client can close a connection that has already been returned to the pool and handed to the next request. That request fails with Failed to write request to channel.
  • 🧟 Response lifecycle leaks. The Coordinator's segment-metadata response (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 a sys.segments scan 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:

WARN ResourcePool - Resource at key[https://host:8283] was returned multiple times?
WARN JsonParserIterator - Query [...] to host [host:8283] interrupted
io.netty.channel.ChannelException: [POST https://host:8283/druid/v2/] Failed to write request to channel
Caused by: io.netty.handler.ssl.SslClosedEngineException: SSLEngine closed already

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 stream close() calls TrafficCop.abort(), which is channel.close(). JsonParserIterator closes the stream as soon as it reads END_ARRAY, which is usually before the event loop has processed the response's LastHttpContent.

  • Consumer thread: close() sees done == false and calls abort(). channel.close() is called off the event loop, so it is only queued.
  • Event loop: LastHttpContent → finishRequest() → the channel returns to the pool.
  • Next request: takes the channel. isGood() still passes because the close hasn't run yet, so its write is queued behind the close.
  • Event loop again: the close runs, then the write fails. The returned multiple times? warning is the victim's own double return: once from the write listener and once from channelInactive.

The done check in close() narrows the window but does not close it, because the check and abort() are not atomic with respect to finishRequest(). The race was introduced with abort-on-close (#19607). It is much more likely with druid.*.http.poolImplementation=adaptive (#20273): under light load that pool shrinks to demand, so the channel just returned is the next one handed out. retaining keeps numConnections per 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, TrafficCop calls 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: stream close() always discards buffered chunks and calls abort(), 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 AppendableByteArrayInputStream instances. Five /metadata/segments responses 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 reported available = -1009272128.

Causes and changes

Defect Change
AppendableByteArrayInputStream.exceptionCaught keeps the queued bytes, which are unreachable once the stream has failed. Clears the queue and the available count. Chunks arriving after completion are dropped.
AppendableByteArrayInputStream.close() is the InputStream no-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.
available is an int, so above 2 GiB it wraps and available() goes negative. It is now a long, and available() saturates at Integer.MAX_VALUE. Bytes pulled concurrently with a failure no longer drive it negative.
CoordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus returns a CloseableIterator over a JsonParserIterator whose underlying InputStream is only closed when the parser reaches END_ARRAY. It returns a 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. It closes in 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 plain Iterator, and SystemSchema's sys.segments scan never closes it (LIMIT, errors, cancellation). getSegments() returns a CloseableIterator, including when filtered. The sys.segments scan fetches lazily on first read and closes the response when the enumerator closes.
QueryHandler never closes the Calcite Interpreter returned by BindableRel.bind(), so system-table resources (such as the response above) are not released. The bindable path creates the enumerator lazily and closes it in cleanup. It also closes the bound Enumerable when it is AutoCloseable, including when enumeration fails.

⚙️ Coordinator client configuration

BrokerViewOfBrokerConfig and BrokerViewOfCoordinatorConfig each built a private Coordinator client with 15 attempts. They now use the injected CoordinatorClient, so there is one client and one retry policy for Coordinator calls. The attempt count comes from the new druid.client.coordinator.maxAttempts, default 15, so Broker startup keeps its previous retry budget (about 200 s of backoff). Other Coordinator calls go from 6 attempts to 15. Values below 1 are rejected at startup; previously a negative value would have meant unlimited retries. The new CoordinatorClientConfig gives future Coordinator-client settings a home.

🧪 Tests

  • FriendlyServersTest: abort and resume after completion leave the pooled channel and the next request untouched.
  • AppendableByteArrayInputStreamTest: saturation of available, 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 when sys.segments is closed early. The fetch is deferred until rows are read. The interpreter is closed on success and on failure.

🚧 Out of scope / known limitations

  • Back-pressure for the segment-metadata response (a bounded buffer that suspends reads when the consumer is slower than the network) is not included. The 3.3 GB in-flight response above remains possible. Follow-up.
  • If a stream is closed after its last row but before LastHttpContent is processed, the connection is still closed rather than reused. That is safe, but it costs a reconnect.
  • A server-initiated close (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...

cleanup

remove ref

up

u
);
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 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

🟢 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 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

🟢 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)

@cryptoe cryptoe added this to the 39.0.0 milestone Oct 6, 2026
@cryptoe

cryptoe commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

This seems like a blocker for Druid 39.

kgyrtkirk added a commit to kgyrtkirk/druid that referenced this pull request Oct 8, 2026
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>
@kgyrtkirk

Copy link
Copy Markdown
Member Author

split up the changes of this into the following pieces:

  • small preparation pr to later have controls over the client buffer - the retries was burned in ; this makes it configurable....so the next PR will just add a new config to it
  • fix the int overflow and other minor stuff
  • the trafficCop could have aborted connection which were already returned - and might be already taken out
  • the close fixes to actually call the close on the buffer

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants