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..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 @@ -86,6 +86,40 @@ 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 {