From 90066d8c1d7299aea88acfa4361da0b0fec44940 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?JB=20Onofr=C3=A9?= Date: Mon, 5 Oct 2026 18:56:23 +0200 Subject: [PATCH 1/2] GH-384: Wait for the server executor in FlightServer.close() gRPC reports the server as terminated once its transports are closed, without waiting for the calls still running on the executor. Such a call can still allocate or hold buffers (for instance while parsing an incoming message that is then discarded), so closing the allocator right after the server could report a leak. This is what made TestDoExchange.tearDown flaky. FlightServer.close() now also waits for the executor it created, within the existing 6 seconds budget. The workaround in TestDoExchange.testClientClose, which leaked the allocator on purpose because of the same race, is not needed anymore. --- .../org/apache/arrow/flight/FlightServer.java | 32 +++++++++++------- .../apache/arrow/flight/TestDoExchange.java | 9 ----- .../arrow/flight/TestServerOptions.java | 33 +++++++++++++++++++ 3 files changed, 53 insertions(+), 21 deletions(-) diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java index ac761457f5..c003b541b8 100644 --- a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java +++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java @@ -145,25 +145,33 @@ public boolean awaitTermination(final long timeout, final TimeUnit unit) /** Shutdown the server, waits for up to 6 seconds for successful shutdown before returning. */ @Override public void close() throws InterruptedException { + final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(6); shutdown(); final boolean terminated = awaitTermination(3000, TimeUnit.MILLISECONDS); if (terminated) { logger.debug("Server was terminated within 3s"); - return; - } - - // get more aggressive in termination. - server.shutdownNow(); + } else { + // get more aggressive in termination. + server.shutdownNow(); + + int count = 0; + while (!server.isTerminated() && count < 30) { + count++; + logger.debug("Waiting for termination"); + Thread.sleep(100); + } - int count = 0; - while (!server.isTerminated() && count < 30) { - count++; - logger.debug("Waiting for termination"); - Thread.sleep(100); + if (!server.isTerminated()) { + logger.warn("Couldn't shutdown server, resources likely will be leaked."); + } } - if (!server.isTerminated()) { - logger.warn("Couldn't shutdown server, resources likely will be leaked."); + // gRPC reports the server as terminated once its transports are closed, without waiting for + // the calls still running on the executor. Those calls may still allocate or hold buffers, so + // wait for them too; otherwise closing the allocator right after the server can report a leak. + if (grpcExecutor != null + && !grpcExecutor.awaitTermination(deadline - System.nanoTime(), TimeUnit.NANOSECONDS)) { + logger.warn("Couldn't shutdown server executor, resources likely will be leaked."); } } diff --git a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestDoExchange.java b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestDoExchange.java index 1cd6e95e3d..8967c92b9d 100644 --- a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestDoExchange.java +++ b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestDoExchange.java @@ -424,15 +424,6 @@ public void testClientClose() throws Exception { client.doExchange(FlightDescriptor.command(EXCHANGE_DO_GET))) { assertEquals(Producer.SCHEMA, stream.getReader().getSchema()); } - // Intentionally leak the allocator in this test. gRPC has a bug where it does not wait for all - // calls to complete - // when shutting down the server, so this test will fail otherwise because it closes the - // allocator while the - // server-side call still has memory allocated. - // TODO(ARROW-9586): fix this once we track outstanding RPCs outside of gRPC. - // https://stackoverflow.com/questions/46716024/ - allocator = null; - client = null; } /** Test closing with Metadata can't lead to error. */ diff --git a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java index c39ac922cf..bcb0ac4395 100644 --- a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java +++ b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java @@ -86,6 +86,39 @@ public void defaultExecutorClosed() throws Exception { assertTrue(executor.isShutdown()); } + /** + * Make sure that closing the server waits for the calls still running on the default executor, as + * they may still be using the allocator. + */ + @Test + public void closeWaitsForDefaultExecutor() throws Exception { + final AtomicBoolean callFinished = new AtomicBoolean(); + final FlightProducer producer = + new NoOpFlightProducer() { + @Override + public void doAction(CallContext context, Action action, StreamListener listener) { + // End the RPC, so that gRPC can shut down, but keep the executor thread busy + listener.onCompleted(); + try { + Thread.sleep(200); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + callFinished.set(true); + } + }; + try (BufferAllocator a = new RootAllocator(Long.MAX_VALUE); + FlightServer server = + FlightServer.builder(a, forGrpcInsecure(LOCALHOST, 0), producer).build().start()) { + try (FlightClient client = FlightClient.builder(a, server.getLocation()).build()) { + client.doAction(new Action("")).forEachRemaining(result -> {}); + } + server.close(); + assertTrue(callFinished.get()); + } + } + /** Make sure that if the user provides an executor to gRPC, then Flight does not close it. */ @Test public void suppliedExecutorNotClosed() throws Exception { From d991b4d953821e21549347dd4d8283587561b1ac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?JB=20Onofr=C3=A9?= Date: Tue, 6 Oct 2026 06:48:39 +0200 Subject: [PATCH 2/2] GH-384: Fix formatting in TestServerOptions --- .../test/java/org/apache/arrow/flight/TestServerOptions.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java index bcb0ac4395..b18e6cfefc 100644 --- a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java +++ b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestServerOptions.java @@ -96,7 +96,8 @@ public void closeWaitsForDefaultExecutor() throws Exception { final FlightProducer producer = new NoOpFlightProducer() { @Override - public void doAction(CallContext context, Action action, StreamListener listener) { + public void doAction( + CallContext context, Action action, StreamListener listener) { // End the RPC, so that gRPC can shut down, but keep the executor thread busy listener.onCompleted(); try {