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]

Reply via email to