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();