From 9131f94503f8849680235be20a00001598063453 Mon Sep 17 00:00:00 2001 From: Renato Haeberli Date: Thu, 23 Jul 2026 16:15:55 +0200 Subject: [PATCH 1/5] SOLR-18312: introduce dedicated thread pool executor for httpClientBuilder --- .../client/solrj/impl/HttpJdkSolrClient.java | 41 +++++++++++++---- .../solrj/impl/HttpJdkSolrClientTest.java | 46 +++++++++++++++++++ 2 files changed, 78 insertions(+), 9 deletions(-) diff --git a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java index 111bfd1bc92..5d53f0b5420 100644 --- a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java +++ b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java @@ -42,6 +42,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.SynchronousQueue; import java.util.concurrent.TimeUnit; import java.util.regex.Matcher; import java.util.regex.Pattern; @@ -78,28 +79,36 @@ public class HttpJdkSolrClient extends HttpSolrClient { protected HttpClient httpClient; + /** + * Executor used to stream (produce) request bodies into the pipe consumed by the JDK HttpClient. + * This is the "producer" side and may be supplied by the caller. + */ protected ExecutorService executor; + /** Dedicated executor handed to the JDK HttpClient */ + protected ExecutorService httpClientExecutor; + private boolean forceHttp11; private final boolean shutdownExecutor; protected HttpJdkSolrClient(String serverBaseUrl, HttpJdkSolrClient.Builder builder) { super(serverBaseUrl, builder); - HttpClient.Builder b = HttpClient.newBuilder(); + HttpClient.Builder httpClientBuilder = HttpClient.newBuilder(); HttpClient.Redirect followRedirects = Boolean.TRUE.equals(builder.getFollowRedirects()) ? HttpClient.Redirect.NORMAL : HttpClient.Redirect.NEVER; - b.followRedirects(followRedirects); + httpClientBuilder.followRedirects(followRedirects); - b.connectTimeout(Duration.of(builder.getConnectionTimeoutMillis(), ChronoUnit.MILLIS)); + httpClientBuilder.connectTimeout( + Duration.of(builder.getConnectionTimeoutMillis(), ChronoUnit.MILLIS)); // note: idle timeout isn't used for the JDK client // note: request timeout is set per request if (builder.sslContext != null) { - b.sslContext(builder.sslContext); + httpClientBuilder.sslContext(builder.sslContext); } if (builder.getExecutor() != null) { @@ -117,15 +126,23 @@ protected HttpJdkSolrClient(String serverBaseUrl, HttpJdkSolrClient.Builder buil new SolrNamedThreadFactory(this.getClass().getSimpleName())); this.shutdownExecutor = true; } - b.executor(this.executor); + this.httpClientExecutor = + new ExecutorUtil.MDCAwareThreadPoolExecutor( + 0, + Integer.MAX_VALUE, + 60, + TimeUnit.SECONDS, + new SynchronousQueue<>(), + new SolrNamedThreadFactory(this.getClass().getSimpleName() + "-http")); + httpClientBuilder.executor(this.httpClientExecutor); if (builder.shouldUseHttp1_1()) { this.forceHttp11 = true; - b.version(HttpClient.Version.HTTP_1_1); + httpClientBuilder.version(HttpClient.Version.HTTP_1_1); } if (builder.cookieHandler != null) { - b.cookieHandler(builder.cookieHandler); + httpClientBuilder.cookieHandler(builder.cookieHandler); } if (builder.getProxyHost() != null) { @@ -133,10 +150,10 @@ protected HttpJdkSolrClient(String serverBaseUrl, HttpJdkSolrClient.Builder buil log.warn( "Socks4 is likely not supported by this client. See https://bugs.openjdk.org/browse/JDK-8214516"); } - b.proxy( + httpClientBuilder.proxy( ProxySelector.of(new InetSocketAddress(builder.getProxyHost(), builder.getProxyPort()))); } - this.httpClient = b.build(); + this.httpClient = httpClientBuilder.build(); assert ObjectReleaseTracker.track(this); } @@ -545,6 +562,12 @@ public void close() throws IOException { } executor = null; + // The http client executor is always created and owned by this instance. + if (httpClientExecutor != null) { + ExecutorUtil.shutdownAndAwaitTermination(httpClientExecutor); + httpClientExecutor = null; + } + assert ObjectReleaseTracker.release(this); } diff --git a/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java b/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java index 79bfdb63e9c..9aef536d370 100644 --- a/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java +++ b/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java @@ -48,6 +48,7 @@ import org.apache.solr.client.solrj.request.SolrQuery; import org.apache.solr.client.solrj.request.UpdateRequest; import org.apache.solr.client.solrj.request.XMLRequestWriter; +import org.apache.solr.client.solrj.request.json.JsonQueryRequest; import org.apache.solr.client.solrj.response.JavaBinResponseParser; import org.apache.solr.client.solrj.response.ResponseParser; import org.apache.solr.client.solrj.response.SolrPingResponse; @@ -637,6 +638,51 @@ public void testMaybeTryHeadRequestHasContentType() throws Exception { } } + @Test(timeout = 30000) + public void testConcurrentStreamedBodiesDoNotDeadlockWithHttp1() throws Exception { + DebugServlet.clear(); + String url = solrTestRule.getBaseUrl() + DEBUG_SERVLET_PATH; + + int concurrency = 8; + ExecutorService callers = + ExecutorUtil.newMDCAwareFixedThreadPool(concurrency, new NamedThreadFactory("test-caller")); + + try (HttpJdkSolrClient client = builder(url).useHttp1_1(true).build()) { + List> futures = new ArrayList<>(concurrency); + for (int i = 0; i < concurrency; i++) { + futures.add( + CompletableFuture.runAsync( + () -> { + JsonQueryRequest q = buildLargeBodyQuery(); + try { + q.process(client); + } catch (Exception ignored) { + } + }, + callers)); + } + CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])) + .get(45, TimeUnit.SECONDS); + } finally { + ExecutorUtil.shutdownAndAwaitTermination(callers); + } + } + + private static JsonQueryRequest buildLargeBodyQuery() { + StringBuilder filter = new StringBuilder("id:("); + for (int i = 0; i < 400; i++) { + if (i > 0) { + filter.append(" OR "); + } + filter.append("value_").append(i); + } + filter.append(')'); + JsonQueryRequest q = new JsonQueryRequest(); + q.setQuery("*:*"); + q.withFilter(filter.toString()); + return q; + } + /** * This is not required for any test, but there appears to be a bug in the JDK client where it * does not release all threads if the client has not performed any queries, even after a forced From e79a86d9319c8e3f70c9389b28ca009f5f51e2a1 Mon Sep 17 00:00:00 2001 From: Renato Haeberli Date: Sat, 25 Jul 2026 12:22:18 +0200 Subject: [PATCH 2/5] SOLR-18312: improving test not to swallow exceptions --- .../solr/client/solrj/impl/HttpJdkSolrClientTest.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java b/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java index 9aef536d370..9f9233f375e 100644 --- a/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java +++ b/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java @@ -641,6 +641,8 @@ public void testMaybeTryHeadRequestHasContentType() throws Exception { @Test(timeout = 30000) public void testConcurrentStreamedBodiesDoNotDeadlockWithHttp1() throws Exception { DebugServlet.clear(); + DebugServlet.addResponseHeader("Content-Type", "application/octet-stream"); + DebugServlet.responseBodyByQueryFragment.put("", javabinResponse()); String url = solrTestRule.getBaseUrl() + DEBUG_SERVLET_PATH; int concurrency = 8; @@ -656,13 +658,14 @@ public void testConcurrentStreamedBodiesDoNotDeadlockWithHttp1() throws Exceptio JsonQueryRequest q = buildLargeBodyQuery(); try { q.process(client); - } catch (Exception ignored) { + } catch (SolrServerException | IOException e) { + throw new RuntimeException(e); } }, callers)); } CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])) - .get(45, TimeUnit.SECONDS); + .get(30, TimeUnit.SECONDS); } finally { ExecutorUtil.shutdownAndAwaitTermination(callers); } From 7e3956bdafb3549f62ec0cd20d70284c351361e5 Mon Sep 17 00:00:00 2001 From: David Smiley Date: Sat, 25 Jul 2026 23:21:51 -0400 Subject: [PATCH 3/5] Use newMDCAwareCachedThreadPool --- .../client/solrj/impl/HttpJdkSolrClient.java | 21 +++---------------- 1 file changed, 3 insertions(+), 18 deletions(-) diff --git a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java index 5d53f0b5420..61964ca7eff 100644 --- a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java +++ b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java @@ -37,13 +37,9 @@ import java.util.HashMap; import java.util.Locale; import java.util.Map; -import java.util.concurrent.BlockingQueue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.SynchronousQueue; -import java.util.concurrent.TimeUnit; import java.util.regex.Matcher; import java.util.regex.Pattern; import javax.net.ssl.SSLContext; @@ -115,24 +111,13 @@ protected HttpJdkSolrClient(String serverBaseUrl, HttpJdkSolrClient.Builder buil this.executor = builder.getExecutor(); this.shutdownExecutor = false; } else { - BlockingQueue queue = new LinkedBlockingQueue<>(1024); this.executor = - new ExecutorUtil.MDCAwareThreadPoolExecutor( - 4, - 256, - 60, - TimeUnit.SECONDS, - queue, - new SolrNamedThreadFactory(this.getClass().getSimpleName())); + ExecutorUtil.newMDCAwareCachedThreadPool( + new SolrNamedThreadFactory(this.getClass().getSimpleName() + "-reqBody")); this.shutdownExecutor = true; } this.httpClientExecutor = - new ExecutorUtil.MDCAwareThreadPoolExecutor( - 0, - Integer.MAX_VALUE, - 60, - TimeUnit.SECONDS, - new SynchronousQueue<>(), + ExecutorUtil.newMDCAwareCachedThreadPool( new SolrNamedThreadFactory(this.getClass().getSimpleName() + "-http")); httpClientBuilder.executor(this.httpClientExecutor); From e534b380e080ee38d8af75d32f07884df6e1d02d Mon Sep 17 00:00:00 2001 From: David Smiley Date: Thu, 30 Jul 2026 01:18:06 -0400 Subject: [PATCH 4/5] Defer body thread usage until actually requested. --- .../client/solrj/impl/HttpJdkSolrClient.java | 103 ++++++++++-------- 1 file changed, 55 insertions(+), 48 deletions(-) diff --git a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java index 61964ca7eff..db810f23b8c 100644 --- a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java +++ b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java @@ -21,6 +21,7 @@ import java.io.InputStream; import java.io.PipedInputStream; import java.io.PipedOutputStream; +import java.io.UncheckedIOException; import java.lang.invoke.MethodHandles; import java.net.CookieHandler; import java.net.InetSocketAddress; @@ -149,7 +150,7 @@ protected CompletableFuture> requestInputStreamAsync( PreparedRequest pReq = prepareRequest(baseUrl, solrRequest, collection); return httpClient .sendAsync(pReq.reqb.build(), HttpResponse.BodyHandlers.ofInputStream()) - .whenComplete((httpResponse, throwable) -> releaseContentWriting(pReq)); + .whenComplete((httpResponse, throwable) -> pReq.releaseContentWriting()); } catch (Exception e) { CompletableFuture> cf = new CompletableFuture<>(); cf.completeExceptionally(e); @@ -164,7 +165,7 @@ public CompletableFuture> requestAsync( PreparedRequest pReq = prepareRequest(null, solrRequest, collection); return httpClient .sendAsync(pReq.reqb.build(), HttpResponse.BodyHandlers.ofInputStream()) - .whenComplete((httpResponse, throwable) -> releaseContentWriting(pReq)) + .whenComplete((httpResponse, throwable) -> pReq.releaseContentWriting()) .thenApply( httpResponse -> { try { @@ -181,21 +182,6 @@ public CompletableFuture> requestAsync( } } - private void releaseContentWriting(PreparedRequest pReq) { - if (pReq.contentWritingFuture != null) { - pReq.contentWritingFuture.cancel(true); - } - // Closing the sink is what unblocks a writer already stuck in the pipe; cancel() alone does - // not. - if (pReq.contentWritingSink != null) { - try { - pReq.contentWritingSink.close(); - } catch (IOException e) { - log.warn("Could not close content-writing pipe", e); - } - } - } - @Override public NamedList requestWithBaseUrl( String baseUrl, SolrRequest solrRequest, String collection) @@ -216,9 +202,7 @@ public NamedList requestWithBaseUrl( } catch (RuntimeException e) { throw new SolrServerException(e); } finally { - if (pReq.contentWritingFuture != null) { - pReq.contentWritingFuture.cancel(true); - } + pReq.releaseContentWriting(); // See // https://docs.oracle.com/en/java/javase/17/docs/api/java.net.http/java/net/http/HttpResponse.BodySubscribers.html#ofInputStream() @@ -254,7 +238,7 @@ protected PreparedRequest prepareRequest( ResponseParser parserToUse = responseParser(solrRequest); ModifiableSolrParams queryParams = initializeSolrParams(solrRequest, parserToUse); var reqb = HttpRequest.newBuilder(); - PreparedRequest pReq = null; + PreparedRequest pReq; try { switch (solrRequest.getMethod()) { case GET: @@ -291,7 +275,7 @@ private PreparedRequest prepareGet( reqb.GET(); decorateRequest(reqb, solrRequest); reqb.uri(new URI(url + queryParams.toQueryString())); - return new PreparedRequest(reqb, null, null); + return new PreparedRequest(reqb); } private PreparedRequest preparePutOrPost( @@ -322,28 +306,16 @@ private PreparedRequest preparePutOrPost( } HttpRequest.BodyPublisher bodyPublisher; - Future contentWritingFuture = null; - PipedInputStream contentWritingSink = null; + PreparedRequest pReq = new PreparedRequest(reqb); if (contentWriter != null) { boolean success = maybeTryHeadRequest(url); if (!success) { reqb.version(HttpClient.Version.HTTP_1_1); } - final PipedOutputStream source = new PipedOutputStream(); - contentWritingSink = new PipedInputStream(source); - final PipedInputStream sink = contentWritingSink; - bodyPublisher = HttpRequest.BodyPublishers.ofInputStream(() -> sink); - - contentWritingFuture = - executor.submit( - () -> { - try (source) { - contentWriter.write(source); - } catch (Exception e) { - log.error("Cannot write Content Stream", e); - } - }); + bodyPublisher = + HttpRequest.BodyPublishers.ofInputStream( + () -> pReq.beginContentWriting(contentWriter, this.executor)); } else if (streams != null && streams.size() == 1) { boolean success = maybeTryHeadRequest(url); if (!success) { @@ -374,25 +346,60 @@ private PreparedRequest preparePutOrPost( URI uriWithQueryParams = new URI(url + queryParams.toQueryString()); reqb.uri(uriWithQueryParams); - return new PreparedRequest(reqb, contentWritingFuture, contentWritingSink); + return pReq; } protected static class PreparedRequest { - Future contentWritingFuture; - PipedInputStream contentWritingSink; - HttpRequest.Builder reqb; + final HttpRequest.Builder reqb; ResponseParser parserToUse; String url; - PreparedRequest( - HttpRequest.Builder reqb, - Future contentWritingFuture, - PipedInputStream contentWritingSink) { + // Both remain null if the request has no streamed content, or if the body is never requested + // (e.g. the connection failed before sending it). Filled in lazily by + // beginContentWriting once the JDK HttpClient actually requests the body. + volatile PipedInputStream contentWritingSink; + volatile Future contentWritingFuture; + + PreparedRequest(HttpRequest.Builder reqb) { this.reqb = reqb; - this.contentWritingFuture = contentWritingFuture; - this.contentWritingSink = contentWritingSink; + } + + private PipedInputStream beginContentWriting( + RequestWriter.ContentWriter contentWriter, ExecutorService bodyExecutor) { + final PipedOutputStream source = new PipedOutputStream(); + try { + contentWritingSink = new PipedInputStream(source); + } catch (IOException e) { + throw new UncheckedIOException(e); + } + contentWritingFuture = + bodyExecutor.submit( + () -> { + try (source) { + contentWriter.write(source); + } catch (Exception e) { + log.error("Cannot write Content Stream", e); + } + }); + return contentWritingSink; + } + + private void releaseContentWriting() { + if (contentWritingFuture != null) { + contentWritingFuture.cancel(true); + } + // Closing the sink is what unblocks a writer already stuck in the pipe; cancel() alone does + // not. + PipedInputStream sink = contentWritingSink; + if (sink != null) { + try { + sink.close(); + } catch (IOException e) { + log.warn("Could not close content-writing pipe", e); + } + } } } From 0c5768f9fe6b5968e10b93fcc926b2f58a87130a Mon Sep 17 00:00:00 2001 From: David Smiley Date: Thu, 30 Jul 2026 10:04:16 -0400 Subject: [PATCH 5/5] Switch from volatile to synchronized --- .../client/solrj/impl/HttpJdkSolrClient.java | 16 +++++++++------- 1 file changed, 9 insertions(+), 7 deletions(-) diff --git a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java index db810f23b8c..64ae67c3182 100644 --- a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java +++ b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java @@ -359,14 +359,14 @@ protected static class PreparedRequest { // Both remain null if the request has no streamed content, or if the body is never requested // (e.g. the connection failed before sending it). Filled in lazily by // beginContentWriting once the JDK HttpClient actually requests the body. - volatile PipedInputStream contentWritingSink; - volatile Future contentWritingFuture; + private PipedInputStream contentWritingSink; + private Future contentWritingFuture; PreparedRequest(HttpRequest.Builder reqb) { this.reqb = reqb; } - private PipedInputStream beginContentWriting( + synchronized PipedInputStream beginContentWriting( RequestWriter.ContentWriter contentWriter, ExecutorService bodyExecutor) { final PipedOutputStream source = new PipedOutputStream(); try { @@ -374,9 +374,11 @@ private PipedInputStream beginContentWriting( } catch (IOException e) { throw new UncheckedIOException(e); } + contentWritingFuture = bodyExecutor.submit( () -> { + // note: doesn't need to synchronize with PreparedRequest.this try (source) { contentWriter.write(source); } catch (Exception e) { @@ -386,16 +388,16 @@ private PipedInputStream beginContentWriting( return contentWritingSink; } - private void releaseContentWriting() { + synchronized void releaseContentWriting() { if (contentWritingFuture != null) { contentWritingFuture.cancel(true); } + // Closing the sink is what unblocks a writer already stuck in the pipe; cancel() alone does // not. - PipedInputStream sink = contentWritingSink; - if (sink != null) { + if (contentWritingSink != null) { try { - sink.close(); + contentWritingSink.close(); } catch (IOException e) { log.warn("Could not close content-writing pipe", e); }