This is an automated email from the ASF dual-hosted git repository.

lidavidm pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-java.git


The following commit(s) were added to refs/heads/main by this push:
     new 91b4a2ca4 GH-1067: Close cached HDFS FileSystem instances (#1141)
91b4a2ca4 is described below

commit 91b4a2ca418ecb47cb8871d97588b19444337eb4
Author: Hélder Gregório <[email protected]>
AuthorDate: Wed Sep 9 21:41:27 2026 -0700

    GH-1067: Close cached HDFS FileSystem instances (#1141)
    
    ## What's Changed
    
    - This PR fixes JVM shutdown hangs after reading HDFS datasets through
    Arrow Java.
    - FileSystemDatasetFactory now tracks hdfs:// URIs used to create the
    factory. On close(), after releasing the native dataset factory, it
    best-effort closes the matching Hadoop FileSystem instances.
    - The Hadoop cleanup is done via reflection so Arrow Java does not add a
    production dependency on Hadoop. Non-HDFS URIs are ignored.
    
    Closes #1067 .
---
 dataset/pom.xml                                    |  65 ++++++++++
 dataset/src/main/java/module-info.java             |   1 +
 .../dataset/file/FileSystemDatasetFactory.java     |  76 ++++++++++++
 .../dataset/file/TestHdfsFileSystemCleanup.java    | 131 +++++++++++++++++++++
 4 files changed, 273 insertions(+)

diff --git a/dataset/pom.xml b/dataset/pom.xml
index 47b508073..2b1d8025f 100644
--- a/dataset/pom.xml
+++ b/dataset/pom.xml
@@ -54,6 +54,10 @@ under the License.
       <artifactId>arrow-c-data</artifactId>
       <scope>compile</scope>
     </dependency>
+    <dependency>
+      <groupId>org.slf4j</groupId>
+      <artifactId>slf4j-api</artifactId>
+    </dependency>
     <dependency>
       <groupId>org.immutables</groupId>
       <artifactId>value-annotations</artifactId>
@@ -112,6 +116,67 @@ under the License.
         </exclusion>
       </exclusions>
     </dependency>
+    <dependency>
+      <groupId>org.apache.hadoop</groupId>
+      <artifactId>hadoop-hdfs</artifactId>
+      <version>${dep.hadoop.version}</version>
+      <scope>test</scope>
+      <exclusions>
+        <exclusion>
+          <groupId>commons-logging</groupId>
+          <artifactId>commons-logging</artifactId>
+        </exclusion>
+        <exclusion>
+          <groupId>log4j</groupId>
+          <artifactId>log4j</artifactId>
+        </exclusion>
+        <exclusion>
+          <groupId>org.slf4j</groupId>
+          <artifactId>slf4j-log4j12</artifactId>
+        </exclusion>
+      </exclusions>
+    </dependency>
+    <dependency>
+      <groupId>org.apache.hadoop</groupId>
+      <artifactId>hadoop-hdfs</artifactId>
+      <version>${dep.hadoop.version}</version>
+      <type>test-jar</type>
+      <scope>test</scope>
+      <exclusions>
+        <exclusion>
+          <groupId>commons-logging</groupId>
+          <artifactId>commons-logging</artifactId>
+        </exclusion>
+        <exclusion>
+          <groupId>log4j</groupId>
+          <artifactId>log4j</artifactId>
+        </exclusion>
+        <exclusion>
+          <groupId>org.slf4j</groupId>
+          <artifactId>slf4j-log4j12</artifactId>
+        </exclusion>
+      </exclusions>
+    </dependency>
+    <dependency>
+      <groupId>org.apache.hadoop</groupId>
+      <artifactId>hadoop-minicluster</artifactId>
+      <version>${dep.hadoop.version}</version>
+      <scope>test</scope>
+      <exclusions>
+        <exclusion>
+          <groupId>commons-logging</groupId>
+          <artifactId>commons-logging</artifactId>
+        </exclusion>
+        <exclusion>
+          <groupId>log4j</groupId>
+          <artifactId>log4j</artifactId>
+        </exclusion>
+        <exclusion>
+          <groupId>org.slf4j</groupId>
+          <artifactId>slf4j-log4j12</artifactId>
+        </exclusion>
+      </exclusions>
+    </dependency>
     <dependency>
       <groupId>com.google.guava</groupId>
       <artifactId>guava</artifactId>
diff --git a/dataset/src/main/java/module-info.java 
b/dataset/src/main/java/module-info.java
index 3092281a9..3b60dccee 100644
--- a/dataset/src/main/java/module-info.java
+++ b/dataset/src/main/java/module-info.java
@@ -26,4 +26,5 @@ open module org.apache.arrow.dataset {
   requires org.apache.arrow.c;
   requires org.apache.arrow.memory.core;
   requires org.apache.arrow.vector;
+  requires org.slf4j;
 }
diff --git 
a/dataset/src/main/java/org/apache/arrow/dataset/file/FileSystemDatasetFactory.java
 
b/dataset/src/main/java/org/apache/arrow/dataset/file/FileSystemDatasetFactory.java
index fcf124a61..9b6e8731e 100644
--- 
a/dataset/src/main/java/org/apache/arrow/dataset/file/FileSystemDatasetFactory.java
+++ 
b/dataset/src/main/java/org/apache/arrow/dataset/file/FileSystemDatasetFactory.java
@@ -16,18 +16,30 @@
  */
 package org.apache.arrow.dataset.file;
 
+import java.lang.reflect.Method;
+import java.net.URI;
+import java.net.URISyntaxException;
+import java.util.LinkedHashSet;
 import java.util.Optional;
+import java.util.Set;
 import org.apache.arrow.dataset.jni.NativeDatasetFactory;
 import org.apache.arrow.dataset.jni.NativeMemoryPool;
 import org.apache.arrow.dataset.scanner.FragmentScanOptions;
 import org.apache.arrow.memory.BufferAllocator;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /** Java binding of the C++ FileSystemDatasetFactory. */
 public class FileSystemDatasetFactory extends NativeDatasetFactory {
 
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(FileSystemDatasetFactory.class);
+
+  private final Set<URI> hdfsFileSystems;
+
   public FileSystemDatasetFactory(
       BufferAllocator allocator, NativeMemoryPool memoryPool, FileFormat 
format, String uri) {
     super(allocator, memoryPool, createNative(format, uri, Optional.empty()));
+    this.hdfsFileSystems = toHdfsFileSystems(uri);
   }
 
   public FileSystemDatasetFactory(
@@ -37,11 +49,13 @@ public class FileSystemDatasetFactory extends 
NativeDatasetFactory {
       String uri,
       Optional<FragmentScanOptions> fragmentScanOptions) {
     super(allocator, memoryPool, createNative(format, uri, 
fragmentScanOptions));
+    this.hdfsFileSystems = toHdfsFileSystems(uri);
   }
 
   public FileSystemDatasetFactory(
       BufferAllocator allocator, NativeMemoryPool memoryPool, FileFormat 
format, String[] uris) {
     super(allocator, memoryPool, createNative(format, uris, Optional.empty()));
+    this.hdfsFileSystems = toHdfsFileSystems(uris);
   }
 
   public FileSystemDatasetFactory(
@@ -51,6 +65,68 @@ public class FileSystemDatasetFactory extends 
NativeDatasetFactory {
       String[] uris,
       Optional<FragmentScanOptions> fragmentScanOptions) {
     super(allocator, memoryPool, createNative(format, uris, 
fragmentScanOptions));
+    this.hdfsFileSystems = toHdfsFileSystems(uris);
+  }
+
+  /**
+   * Close this factory and release the native instance. For HDFS URIs, also 
closes the cached
+   * Hadoop FileSystem to release non-daemon threads that would otherwise 
prevent JVM exit. See <a
+   * href="https://github.com/apache/arrow-java/issues/1067";>#1067</a>.
+   */
+  @Override
+  public synchronized void close() {
+    try {
+      super.close();
+    } finally {
+      hdfsFileSystems.forEach(FileSystemDatasetFactory::closeHadoopFileSystem);
+    }
+  }
+
+  /**
+   * For each {@code hdfs://} URI, close the cached Hadoop FileSystem. When 
Arrow C++ accesses HDFS
+   * via libhdfs, the Hadoop Java client creates cached FileSystem instances 
with non-daemon threads
+   * (IPC connections, lease renewers) that prevent JVM exit. Closing the 
FileSystem terminates
+   * these connections. Uses reflection to avoid a compile-time dependency on 
hadoop-common.
+   */
+  static void closeHadoopFileSystemsIfHdfs(String... uris) {
+    
toHdfsFileSystems(uris).forEach(FileSystemDatasetFactory::closeHadoopFileSystem);
+  }
+
+  private static Set<URI> toHdfsFileSystems(String... uris) {
+    Set<URI> hdfsFileSystems = new LinkedHashSet<>();
+    if (uris == null) {
+      return hdfsFileSystems;
+    }
+    for (String uri : uris) {
+      if (uri == null) {
+        continue;
+      }
+      try {
+        URI parsedUri = new URI(uri);
+        if ("hdfs".equalsIgnoreCase(parsedUri.getScheme())) {
+          hdfsFileSystems.add(
+              new URI(parsedUri.getScheme(), parsedUri.getAuthority(), null, 
null, null));
+        }
+      } catch (URISyntaxException e) {
+        // Ignore here; native factory creation reports invalid user URIs.
+      }
+    }
+    return hdfsFileSystems;
+  }
+
+  private static void closeHadoopFileSystem(URI hdfsUri) {
+    try {
+      Class<?> confClass = 
Class.forName("org.apache.hadoop.conf.Configuration");
+      Object conf = confClass.getDeclaredConstructor().newInstance();
+      Class<?> fsClass = Class.forName("org.apache.hadoop.fs.FileSystem");
+      Method getMethod = fsClass.getMethod("get", URI.class, confClass);
+      Object fs = getMethod.invoke(null, hdfsUri, conf);
+      Method closeMethod = fsClass.getMethod("close");
+      closeMethod.invoke(fs);
+    } catch (Exception e) {
+      // Best-effort cleanup; Hadoop may not be on classpath or FileSystem 
already closed.
+      LOGGER.debug("Failed to close Hadoop FileSystem for {}", hdfsUri, e);
+    }
   }
 
   private static long createNative(
diff --git 
a/dataset/src/test/java/org/apache/arrow/dataset/file/TestHdfsFileSystemCleanup.java
 
b/dataset/src/test/java/org/apache/arrow/dataset/file/TestHdfsFileSystemCleanup.java
new file mode 100644
index 000000000..7440d1b44
--- /dev/null
+++ 
b/dataset/src/test/java/org/apache/arrow/dataset/file/TestHdfsFileSystemCleanup.java
@@ -0,0 +1,131 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.arrow.dataset.file;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.concurrent.TimeUnit;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hdfs.MiniDFSCluster;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/** Regression test for <a 
href="https://github.com/apache/arrow-java/issues/1067";>#1067</a>. */
+public class TestHdfsFileSystemCleanup {
+
+  private static final int CHILD_TIMEOUT_SECONDS = 10;
+
+  private static MiniDFSCluster cluster;
+
+  @TempDir static File clusterDir;
+
+  @BeforeAll
+  static void startCluster() throws IOException {
+    Configuration conf = new Configuration();
+    conf.set(MiniDFSCluster.HDFS_MINIDFS_BASEDIR, 
clusterDir.getAbsolutePath());
+    cluster = new MiniDFSCluster.Builder(conf).numDataNodes(1).build();
+    cluster.waitActive();
+  }
+
+  @AfterAll
+  static void stopCluster() {
+    if (cluster != null) {
+      cluster.shutdown();
+    }
+  }
+
+  @Test
+  void testJvmHangsWithoutCleanup() throws Exception {
+    Process child = forkChildProcess(false);
+    assertFalse(waitForExit(child), "JVM should hang when HDFS is not cleaned 
up");
+  }
+
+  @Test
+  void testJvmExitsWithCleanup() throws Exception {
+    Process child = forkChildProcess(true);
+    assertTrue(waitForExit(child), "JVM should exit when 
FileSystemDatasetFactory cleanup runs");
+    assertEquals(0, child.exitValue(), "Child process should exit cleanly 
(exit code 0)");
+  }
+
+  private boolean waitForExit(Process child) throws InterruptedException {
+    boolean exited = child.waitFor(CHILD_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+    if (!exited) {
+      child.destroyForcibly();
+    }
+    return exited;
+  }
+
+  private Process forkChildProcess(boolean withCleanup) throws IOException {
+    String classpath = System.getProperty("java.class.path");
+    int port = cluster.getNameNodePort();
+    ProcessBuilder pb =
+        new ProcessBuilder(
+            ProcessHandle.current().info().command().orElse("java"),
+            "-cp",
+            classpath,
+            HdfsClientSimulator.class.getName(),
+            String.valueOf(port),
+            String.valueOf(withCleanup));
+    pb.redirectError(ProcessBuilder.Redirect.DISCARD);
+    pb.redirectOutput(ProcessBuilder.Redirect.DISCARD);
+    return pb.start();
+  }
+
+  /** Simulates libhdfs leaving a non-daemon thread attached to an HDFS 
connection. */
+  public static class HdfsClientSimulator {
+    public static void main(String[] args) throws Exception {
+      int port = Integer.parseInt(args[0]);
+      boolean withCleanup = Boolean.parseBoolean(args[1]);
+
+      Configuration conf = new Configuration();
+      String hdfsUri = "hdfs://localhost:" + port;
+      conf.set("fs.defaultFS", hdfsUri);
+
+      FileSystem fs = FileSystem.get(conf);
+      fs.exists(new Path("/"));
+
+      Thread connectionThread =
+          new Thread(
+              () -> {
+                while (true) {
+                  try {
+                    fs.getFileStatus(new Path("/"));
+                    Thread.sleep(500);
+                  } catch (Exception e) {
+                    break;
+                  }
+                }
+              },
+              "simulated-libhdfs-ipc-thread");
+      connectionThread.setDaemon(false);
+      connectionThread.start();
+
+      if (withCleanup) {
+        FileSystemDatasetFactory.closeHadoopFileSystemsIfHdfs(hdfsUri);
+        connectionThread.join(5000);
+      }
+    }
+  }
+}

Reply via email to