goutamadwant opened a new pull request, #12480:
URL: https://github.com/apache/seatunnel/pull/12480
### Purpose of this pull request
Fixes #12132.
`kudu-connector-it` is often cancelled at its 90-minute job limit. In
`Build` runs on `dev`
pushes from 2026-05-30 to 2026-09-24, 114 Kudu legs were cancelled after
more than 60 minutes,
against 464 that passed (median 33 minutes). In 91 of the 92 runs where only
one JDK leg hung,
the other leg passed. The cause is a Java-level deadlock in the Flink
JobManager.
It starts when the Kudu client is created on JDK 11.
**Root cause.** Kudu 1.15 `SecurityContext.setupSubject()` looks at the
caller's `Subject`.
If that Subject has no Kerberos principal, Kudu ignores it, but first it
builds a debug message
with `subject.toString()`. That call runs even when debug logging is off.
`Subject#toString` holds the Subject's principal-set lock while it calls
`Principal#toString`.
The first time this happens in a JVM, `UnixPrincipal#toString` loads a
resource bundle. On
JDK 11 that load goes through `AccessController.getContext()` and
`SubjectDomainCombiner.combine()`, which needs the combiner lock.
Flink's JobManager and TaskManager run inside such a Subject (Hadoop login
user `flink`, no
Kerberos). Every `new Thread()` in those JVMs calls
`AccessController.getContext()`. That takes
the combiner lock first and then the principal-set lock, which is the
opposite order. JDK-8166124
reports a deadlock on the same two locks (closed as "Cannot Reproduce").
When a thread is
created at the same moment the Kudu enumerator builds its client, the two
threads deadlock. From
then on, every thread creation in the JobManager blocks on the combiner lock.
Thread dump from a local reproduction (Flink 1.18.0 image, OpenJDK 11.0.21,
`kudu_to_assert.conf`):
```
Found one Java-level deadlock:
"BLOB Server listener at 6124":
waiting to lock ... (a
javax.security.auth.SubjectDomainCombiner$WeakKeyValueMap),
which is held by "flink-rest-server-netty-boss-thread-1"
"flink-rest-server-netty-boss-thread-1":
waiting to lock ... (a java.util.Collections$SynchronizedSet),
which is held by "SourceCoordinator-Source: Kudu-Source"
"SourceCoordinator-Source: Kudu-Source":
waiting to lock ... (a
javax.security.auth.SubjectDomainCombiner$WeakKeyValueMap),
which is held by "flink-rest-server-netty-boss-thread-1"
"SourceCoordinator-Source: Kudu-Source":
at javax.security.auth.SubjectDomainCombiner.combine
at java.security.AccessController.getContext
at java.security.AccessController.doPrivileged
at
java.util.ResourceBundle$ResourceBundleProviderHelper.loadResourceBundle
...
at sun.security.util.ResourcesMgr.getAuthResourceString
at com.sun.security.auth.UnixPrincipal.toString
at javax.security.auth.Subject.toString
- locked <...> (a java.util.Collections$SynchronizedSet)
at
org.apache.kudu.client.SecurityContext.setupSubject(SecurityContext.java:162)
at
org.apache.kudu.client.AsyncKuduClient.<init>(AsyncKuduClient.java:427)
at
org.apache.seatunnel.connectors.seatunnel.kudu.util.KuduUtil.getKuduClientInternal(KuduUtil.java:164)
...
at
org.apache.seatunnel.connectors.seatunnel.kudu.source.KuduSourceSplitEnumerator.open(KuduSourceSplitEnumerator.java:101)
at
org.apache.seatunnel.translation.flink.source.FlinkSourceEnumerator.start(FlinkSourceEnumerator.java:79)
```
In the dump, the BLOB server listener is stuck creating a connection thread,
so the
TaskManager's jar download fails with `GET operation failed: Read timed
out`. The REST server's
boss thread is stuck the same way, so `flink run` never gets a job status.
In the CI logs the
JobMaster also stops answering (`Slot offering to JobManager
did not finish in time`, heartbeat timeout after 120 s), and the
ResourceManager logs `Missing
resources ... numberOfRequiredSlots=1` until the workflow is cancelled. This
matches every
sighting in #12132.
It explains the pattern in the CI logs:
- Only the JDK 11 Flink images hang (1.15.3 on 11.0.17 and 1.18.0 on
11.0.21). The 1.13.6
and 1.20.1 images run JDK 8 and never hung in these logs. I verified the
deadlock path only on
JDK 11.
- It hits the first Kudu job on a freshly started JobManager, because the
bundle is loaded once
per JVM, and each `@TestTemplate` invocation starts a new Flink cluster.
- I read 86 logs of cancelled `kudu-connector-it` legs (`dev` pushes plus
one fork, since
2026-08-20). In 85 of them the last job is on Flink 1.18.0 (73) or 1.15.3
(12), and the Kudu
enumerator's `EnumeratorOpenEvent` is never logged. In passing runs it is
logged within about
0.3 s of `Starting split enumerator`. So the enumerator is stuck while it
creates the Kudu
client. The remaining log (Flink 1.13.6) is a different failure.
**Change.** `KuduUtil` builds the Kudu client outside the caller's Subject
when that Subject has
no Kerberos principal, using `Subject.doAs(null, ...)`. Kudu makes the same
choice without a
caller Subject: it ignores such a Subject and falls back to the ticket
cache. So Kudu's
authentication behaviour does not change, and `Subject#toString` is no
longer called. If the
caller Subject has a Kerberos principal, or there is no Subject, the client
is built exactly as
before. This covers every place a Kudu client is created (source enumerator
and reader, sink,
catalog), because they all go through `KuduUtil`. The only new requirement is
`AuthPermission("doAs")` for `Subject.doAs`, and only when a SecurityManager
is installed, which
Flink and SeaTunnel do not do by default.
This is a product fix, not only a test fix. The same deadlock can hit a real
Flink 1.15+
deployment on JDK 11 that reads from or writes to Kudu without Kerberos.
### Does this PR introduce _any_ user-facing change?
No behaviour or configuration change. Kudu jobs on Flink with JDK 11 no
longer risk freezing the
JobManager when the Kudu client is first created.
### How was this patch tested?
**Unit tests.** Added two tests to `KuduClientResourceTest`:
- `shouldNotInspectCallerSubjectWithoutKerberosCredentials` creates a client
through
`KuduUtil.getKuduClientResource` inside a Subject whose only principal
counts `toString()`
calls.
- `shouldPassCallerSubjectWithKerberosPrincipalToKudu` checks that a caller
Subject holding a
`KerberosPrincipal` still reaches Kudu's `SecurityContext` unchanged. It
fails if the new code
also drops such a Subject.
For the first test:
- Before the fix: `expected: <0> but was: <1>` on JDK 8 and JDK 11.
- After the fix: both tests pass on JDK 8 and JDK 11, together with the
other Kudu connector
tests (11 tests).
```
./mvnw -B verify -DskipUT=false -DskipIT=true -pl
seatunnel-connectors-v2/connector-kudu -am -Dtest='Kudu*Test'
-DfailIfNoTests=false
```
**E2E, forced race (local only, not part of this PR).** The natural hang
rate is too low for a
short local loop: 1 hang in 40 fresh Flink 1.18 clusters. So a temporary
hook in
`KuduSourceSplitEnumerator#open` kept creating threads under the
JobManager's Subject while
the Kudu client was built. Each iteration starts a fresh Flink 1.18.0
JobManager/TaskManager and
runs `kudu_to_assert.conf` against the Kudu 1.15 containers from `KuduIT`:
- `dev`: 3 of 3 runs hung. Each thread dump shows the same deadlock between
the Kudu enumerator
thread (inside `SecurityContext.setupSubject` -> `Subject.toString`) and a
thread in
`Thread.<init>`. In one dump the Pekko RPC dispatcher is also blocked on
the combiner lock.
- this PR: 15 of 15 runs passed in 30â35 s each. During each client build
the hook created
260kâ455k threads, and no deadlock was reported.
Without the hook, `dev` hung once in 40 runs of the same loop. That dump is
the one quoted above.
**E2E, unchanged `KuduIT`.** With this fix, 29 of 29 invocations passed on
Zeta, Flink 1.18.0,
Flink 1.20.1 and Spark 3.3.0 (`TEST_IN_PR=true RUN_ALL_CONTAINER=false`,
host JDK 8, 1,883 s).
Flink 1.13.6 and 1.15.3 were not run locally.
```
./mvnw -B verify -DskipUT=true -DskipIT=false -pl :connector-kudu-e2e -am
-Dit.test=KuduIT
```
**Not verified:** a real Kerberos-enabled Kudu cluster. The Kerberos path is
unchanged,
because a Subject that holds a Kerberos principal is passed to Kudu as
before. The fork CI run of
this branch is still pending. KuduIT itself is unchanged. The hang is timing
dependent, so the
unit test pins the mechanism instead.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]