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

mattcasters pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git


The following commit(s) were added to refs/heads/main by this push:
     new a604ca3c50 Issue #8440 : Do not hold the Kerberos lock during HDFS 
streaming PUTs (#8487)
a604ca3c50 is described below

commit a604ca3c50abdce4f2e1493d165658b0de377aff
Author: Matt Casters <[email protected]>
AuthorDate: Mon Sep 21 11:12:13 2026 +0200

    Issue #8440 : Do not hold the Kerberos lock during HDFS streaming PUTs 
(#8487)
    
    Parallel Parquet writes over WebHDFS/Kerberos stalled because doAs()
    kept JVM_KERBEROS for the whole file upload. SPNEGO stays serialized;
    HTTP execute, the connection pool, and stream close are no longer
    able to deadlock the other transform copies.
    
    Fix validated with 20 threads torture test and a few hundred million rows 
in parquet files.
---
 .../ROOT/pages/metadata-types/hdfs-connection.adoc |   3 +
 .../apache/hop/vfs/hdfs/HdfsConnectionTester.java  |   2 +-
 .../java/org/apache/hop/vfs/hdfs/HdfsHttp.java     |  20 +++-
 .../vfs/hdfs/client/HdfsStreamingOutputStream.java |  87 +++++++++++------
 .../hop/vfs/hdfs/client/HdfsWebHdfsClient.java     |  56 +++++------
 .../hop/vfs/hdfs/kerberos/HdfsKerberosSession.java |  60 +++++++++---
 .../vfs/hdfs/messages/messages_en_US.properties    |   1 +
 .../hdfs/client/HdfsStreamingOutputStreamTest.java |  39 +++++++-
 .../hop/vfs/hdfs/client/HdfsWebHdfsClientTest.java | 104 +++++++++++++++++++++
 .../hop/vfs/hdfs/client/WebHdfsTestServer.java     |  44 ++++++++-
 .../vfs/hdfs/kerberos/HdfsKerberosSessionTest.java |  48 ++++++++++
 11 files changed, 380 insertions(+), 84 deletions(-)

diff --git 
a/docs/hop-user-manual/modules/ROOT/pages/metadata-types/hdfs-connection.adoc 
b/docs/hop-user-manual/modules/ROOT/pages/metadata-types/hdfs-connection.adoc
index a7d8dfa9cd..ca41f730bf 100644
--- 
a/docs/hop-user-manual/modules/ROOT/pages/metadata-types/hdfs-connection.adoc
+++ 
b/docs/hop-user-manual/modules/ROOT/pages/metadata-types/hdfs-connection.adoc
@@ -85,6 +85,9 @@ That is the library-conflict and OpenShift-network path this 
plugin exists to av
 The plugin re-logs from the keytab in Java (JAAS).
 It does not call the `kinit` binary and does not use Hadoop 
`UserGroupInformation`.
 
+Parallel Parquet File Output copies (and other concurrent VFS writers) share 
one HTTP connection pool per named connection.
+Kerberos SPNEGO is generated per NameNode / HttpFS / Knox request and does 
**not** hold the login lock while file bytes are uploaded, so several files can 
be written at once.
+
 === TLS
 
 |===
diff --git 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsConnectionTester.java
 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsConnectionTester.java
index c970a513d6..9b5eb92cea 100644
--- 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsConnectionTester.java
+++ 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsConnectionTester.java
@@ -95,7 +95,7 @@ public final class HdfsConnectionTester {
           throw failed(succeeded, "Kerberos", probeHost, e);
         }
         try {
-          session.doAs(() -> HdfsSpnego.authorizationHeader(probeHost));
+          session.gss(() -> HdfsSpnego.authorizationHeader(probeHost));
           succeeded.add(BaseMessages.getString(PKG, "Hdfs.Test.Ok.Spnego", 
probeHost));
         } catch (Exception e) {
           throw failed(succeeded, "SPNEGO", "HTTP@" + probeHost, e);
diff --git 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsHttp.java 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsHttp.java
index a742c8c3f5..6ea280c5d5 100644
--- a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsHttp.java
+++ b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/HdfsHttp.java
@@ -22,12 +22,14 @@ import javax.net.ssl.HostnameVerifier;
 import javax.net.ssl.SSLContext;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.commons.vfs2.FileSystemException;
+import org.apache.hc.client5.http.config.ConnectionConfig;
 import org.apache.hc.client5.http.config.RequestConfig;
 import org.apache.hc.client5.http.impl.classic.CloseableHttpClient;
 import org.apache.hc.client5.http.impl.classic.HttpClientBuilder;
 import 
org.apache.hc.client5.http.impl.io.PoolingHttpClientConnectionManagerBuilder;
 import org.apache.hc.client5.http.ssl.DefaultClientTlsStrategy;
 import org.apache.hc.client5.http.ssl.NoopHostnameVerifier;
+import org.apache.hc.core5.util.TimeValue;
 import org.apache.hc.core5.util.Timeout;
 import org.apache.hop.core.Const;
 import org.apache.hop.core.logging.LogChannel;
@@ -54,12 +56,24 @@ public final class HdfsHttp {
     builder.disableContentCompression();
     // CREATE/OPEN 307s are followed explicitly so NameNode SPNEGO is not sent 
to a DataNode.
     builder.disableRedirectHandling();
+    builder.evictIdleConnections(TimeValue.ofSeconds(30));
     builder.setDefaultRequestConfig(
         RequestConfig.custom()
-            .setConnectTimeout(Timeout.ofSeconds(30))
+            .setConnectionRequestTimeout(Timeout.ofSeconds(30))
             .setResponseTimeout(Timeout.ofSeconds(300))
+            .setExpectContinueEnabled(false)
             .build());
     try {
+      PoolingHttpClientConnectionManagerBuilder pool =
+          PoolingHttpClientConnectionManagerBuilder.create()
+              .setMaxConnTotal(128)
+              .setMaxConnPerRoute(64)
+              .setDefaultConnectionConfig(
+                  ConnectionConfig.custom()
+                      .setConnectTimeout(Timeout.ofSeconds(30))
+                      .setSocketTimeout(Timeout.ofSeconds(300))
+                      .setValidateAfterInactivity(TimeValue.ofSeconds(10))
+                      .build());
       if (https) {
         SSLContext sslContext = HdfsTls.sslContext(variables, meta);
         HostnameVerifier verifier =
@@ -68,9 +82,9 @@ public final class HdfsHttp {
             verifier == null
                 ? new DefaultClientTlsStrategy(sslContext)
                 : new DefaultClientTlsStrategy(sslContext, verifier);
-        builder.setConnectionManager(
-            
PoolingHttpClientConnectionManagerBuilder.create().setTlsSocketStrategy(tls).build());
+        pool.setTlsSocketStrategy(tls);
       }
+      builder.setConnectionManager(pool.build());
       return builder.build();
     } catch (Exception e) {
       throw new FileSystemException("Unable to create HTTP client for HDFS 
VFS", e);
diff --git 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStream.java
 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStream.java
index c848080cc4..4935ce51b4 100644
--- 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStream.java
+++ 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStream.java
@@ -24,6 +24,9 @@ import java.io.PipedOutputStream;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import org.apache.hc.client5.http.classic.methods.HttpUriRequestBase;
 import org.apache.hop.i18n.BaseMessages;
 import org.apache.hop.vfs.hdfs.HdfsTransport;
 
@@ -34,47 +37,63 @@ import org.apache.hop.vfs.hdfs.HdfsTransport;
  */
 public class HdfsStreamingOutputStream extends OutputStream {
   private static final Class<?> PKG = HdfsTransport.class;
-  private static final int PIPE_BUFFER = 1024 * 1024;
+  private static final int PIPE_BUFFER = 8 * 1024 * 1024;
+  static final long CLOSE_TIMEOUT_SECONDS = 120;
 
   @FunctionalInterface
   public interface Uploader {
-    void upload(InputStream body) throws IOException;
+    void upload(InputStream body, HdfsStreamingOutputStream stream) throws 
IOException;
   }
 
   private final PipedOutputStream pipe;
-  private final Future<Void> upload;
   private final String path;
+  private final long closeTimeoutSeconds;
+  private Future<Void> upload;
+  private volatile HttpUriRequestBase inflight;
   private volatile IOException uploadError;
   private boolean closed;
 
   public static HdfsStreamingOutputStream start(
       ExecutorService executor, Uploader uploader, String path) throws 
IOException {
+    return start(executor, uploader, path, CLOSE_TIMEOUT_SECONDS);
+  }
+
+  static HdfsStreamingOutputStream start(
+      ExecutorService executor, Uploader uploader, String path, long 
closeTimeoutSeconds)
+      throws IOException {
     PipedInputStream in = new PipedInputStream(PIPE_BUFFER);
     PipedOutputStream out = new PipedOutputStream(in);
     HdfsStreamingOutputStream stream =
-        new HdfsStreamingOutputStream(out, executor, uploader, in, path);
-    return stream;
-  }
-
-  private HdfsStreamingOutputStream(
-      PipedOutputStream pipe,
-      ExecutorService executor,
-      Uploader uploader,
-      PipedInputStream in,
-      String path) {
-    this.pipe = pipe;
-    this.path = path;
-    this.upload =
+        new HdfsStreamingOutputStream(out, path, closeTimeoutSeconds);
+    stream.upload =
         executor.submit(
             () -> {
               try (PipedInputStream body = in) {
-                uploader.upload(body);
+                uploader.upload(body, stream);
               } catch (IOException e) {
-                uploadError = e;
+                stream.uploadError = e;
                 throw e;
               }
               return null;
             });
+    return stream;
+  }
+
+  private HdfsStreamingOutputStream(PipedOutputStream pipe, String path, long 
closeTimeoutSeconds) {
+    this.pipe = pipe;
+    this.path = path;
+    this.closeTimeoutSeconds = closeTimeoutSeconds;
+  }
+
+  void watch(HttpUriRequestBase request) {
+    this.inflight = request;
+  }
+
+  private void abortInflight() {
+    HttpUriRequestBase request = inflight;
+    if (request != null) {
+      request.cancel();
+    }
   }
 
   @Override
@@ -110,18 +129,28 @@ public class HdfsStreamingOutputStream extends 
OutputStream {
     try {
       pipe.close();
     } finally {
-      try {
-        upload.get();
-      } catch (InterruptedException e) {
-        Thread.currentThread().interrupt();
-        throw new IOException(BaseMessages.getString(PKG, 
"Hdfs.Error.UploadFailed", path), e);
-      } catch (ExecutionException e) {
-        Throwable cause = e.getCause() != null ? e.getCause() : e;
-        if (cause instanceof IOException io) {
-          throw io;
-        }
-        throw new IOException(BaseMessages.getString(PKG, 
"Hdfs.Error.UploadFailed", path), cause);
+      awaitUpload();
+    }
+  }
+
+  private void awaitUpload() throws IOException {
+    try {
+      upload.get(closeTimeoutSeconds, TimeUnit.SECONDS);
+    } catch (TimeoutException e) {
+      abortInflight();
+      upload.cancel(true);
+      throw new IOException(BaseMessages.getString(PKG, 
"Hdfs.Error.UploadTimedOut", path), e);
+    } catch (InterruptedException e) {
+      abortInflight();
+      upload.cancel(true);
+      Thread.currentThread().interrupt();
+      throw new IOException(BaseMessages.getString(PKG, 
"Hdfs.Error.UploadFailed", path), e);
+    } catch (ExecutionException e) {
+      Throwable cause = e.getCause() != null ? e.getCause() : e;
+      if (cause instanceof IOException io) {
+        throw io;
       }
+      throw new IOException(BaseMessages.getString(PKG, 
"Hdfs.Error.UploadFailed", path), cause);
     }
     if (uploadError != null) {
       throw uploadError;
diff --git 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClient.java
 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClient.java
index 2273e6fe4a..cb46507f15 100644
--- 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClient.java
+++ 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClient.java
@@ -27,7 +27,6 @@ import java.net.URI;
 import java.net.URISyntaxException;
 import java.net.URLEncoder;
 import java.nio.charset.StandardCharsets;
-import java.security.PrivilegedExceptionAction;
 import java.util.ArrayList;
 import java.util.LinkedHashMap;
 import java.util.List;
@@ -188,7 +187,7 @@ public class HdfsWebHdfsClient {
         String uri = buildUri(endpoint, path, params);
         HttpGet request = new HttpGet(uri);
         addSpnego(request, uri);
-        response = privileged(() -> httpClient.execute(request));
+        response = httpClient.execute(request);
         int code = response.getCode();
         if (code >= 400) {
           String error = readBody(response);
@@ -239,7 +238,8 @@ public class HdfsWebHdfsClient {
     if (transport.dataOnCreateRequest()) {
       params.put("data", "true");
       String uri = firstWorkingUri(path, params);
-      return HdfsStreamingOutputStream.start(executor, in -> putStream(uri, 
in), path);
+      return HdfsStreamingOutputStream.start(
+          executor, (in, stream) -> putStream(uri, in, stream), path);
     }
     params.put("noredirect", "true");
     String body = executeString("PUT", path, params, null);
@@ -247,7 +247,8 @@ public class HdfsWebHdfsClient {
     if (location == null || location.isBlank()) {
       throw new IOException("WebHDFS CREATE did not return a DataNode Location 
for " + path);
     }
-    return HdfsStreamingOutputStream.start(executor, in -> putStream(location, 
in), path);
+    return HdfsStreamingOutputStream.start(
+        executor, (in, stream) -> putStream(location, in, stream), path);
   }
 
   private String locationFromCreate(String body) throws IOException {
@@ -270,9 +271,11 @@ public class HdfsWebHdfsClient {
     return status;
   }
 
-  private void putStream(String uri, InputStream body) throws IOException {
+  private void putStream(String uri, InputStream body, 
HdfsStreamingOutputStream stream)
+      throws IOException {
     rejectHttpDowngrade(uri);
     HttpPut put = new HttpPut(uri);
+    stream.watch(put);
     put.setEntity(new InputStreamEntity(body, 
ContentType.APPLICATION_OCTET_STREAM));
     put.setHeader("Content-Type", "application/octet-stream");
     if (shouldSpnegoForLocation(uri)) {
@@ -289,7 +292,7 @@ public class HdfsWebHdfsClient {
     }
     CloseableHttpResponse response = null;
     try {
-      response = privileged(() -> httpClient.execute(get));
+      response = httpClient.execute(get);
       int code = response.getCode();
       if (code >= 400) {
         String error = readBody(response);
@@ -461,32 +464,17 @@ public class HdfsWebHdfsClient {
   }
 
   private String executeRequest(ClassicHttpRequest request) throws IOException 
{
-    try {
-      return privileged(
-          () ->
-              httpClient.execute(
-                  request,
-                  response -> {
-                    int code = response.getCode();
-                    String body = readBody(response);
-                    if (code >= 400) {
-                      throw new IOException(
-                          errorMessage(request.getMethod(), 
request.getRequestUri(), code, body));
-                    }
-                    return body == null ? "" : body;
-                  }));
-    } catch (IOException e) {
-      throw e;
-    } catch (Exception e) {
-      throw new IOException(e);
-    }
-  }
-
-  private <T> T privileged(PrivilegedExceptionAction<T> action) throws 
Exception {
-    if (kerberos && kerberosSession != null) {
-      return kerberosSession.doAs(action);
-    }
-    return action.run();
+    return httpClient.execute(
+        request,
+        response -> {
+          int code = response.getCode();
+          String body = readBody(response);
+          if (code >= 400) {
+            throw new IOException(
+                errorMessage(request.getMethod(), request.getRequestUri(), 
code, body));
+          }
+          return body == null ? "" : body;
+        });
   }
 
   private HttpUriRequestBase request(String method, String uri) throws 
IOException {
@@ -513,8 +501,8 @@ public class HdfsWebHdfsClient {
     }
     try {
       // GSS reads the TGT from the current Subject. useSubjectCredsOnly=true, 
so this
-      // must run inside session.doAs, not on the calling thread after a 
separate login.
-      String header = privileged(() -> HdfsSpnego.authorizationHeader(host));
+      // must run inside session.gss / Subject.doAs. HTTP execute stays 
outside that lock.
+      String header = kerberosSession.gss(() -> 
HdfsSpnego.authorizationHeader(host));
       request.setHeader("Authorization", header);
     } catch (GSSException e) {
       throw new IOException(
diff --git 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSession.java
 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSession.java
index 974e101b81..7f8aeb5cc7 100644
--- 
a/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSession.java
+++ 
b/plugins/tech/hadoop/src/main/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSession.java
@@ -57,6 +57,9 @@ public class HdfsKerberosSession {
   private volatile LoginContext loginContext;
   private volatile long loginTimeMillis;
 
+  /** Test seam: skip JAAS login and use this Subject in {@link #doAs}. */
+  private volatile Subject subjectOverride;
+
   public HdfsKerberosSession(IVariables variables, HdfsMeta meta) {
     this(variables, meta, new KerberosUtil());
   }
@@ -111,24 +114,59 @@ public class HdfsKerberosSession {
     LogChannel.GENERAL.logBasic(BaseMessages.getString(PKG, 
"Hdfs.Log.KerberosRenew", principal));
   }
 
+  /**
+   * Run {@code action} as the logged-in Subject without holding {@link 
#JVM_KERBEROS} for the
+   * duration. HTTP transfers must not wrap execute in this method; use {@link 
#gss} for SPNEGO.
+   */
   @SuppressWarnings("removal")
   public <T> T doAs(PrivilegedExceptionAction<T> action) throws Exception {
+    Subject subject;
     synchronized (JVM_KERBEROS) {
-      if (loginContext == null) {
-        login();
-      }
-      try {
-        return Subject.doAs(loginContext.getSubject(), action);
-      } catch (PrivilegedActionException e) {
-        Throwable cause = e.getCause() != null ? e.getCause() : e;
-        if (cause instanceof Exception exception) {
-          throw exception;
-        }
-        throw new Exception(cause);
+      subject = currentSubject();
+    }
+    return callAs(subject, action);
+  }
+
+  /**
+   * Run a short GSS action (SPNEGO token). Holds {@link #JVM_KERBEROS} only 
for this call so
+   * concurrent {@code initSecContext} and keytab re-login cannot overlap. 
HTTP transfers must not
+   * use this.
+   */
+  public <T> T gss(PrivilegedExceptionAction<T> action) throws Exception {
+    synchronized (JVM_KERBEROS) {
+      return doAs(action);
+    }
+  }
+
+  @SuppressWarnings("removal")
+  private static <T> T callAs(Subject subject, PrivilegedExceptionAction<T> 
action)
+      throws Exception {
+    try {
+      return Subject.doAs(subject, action);
+    } catch (PrivilegedActionException e) {
+      Throwable cause = e.getCause() != null ? e.getCause() : e;
+      if (cause instanceof Exception exception) {
+        throw exception;
       }
+      throw new Exception(cause);
     }
   }
 
+  private Subject currentSubject() throws LoginException {
+    Subject override = subjectOverride;
+    if (override != null) {
+      return override;
+    }
+    if (loginContext == null) {
+      login();
+    }
+    return loginContext.getSubject();
+  }
+
+  void useSubjectForTest(Subject subject) {
+    this.subjectOverride = subject;
+  }
+
   public void close() {
     HdfsKerberosRenewer.getInstance().unregister(this);
   }
diff --git 
a/plugins/tech/hadoop/src/main/resources/org/apache/hop/vfs/hdfs/messages/messages_en_US.properties
 
b/plugins/tech/hadoop/src/main/resources/org/apache/hop/vfs/hdfs/messages/messages_en_US.properties
index 8227bbfedb..bfbd7c63a2 100644
--- 
a/plugins/tech/hadoop/src/main/resources/org/apache/hop/vfs/hdfs/messages/messages_en_US.properties
+++ 
b/plugins/tech/hadoop/src/main/resources/org/apache/hop/vfs/hdfs/messages/messages_en_US.properties
@@ -20,6 +20,7 @@ Hdfs.Error.MissingHost=HDFS VFS: endpoint hostname is not set 
on connection '{0}
 Hdfs.Error.KerberosLogin=HDFS VFS: Kerberos login failed for principal '{0}'
 Hdfs.Error.AppendNotSupported=Appending to HDFS files is not supported
 Hdfs.Error.UploadFailed=HDFS VFS: streaming upload of '{0}' failed
+Hdfs.Error.UploadTimedOut=HDFS VFS: streaming upload of '{0}' timed out; the 
HTTP PUT was cancelled
 Hdfs.Log.KerberosLogin=HDFS VFS: Kerberos login for principal {0}
 Hdfs.Log.KerberosRenew=HDFS VFS: renewed Kerberos ticket for principal {0}
 Hdfs.Test.Cluster.Success=Reached {0} at {1}:{2} (GETFILESTATUS {3}).
diff --git 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStreamTest.java
 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStreamTest.java
index 92d07b7108..6c2b3f8c54 100644
--- 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStreamTest.java
+++ 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsStreamingOutputStreamTest.java
@@ -18,9 +18,13 @@ package org.apache.hop.vfs.hdfs.client;
 
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 import java.io.ByteArrayOutputStream;
+import java.io.IOException;
 import java.nio.charset.StandardCharsets;
+import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.Executors;
 import org.junit.jupiter.api.Test;
 
@@ -32,7 +36,7 @@ class HdfsStreamingOutputStreamTest {
     var executor = Executors.newSingleThreadExecutor();
     try (HdfsStreamingOutputStream out =
         HdfsStreamingOutputStream.start(
-            executor, in -> in.transferTo(received), "/it/file.parquet")) {
+            executor, (in, stream) -> in.transferTo(received), 
"/it/file.parquet")) {
       out.write("abc".getBytes(StandardCharsets.UTF_8));
       out.flush();
     }
@@ -41,11 +45,42 @@ class HdfsStreamingOutputStreamTest {
     // already closed by try-with-resources; construct another and close twice
     ByteArrayOutputStream received2 = new ByteArrayOutputStream();
     HdfsStreamingOutputStream out2 =
-        HdfsStreamingOutputStream.start(executor, in -> 
in.transferTo(received2), "/it/x");
+        HdfsStreamingOutputStream.start(
+            executor, (in, stream) -> in.transferTo(received2), "/it/x");
     out2.write(1);
     out2.close();
     assertDoesNotThrow(out2::close);
     assertDoesNotThrow(out2::flush);
     executor.shutdownNow();
   }
+
+  @Test
+  void closeTimesOutAndUnblocksWhenUploaderNeverReads() throws Exception {
+    CountDownLatch block = new CountDownLatch(1);
+    var executor = Executors.newSingleThreadExecutor();
+    try {
+      HdfsStreamingOutputStream out =
+          HdfsStreamingOutputStream.start(
+              executor,
+              (in, stream) -> {
+                try {
+                  block.await();
+                } catch (InterruptedException e) {
+                  Thread.currentThread().interrupt();
+                  throw new IOException(e);
+                }
+                in.transferTo(new ByteArrayOutputStream());
+              },
+              "/it/stuck.parquet",
+              1);
+      out.write(1);
+      long started = System.nanoTime();
+      IOException error = assertThrows(IOException.class, out::close);
+      assertTrue(error.getMessage().contains("timed out"));
+      assertTrue(System.nanoTime() - started < 10_000_000_000L);
+    } finally {
+      block.countDown();
+      executor.shutdownNow();
+    }
+  }
 }
diff --git 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClientTest.java
 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClientTest.java
index ac8dd11dc4..d907178bd7 100644
--- 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClientTest.java
+++ 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/HdfsWebHdfsClientTest.java
@@ -27,8 +27,11 @@ import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.charset.StandardCharsets;
 import java.security.PrivilegedExceptionAction;
+import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 import javax.security.auth.login.LoginException;
 import org.apache.hc.client5.http.impl.classic.HttpClientBuilder;
 import org.apache.hop.core.variables.Variables;
@@ -272,6 +275,35 @@ class HdfsWebHdfsClientTest {
     }
   }
 
+  /**
+   * Serializes {@link HdfsKerberosSession#doAs} the way a process-wide 
Kerberos lock used to wrap
+   * the whole HTTP PUT. Parallel CREATEs must still overlap; if execute() is 
inside doAs this test
+   * fails.
+   */
+  static class LockingSession extends HdfsKerberosSession {
+    private final Object lock = new Object();
+
+    LockingSession(HdfsMeta meta) {
+      super(new Variables(), meta);
+    }
+
+    @Override
+    public <T> T doAs(PrivilegedExceptionAction<T> action) throws Exception {
+      synchronized (lock) {
+        return action.run();
+      }
+    }
+
+    @Override
+    public <T> T gss(PrivilegedExceptionAction<T> action) throws Exception {
+      synchronized (lock) {
+        @SuppressWarnings("unchecked")
+        T header = (T) "Negotiate dGVzdA==";
+        return header;
+      }
+    }
+  }
+
   static class TrackingSession extends HdfsKerberosSession {
     boolean doAsCalled;
 
@@ -286,6 +318,78 @@ class HdfsWebHdfsClientTest {
     }
   }
 
+  @Test
+  void fourParallelCreatesOverlapWithoutKerberos() throws Exception {
+    server.requireConcurrentPuts(4);
+    var pool = Executors.newFixedThreadPool(4);
+    try {
+      List<Future<?>> futures = new ArrayList<>();
+      for (int i = 0; i < 4; i++) {
+        int n = i;
+        futures.add(
+            pool.submit(
+                () -> {
+                  try (OutputStream out = webhdfs.create("/warehouse/p" + n + 
".parquet", true)) {
+                    out.write(new byte[4096]);
+                  }
+                  return null;
+                }));
+      }
+      for (Future<?> future : futures) {
+        future.get(15, TimeUnit.SECONDS);
+      }
+      assertEquals(4, server.maxConcurrentPuts());
+    } finally {
+      pool.shutdownNow();
+    }
+  }
+
+  @Test
+  void fourParallelCreatesOverlapWhenDoAsIsGloballyLocked() throws Exception {
+    server.requireConcurrentPuts(4);
+    HdfsMeta meta = new HdfsMeta();
+    meta.setPrincipal("[email protected]");
+    meta.setKeytabPath("/tmp/hop.keytab");
+    LockingSession session = new LockingSession(meta);
+    var executor = Executors.newCachedThreadPool();
+    var pool = Executors.newFixedThreadPool(4);
+    try {
+      HdfsWebHdfsClient client =
+          new HdfsWebHdfsClient(
+              HttpClientBuilder.create()
+                  .disableContentCompression()
+                  .disableRedirectHandling()
+                  .build(),
+              HdfsTransport.WebHDFS,
+              List.of(server.endpoint()),
+              "http",
+              "/webhdfs/v1",
+              "hop",
+              true,
+              session,
+              executor);
+      List<Future<?>> futures = new ArrayList<>();
+      for (int i = 0; i < 4; i++) {
+        int n = i;
+        futures.add(
+            pool.submit(
+                () -> {
+                  try (OutputStream out = client.create("/warehouse/k" + n + 
".parquet", true)) {
+                    out.write(new byte[4096]);
+                  }
+                  return null;
+                }));
+      }
+      for (Future<?> future : futures) {
+        future.get(15, TimeUnit.SECONDS);
+      }
+      assertEquals(4, server.maxConcurrentPuts());
+    } finally {
+      pool.shutdownNow();
+      executor.shutdownNow();
+    }
+  }
+
   @Test
   void buildUriAddsSimpleUser() throws Exception {
     String uri =
diff --git 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/WebHdfsTestServer.java
 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/WebHdfsTestServer.java
index 2457138fe2..c49609bf0e 100644
--- 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/WebHdfsTestServer.java
+++ 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/client/WebHdfsTestServer.java
@@ -29,6 +29,7 @@ import java.util.HashMap;
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 
 /** Minimal WebHDFS/HttpFS stub for unit tests. In-memory files, no Hadoop. */
@@ -40,6 +41,10 @@ public class WebHdfsTestServer {
   private volatile boolean standby;
   private volatile boolean openUses307;
   private final AtomicInteger requests = new AtomicInteger();
+  private volatile int barrierPuts;
+  private final AtomicInteger arrivedPuts = new AtomicInteger();
+  private final AtomicInteger currentPuts = new AtomicInteger();
+  private final AtomicInteger maxConcurrentPuts = new AtomicInteger();
 
   public WebHdfsTestServer() {
     dirs.put("/", true);
@@ -83,6 +88,16 @@ public class WebHdfsTestServer {
     return requests.get();
   }
 
+  /** CREATE body handlers wait until this many PUTs are in flight (0 = no 
wait). */
+  public void requireConcurrentPuts(int n) {
+    barrierPuts = n;
+    arrivedPuts.set(0);
+  }
+
+  public int maxConcurrentPuts() {
+    return maxConcurrentPuts.get();
+  }
+
   private void handle(HttpExchange exchange) throws IOException {
     requests.incrementAndGet();
     try {
@@ -188,10 +203,31 @@ public class WebHdfsTestServer {
       send(exchange, 201, "{\"Location\":\"" + location + "\"}");
       return;
     }
-    byte[] body = readAll(exchange.getRequestBody());
-    files.put(path, body);
-    parentDir(path);
-    send(exchange, 201, "");
+    int inFlight = currentPuts.incrementAndGet();
+    maxConcurrentPuts.updateAndGet(seen -> Math.max(seen, inFlight));
+    try {
+      int need = barrierPuts;
+      if (need > 0) {
+        arrivedPuts.incrementAndGet();
+        long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(8);
+        while (arrivedPuts.get() < need && System.nanoTime() < deadline) {
+          Thread.sleep(20);
+        }
+        if (arrivedPuts.get() < need) {
+          send(exchange, 500, "only " + arrivedPuts.get() + " concurrent PUTs, 
expected " + need);
+          return;
+        }
+      }
+      byte[] body = readAll(exchange.getRequestBody());
+      files.put(path, body);
+      parentDir(path);
+      send(exchange, 201, "");
+    } catch (InterruptedException e) {
+      Thread.currentThread().interrupt();
+      send(exchange, 500, "interrupted");
+    } finally {
+      currentPuts.decrementAndGet();
+    }
   }
 
   private void open(HttpExchange exchange, String path, Map<String, String> 
query)
diff --git 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSessionTest.java
 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSessionTest.java
index 33f8aa23c7..2ef64fa448 100644
--- 
a/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSessionTest.java
+++ 
b/plugins/tech/hadoop/src/test/java/org/apache/hop/vfs/hdfs/kerberos/HdfsKerberosSessionTest.java
@@ -21,6 +21,10 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotEquals;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import javax.security.auth.Subject;
 import org.apache.hop.core.variables.Variables;
 import org.apache.hop.vfs.hdfs.metadata.HdfsMeta;
 import org.junit.jupiter.api.AfterEach;
@@ -57,6 +61,50 @@ class HdfsKerberosSessionTest {
     assertEquals("C:/Users/hop/hop.keytab", session.keytabPath());
   }
 
+  @Test
+  void doAsDoesNotHoldJvmKerberosDuringAction() throws Exception {
+    HdfsMeta meta = new HdfsMeta();
+    meta.setName("cdp");
+    meta.setPrincipal("[email protected]");
+    meta.setKeytabPath("/tmp/hop.keytab");
+    HdfsKerberosSession session = new HdfsKerberosSession(new Variables(), 
meta);
+    session.useSubjectForTest(new Subject());
+    CountDownLatch inAction = new CountDownLatch(1);
+    CountDownLatch release = new CountDownLatch(1);
+    AtomicBoolean acquired = new AtomicBoolean();
+    Thread actionThread =
+        new Thread(
+            () -> {
+              try {
+                session.doAs(
+                    () -> {
+                      inAction.countDown();
+                      if (!release.await(5, TimeUnit.SECONDS)) {
+                        throw new IllegalStateException("release");
+                      }
+                      return null;
+                    });
+              } catch (Exception e) {
+                throw new RuntimeException(e);
+              }
+            });
+    actionThread.start();
+    assertTrue(inAction.await(5, TimeUnit.SECONDS));
+    Thread locker =
+        new Thread(
+            () -> {
+              synchronized (HdfsKerberosSession.JVM_KERBEROS) {
+                acquired.set(true);
+              }
+            });
+    locker.start();
+    locker.join(1000);
+    assertTrue(acquired.get(), "JVM_KERBEROS must not be held while doAs 
action runs");
+    release.countDown();
+    actionThread.join(5000);
+    session.close();
+  }
+
   @Test
   void unregisterDropsSessionFromRenewer() {
     HdfsMeta meta = new HdfsMeta();

Reply via email to