Usé Java 16 para realizar solicitudes a una API a través de HTTP. Para acelerar esto en general, lo cargué en un ForkJoinPool personalizado. He compilado un ejemplo de reproducción a continuación.
Desde que se pasó a Java 17 (openjdk build 17.0.1+12-39), esto arroja una RejectedExecutionException:
Caused by: java.util.concurrent.RejectedExecutionException: Thread limit exceeded replacing blocked worker at java.base/java.util.concurrent.ForkJoinPool.tryCompensate(ForkJoinPool.java:1819) at java.base/java.util.concurrent.ForkJoinPool.compensatedBlock(ForkJoinPool.java:3446) at java.base/java.util.concurrent.ForkJoinPool.managedBlock(ForkJoinPool.java:3432) at java.base/java.util.concurrent.CompletableFuture.waitingGet(CompletableFuture.java:1898) at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2072) at java.net.http/jdk.internal.net.http.HttpClientImpl.send(HttpClientImpl.java:553) at java.net.http/jdk.internal.net.http.HttpClientFacade.send(HttpClientFacade.java:119) at Test.lambda$retrieveMany$1(Test.java:30)¿Por qué pasó esto? ¿Cambió algo con respecto a ForkJoinPool que desconozco?
Código
import java.io.IOException; import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse.BodyHandlers; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.ForkJoinPool; import static java.util.concurrent.TimeUnit.MINUTES; import static java.util.stream.Collectors.toList; public class Test { public static void main(String[] args) throws ExecutionException, InterruptedException { final List<String> urls = List.of("https://stackoverflow.com", "https://stackoverflow.com", "https://stackoverflow.com"); // This succeeds on JDK 16, 17 retrieveMany(urls, 4); // This fails on JDK 17, but succeeds on 16 retrieveMany(urls, 3); } private static List<String> retrieveMany(List<String> urls, int threads) throws InterruptedException, ExecutionException { return new ForkJoinPool(threads, ForkJoinPool.defaultForkJoinWorkerThreadFactory, (t, e) -> {}, true, 0, threads, 1, null, 1, MINUTES) .submit(() -> urls.parallelStream() .map(url -> { try { return HttpClient.newBuilder().build().send(HttpRequest.newBuilder(URI.create(url)).build(), BodyHandlers.ofString()).body(); } catch (IOException | InterruptedException aE) { } return null; }) .collect(toList())) .get(); } }Has enviado una tarea, pero usa parallelStream() internamente, que luego ejecuta cada http en diferentes subprocesos del mismo grupo de combinación de bifurcación.
Hay una diferencia en la forma en que JDK16 y 17 tratan la situación en la que todos los subprocesos disponibles en el grupo están en uso: aquí es donde el parámetro saturated se vuelve relevante.
Cuando threads > urls.size() el grupo nunca se satura, pero en su segundo caso threads == urls.size() por lo que todos los subprocesos están en uso. Reemplace null en el constructor de ForkJoinPool por una variable de saturate para ver cuándo se activa la condición de prueba de saturación:
Predicate<? super ForkJoinPool> saturate = pool -> { boolean allow = false; System.out.println(Thread.currentThread().getName()+" saturate:"+allow); return allow; }; En JDK16, el predicado saturate se llama varias veces pero continúa, mientras que en JDK17 el procesamiento se detiene en la primera llamada si devuelve falso. Si cambia allow = true , JDK17 no enviará RejectedExecutionException cuando la cantidad de solicitudes en curso sea la misma que la de los subprocesos en uso para parallelStream() y continuará procesando más solicitudes cuando se completen otras solicitudes.