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

morningman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new dde889fba3b [fix](be) Keep fluss and paimon library threads from 
exiting the BE process (#68710)
dde889fba3b is described below

commit dde889fba3bbe08a5fba1b9d378ef0541fb61ae3
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Fri Oct 9 11:57:11 2026 +0800

    [fix](be) Keep fluss and paimon library threads from exiting the BE process 
(#68710)
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: #66399 (fluss catalog)
    
    Problem Summary:
    
    **In short.** The fluss client and paimon, which BE runs inside its
    embedded JVM for the fluss and paimon catalogs, call `System.exit` when
    one of their worker threads dies, and a full JVM heap makes those
    threads die. Inside BE that exit is a crash: a scan that should fail
    with `OutOfMemoryError` took the whole BE down instead. This PR ships
    Doris copies of the three library classes that do it, which log instead
    of exiting, packaged so that they survive the Maven build cache CI
    builds with. It does not make an out-of-memory error harmless: what the
    error leaves behind in the fluss client can keep a few scans stuck until
    BE is restarted. What it removes is a library thread, hit by chance,
    deciding to crash BE.
    
    **Background**
    
    - BE reads fluss and paimon tables through JNI plugins
    (`fe/be-java-extensions/fluss-scanner`, `paimon-scanner`) that run in
    the one JVM BE embeds (`-Xmx2048m` by default). Each plugin bundles its
    library unmodified.
    - PluginRuntime searches a plugin directory's jars in name order, except
    that a jar whose manifest carries `Doris-Shadows-Classes` is searched
    first (`jni-spi/PROTOCOL.md`). `hadoop-deps` already uses this to ship a
    Doris-patched `org.apache.hadoop.fs.FileSystem`;
    `check_plugin_layout.py` checks that such a jar holds exactly the
    classes it names.
    - `System.exit` called from a JVM thread inside BE is an `::exit()` of
    the BE process: the JVM's shutdown hooks run, then the C++ global
    destructors run while BE is still serving. `jvm_launcher.cpp` describes
    the same failure for the JVM's signal handlers, which `-Xrs` keeps away;
    `System.exit` does not go through a signal.
    - BE's JVM runs with no out-of-memory option (`JAVA_OPTS_FOR_JDK_17` in
    `conf/be.conf`): an `OutOfMemoryError` anywhere else in it (a scan
    thread, JDBC, a Java UDF, libhdfs) fails the work it hit and BE carries
    on, and what keeps the heap from running out is admission
    (`hdfs_file_writer.cpp` for HDFS writes; for JNI reads, the opt-in
    admission in #68713). Where Doris runs a JVM as a process of its own, it
    fails fast instead: FE with `-XX:OnOutOfMemoryError="kill -9 %p"`, the
    CDC client BE forks with `-XX:+ExitOnOutOfMemoryError`.
    
    **The problem, and what it cost**
    
    Three library classes end the process when a thread dies:
    
    | Class | Where it runs | What it does |
    |---|---|---|
    | fluss `org.apache.fluss.utils.FatalExitExceptionHandler` | every
    thread from fluss's `ExecutorThreadFactory`: metadata refresh, remote
    file downloader, security token renewal, lookup and write clients,
    future-timeout delayer; also `FutureUtils.assertNoException`, which
    nothing in fluss-client 1.0.0 calls | `System.exit(-17)` on any uncaught
    exception |
    | fluss `org.apache.fluss.utils.concurrent.ShutdownableThread` |
    `RemoteLogDownloader`'s thread, one per log scanner | `System.exit(-1)`
    on any `Error` from its work (Kafka's original exits only on its own
    `FatalExitError`) |
    | paimon `org.apache.paimon.utils.FatalExitExceptionHandler` | threads
    from paimon's `ExecutorThreadFactory`: `AsyncRecordReader` (merge
    reads), `ParallelExecution` | `System.exit(-17)` |
    
    An `OutOfMemoryError: Java heap space` strikes whatever thread allocates
    at that moment, including these threads while they are idle (the
    download thread allocates while it waits for work). Concurrent fluss and
    paimon reads can fill the default 2 GB heap, and when they did, BE
    aborted instead of failing the query:
    
    ```
    libc++abi: terminating due to uncaught exception of type std::system_error: 
mutex lock failed: Invalid argument
    ```
    
    thrown in BE's compaction thread on a mutex the exit had already
    destroyed. Seen twice on a Release BE with the default heap:
    
    - eight concurrent `SELECT *` union reads of a 30M-row fluss primary-key
    table (100K updated keys per bucket in the log tail, 16 scanners each):
    all 24 queries failed with `OutOfMemoryError`, as they should, and then
    BE aborted;
    - sixteen bucket reads at once over a fluss primary-key table of 2M
    keys, half of them updated since its last kv snapshot: the first query
    failed with the out-of-memory error, and seconds to minutes later BE
    aborted. A JFR `jdk.Shutdown` event named the caller:
    `ShutdownableThread.run()` -> `System.exit(-1)` on a
    `DownloadRemoteLog-[...]` thread.
    
    So one query that runs the JVM heap out takes the whole BE down, with
    every query, load and compaction running on it, and BE comes back only
    if something outside restarts it.
    
    A second, smaller defect in the same class: closing a log scanner (on a
    BE scan thread) shuts its download thread down and waits on a latch with
    no timeout. fluss logs "Starting" before the `try` whose `finally`
    counts that latch down, and under a full heap the JVM can also unwind a
    compiled frame without running its `finally` ("failed reallocation of
    scalar replaced objects"). A thread that died either way left the scan
    thread waiting for ever, holding its query's context; seen once for four
    hours after an out-of-memory run.
    
    **How this PR fixes it**
    
    Two commits, one per library:
    
    1. **fluss**: a new module `fe/be-java-extensions/fluss-client-patch`
    with Doris copies of `FatalExitExceptionHandler` (logs the exception and
    returns) and `ShutdownableThread` (an `Error` from its work is logged
    and ends the thread, what fluss already does with any other `Throwable`;
    "Starting" is logged inside the `try`; `awaitShutdown()` also returns
    once the thread is no longer alive).
    2. **paimon**: a new module `fe/be-java-extensions/paimon-common-patch`
    with a Doris copy of paimon's `FatalExitExceptionHandler` that only
    logs.
    
    The copies declare every member the originals declare, so the library
    code compiled against the originals links to them. Each module's jar
    holds only these classes and names them in `Doris-Shadows-Classes`. The
    plugin depends on the module, declared ahead of the library, so
    `copy-dependencies` puts the jar into the plugin directory, where
    PluginRuntime searches it first, and surefire runs the plugin's tests on
    the copies as well.
    
    Why a module of its own rather than a second jar built by the plugin
    module: CI builds with the Maven build cache
    (`fe/.mvn/maven-build-cache-config.xml`). A cache hit restores a
    module's own jar and nothing else, while `copy-dependencies` is forced
    to run again. A second jar that the plugin module wrote into its own
    `target/lib` was missing from every cache-restored build (checked:
    paimon-scanner's `target/lib` 155 -> 154 jars, fluss-scanner 4 -> 3), so
    the plugin would deploy with the library's original classes. As a
    dependency module it is copied on every build, cached or not.
    
    What it buys:
    
    - BE no longer goes down when an `OutOfMemoryError` lands on one of
    these library threads. The queries whose reads hit the error fail with
    it; what the error leaves behind in the fluss client is listed under Not
    in this PR, and is why BE should still be restarted after one.
    - Closing a fluss log scanner cannot hang a BE scan thread on a download
    thread that already died.
    - The fix survives CI builds that hit the Maven build cache.
    
    **Why log, rather than exit.** An out-of-memory error can leave BE's JVM
    degraded: a dedicated thread it kills is not replaced, and under a full
    heap HotSpot can unwind a frame without running its `finally`, so some
    waits never end. That is a reason to restart BE after one, not to keep
    this exit:
    
    - it is not an orderly shutdown: the shutdown hooks run first, for
    seconds to minutes under a full heap with BE still serving, then the C++
    global destructors run under BE threads that are still working, and BE
    aborts;
    - it fires only when the error happens to land on a thread these three
    classes govern; landing anywhere else in BE's JVM, the same error
    already leaves BE running;
    - it takes every query, load and compaction on the BE down with it.
    
    What remains is degraded, not corrupt: both plugins only read; a bounded
    fluss read ends on offsets, so a reader that stops delivering leaves its
    range waiting, never short; and paimon's async and parallel readers hand
    a failed task back to the scan as an error. Whether BE should go down
    after a JVM out-of-memory error is a choice for all of BE's Java code
    alike, not for a library thread: an operator who wants it can set
    `-XX:+ExitOnOutOfMemoryError` in `JAVA_OPTS_FOR_JDK_17`, which ends the
    process at once, without running shutdown hooks or destructors.
    
    Not in this PR:
    
    - keeping the heap from running out in the first place (the opt-in
    admission in #68713);
    - failing the reads a dead download thread leaves behind. Nothing
    restarts that thread, so a range of its log scanner that still needs a
    remote log segment waits for it indefinitely, and the range's BE scan
    thread, inside `getNext` where it does not see the query's cancellation,
    holds the query's context until BE restarts. Failing those reads has to
    happen in fluss: `RemoteLogDownloader` would have to fail the requests
    it holds, and `RemoteLogDownloadFuture#onComplete` runs its callback
    through `thenRun`, which skips an exceptional completion, so even a
    failed request would not wake the reader;
    - fluss client threads other than these that an `OutOfMemoryError`
    kills: a netty event loop that dies leaves the RPCs on its connection
    waiting, since fluss has no request timeout, so a few scan threads can
    stay stuck after such an error until BE restarts. This and the previous
    item have to be fixed in fluss and will be reported there;
    - FE's fluss connector, which carries the same handler: FE only runs
    admin RPCs on fluss's threads and already exits on `OutOfMemoryError`
    (`-XX:OnOutOfMemoryError`), so it is left as is.
    
    **Results**
    
    Release BE, default 2 GB JVM heap, local fluss 1.0.0 cluster:
    
    | Scenario | master | this PR |
    |---|---|---|
    | 8 concurrent `SELECT *` union reads of a 30M-row fluss primary-key
    table, 100K updated keys per bucket in the tail, 16 scanners each | all
    queries fail with `OutOfMemoryError`, then BE aborts | all queries fail
    with `OutOfMemoryError` (or the attach failure that follows it); BE
    stays up (3 rounds) |
    | 16 bucket reads at once over a fluss primary-key table with 1M of its
    2M keys updated since its kv snapshot, 6 queries in a row | the first
    fails with `OutOfMemoryError`; BE aborts seconds to minutes later | all
    6 fail with `OutOfMemoryError`; BE stays up; 5 download threads log
    `died of an error`; the queries that follow succeed with the expected
    `COUNT(*)` / `SUM(price)`; live heap back to 175 MB |
    | a download thread ends without counting its shutdown latch down, then
    its log scanner is closed | the BE scan thread waits for ever (seen for
    four hours, holding its query's context) | `awaitShutdown()` returns
    once the thread is dead |
    
    BE stays up after these runs, but is not necessarily clean: after later
    out-of-memory runs of the same kind, with the other fluss PRs applied as
    well, an idle BE still held fluss scan threads stuck in the cases listed
    under Not in this PR, with their queries' contexts, four hours on, until
    it was restarted.
    
    | Unit test | run against the library's own classes | run against the
    copies in this PR |
    |---|---|---|
    | `FlussClientThreadDeathTest` (5 tests) | the test JVM exits (status
    239 / 255), or `shutdown()` still waits after 60 s | pass |
    | `PaimonThreadDeathTest` (2 tests) | the test JVM exits (status 239) |
    pass |
    
    **Classes, and how they connect**
    
    - `org.apache.fluss.utils.FatalExitExceptionHandler` (new,
    `fluss-client-patch`): logs at ERROR and returns.
    - `org.apache.fluss.utils.concurrent.ShutdownableThread` (new,
    `fluss-client-patch`): fluss's class with `run()` and `awaitShutdown()`
    changed as above.
    - `org.apache.paimon.utils.FatalExitExceptionHandler` (new,
    `paimon-common-patch`): logs at ERROR and returns.
    - `fluss-scanner/pom.xml`, `paimon-scanner/pom.xml`: depend on the patch
    module, ahead of the library.
    - `be-java-extensions/pom.xml`, `build.sh`, `run-fe-ut.sh`: list the two
    modules (like `hive-udf-shade` and `hive-apache-shade`, they are built
    even when `--be-extension-ignore` leaves their plugin out);
    `fe/check/checkstyle/suppressions.xml`: import control is suppressed for
    the patch modules' fluss and paimon packages (import control is rooted
    at `org.apache.doris`), every other check applies.
    - PluginRuntime (`jni-bootstrap`, untouched): searches a
    `Doris-Shadows-Classes` jar first.
    
    ```
    BE process
     '- embedded JVM
         '- PluginRuntime, plugins/jni/fluss/                    
(plugins/jni/paimon/ likewise)
              |- fluss-client-patch-<v>.jar   Doris-Shadows-Classes: 
FatalExitExceptionHandler, ShutdownableThread   <- searched first
              |- fluss-client-<v>.jar         the same two classes, now never 
loaded
              '- fluss-scanner.jar            FlussJniScanner ...
    
     a fluss ExecutorThreadFactory thread dies   --> 
FatalExitExceptionHandler.uncaughtException
                                                       fluss: System.exit(-17)  
 -> BE aborts
                                                       Doris: log, the pool 
starts another thread
     RemoteLogDownloader thread throws an Error  --> ShutdownableThread.run
                                                       fluss: System.exit(-1)   
 -> BE aborts
                                                       Doris: log, the thread 
ends; nothing restarts it
     BE scan thread closes the log scanner       --> 
ShutdownableThread.awaitShutdown
                                                       Doris: returns when the 
latch is counted down or the thread is dead
    ```
    
    **Merge order.** No code dependency on other PRs. Please merge it before
    #68712, which lets every scan run up to 16 scanners per instance again
    by default: more JNI readers at once make running the JVM heap out more
    likely, and without this PR that becomes a BE crash.
---
 build.sh                                           |   2 +
 fe/be-java-extensions/fluss-client-patch/pom.xml   | 100 +++++++++++++
 .../fluss/utils/FatalExitExceptionHandler.java     |  56 +++++++
 .../fluss/utils/concurrent/ShutdownableThread.java | 165 +++++++++++++++++++++
 fe/be-java-extensions/fluss-scanner/pom.xml        |  22 +++
 .../doris/fluss/FlussClientThreadDeathTest.java    | 154 +++++++++++++++++++
 fe/be-java-extensions/paimon-common-patch/pom.xml  |  78 ++++++++++
 .../paimon/utils/FatalExitExceptionHandler.java    |  57 +++++++
 fe/be-java-extensions/paimon-scanner/pom.xml       |  22 +++
 .../apache/doris/paimon/PaimonThreadDeathTest.java |  68 +++++++++
 fe/be-java-extensions/pom.xml                      |   2 +
 fe/check/checkstyle/suppressions.xml               |   6 +
 run-fe-ut.sh                                       |   3 +-
 13 files changed, 734 insertions(+), 1 deletion(-)

diff --git a/build.sh b/build.sh
index de304892669..c377169b8da 100755
--- a/build.sh
+++ b/build.sh
@@ -884,6 +884,8 @@ if [[ "${BUILD_BE_JAVA_EXTENSIONS}" -eq 1 ]]; then
     # plugin directories. -am would reach them, but they are named here so 
this list stays a
     # complete enumeration.
     modules+=("be-java-extensions/plugin-toolkit")
+    modules+=("be-java-extensions/fluss-client-patch")
+    modules+=("be-java-extensions/paimon-common-patch")
     modules+=("be-java-extensions/hive-udf-shade")
     modules+=("be-java-extensions/hive-apache-shade")
 
diff --git a/fe/be-java-extensions/fluss-client-patch/pom.xml 
b/fe/be-java-extensions/fluss-client-patch/pom.xml
new file mode 100644
index 00000000000..233a6411943
--- /dev/null
+++ b/fe/be-java-extensions/fluss-client-patch/pom.xml
@@ -0,0 +1,100 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+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.
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0";
+    xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance";
+    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
http://maven.apache.org/xsd/maven-4.0.0.xsd";>
+    <parent>
+        <artifactId>be-java-extensions</artifactId>
+        <groupId>org.apache.doris</groupId>
+        <version>${revision}</version>
+    </parent>
+    <modelVersion>4.0.0</modelVersion>
+
+    <artifactId>fluss-client-patch</artifactId>
+    <name>Doris BE fluss client patch</name>
+    <description>
+        Doris's copies of the two fluss-client classes that end the process 
when a client thread
+        dies: org.apache.fluss.utils.FatalExitExceptionHandler and
+        org.apache.fluss.utils.concurrent.ShutdownableThread (their headers 
say why). Not a plugin:
+        the fluss plugin declares this artifact, so its jar lands in the 
plugin directory beside
+        fluss-client, and the jar's Doris-Shadows-Classes manifest entry is 
what makes PluginRuntime
+        search it ahead of fluss-client's copies.
+
+        A module of its own, not a second jar built by fluss-scanner, because 
of the Maven build
+        cache: a cache hit restores a module's own jar and nothing else, while 
copy-dependencies,
+        forced to run on a hit (fe/.mvn/maven-build-cache-config.xml), copies 
the jars of the
+        modules a plugin depends on. A second jar that fluss-scanner wrote 
into its own target/lib
+        was missing from every cache-restored build, and the plugin deployed 
without it. The
+        layout check (tools/be-java-plugins/check_plugin_layout.py) requires a 
shadowing jar to
+        hold nothing but the classes it names, which is all this module builds.
+    </description>
+
+    <properties>
+        <maven.compiler.source>8</maven.compiler.source>
+        <maven.compiler.target>8</maven.compiler.target>
+    </properties>
+
+    <dependencies>
+        <!--
+          The artifact these classes replace, provided by the plugin that 
deploys them. Compiling
+          against it also means compiling against the slf4j 1.7 API shaded 
into it, the one these
+          classes link to inside the plugin; the fe parent pom's slf4j-api 2.x 
is narrowed to test
+          and fluss-client's own slf4j-api declaration excluded, as in 
fluss-scanner.
+        -->
+        <dependency>
+            <groupId>org.apache.fluss</groupId>
+            <artifactId>fluss-client</artifactId>
+            <version>${fluss.version}</version>
+            <scope>provided</scope>
+            <exclusions>
+                <exclusion>
+                    <groupId>org.slf4j</groupId>
+                    <artifactId>slf4j-api</artifactId>
+                </exclusion>
+                <exclusion>
+                    <groupId>com.google.code.findbugs</groupId>
+                    <artifactId>jsr305</artifactId>
+                </exclusion>
+            </exclusions>
+        </dependency>
+        <dependency>
+            <groupId>org.slf4j</groupId>
+            <artifactId>slf4j-api</artifactId>
+            <scope>test</scope>
+        </dependency>
+    </dependencies>
+
+    <build>
+        <plugins>
+            <plugin>
+                <groupId>org.apache.maven.plugins</groupId>
+                <artifactId>maven-jar-plugin</artifactId>
+                <configuration>
+                    <archive>
+                        <!-- Merged with the parent's 
Doris-Jni-Plugin-Api-Version entry, not replacing it. -->
+                        <manifestEntries>
+                            
<Doris-Shadows-Classes>org.apache.fluss.utils.FatalExitExceptionHandler,org.apache.fluss.utils.concurrent.ShutdownableThread</Doris-Shadows-Classes>
+                        </manifestEntries>
+                    </archive>
+                </configuration>
+            </plugin>
+        </plugins>
+    </build>
+</project>
diff --git 
a/fe/be-java-extensions/fluss-client-patch/src/main/java/org/apache/fluss/utils/FatalExitExceptionHandler.java
 
b/fe/be-java-extensions/fluss-client-patch/src/main/java/org/apache/fluss/utils/FatalExitExceptionHandler.java
new file mode 100644
index 00000000000..98a664666b7
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-client-patch/src/main/java/org/apache/fluss/utils/FatalExitExceptionHandler.java
@@ -0,0 +1,56 @@
+// 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.fluss.utils;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Doris's copy of the handler fluss gives its client threads for uncaught 
exceptions: this one logs
+ * them, where fluss's stops the process.
+ *
+ * <p>Fluss hands {@link #INSTANCE} to every thread its {@code 
ExecutorThreadFactory} makes - the
+ * admin's metadata refresh, the remote file downloader, the security token 
renewal, the lookup and
+ * write clients, the delayer behind its future timeouts - and to futures 
{@code FutureUtils} asserts
+ * never fail. Fluss's own version calls {@code System.exit(-17)}, which suits 
a fluss server, a process
+ * of its own. Inside BE it is an {@code ::exit()} from a JVM thread: the C++ 
global destructors run
+ * under a BE that is still serving, the failure {@code jvm_launcher.cpp} 
describes for the JVM's
+ * signal handlers. Exhausting the JVM heap with concurrent union reads did 
exactly that - one of these
+ * threads ran out of memory, and BE aborted in its compaction thread on a 
destroyed mutex. Nothing a
+ * fluss client thread does is worth the BE: the thread that threw is gone, 
the pools fluss runs on
+ * start another, and the scans it was serving fail with errors of their own.
+ *
+ * <p>Found ahead of fluss-client's copy because the jar it ships in names it 
in
+ * {@code Doris-Shadows-Classes} (see this module's pom). It declares every 
member fluss's version does,
+ * so the fluss classes compiled against that one link to this one.
+ */
+public final class FatalExitExceptionHandler implements 
Thread.UncaughtExceptionHandler {
+
+    public static final FatalExitExceptionHandler INSTANCE = new 
FatalExitExceptionHandler();
+
+    /** The status fluss's version exits with. Nothing exits with it here; it 
stays for linkage. */
+    public static final int EXIT_CODE = -17;
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(FatalExitExceptionHandler.class);
+
+    @Override
+    public void uncaughtException(Thread thread, Throwable e) {
+        LOG.error("Fluss client thread '{}' died of an uncaught exception. 
Fluss would stop the process"
+                + " here; inside BE the thread is lost and the process carries 
on.", thread.getName(), e);
+    }
+}
diff --git 
a/fe/be-java-extensions/fluss-client-patch/src/main/java/org/apache/fluss/utils/concurrent/ShutdownableThread.java
 
b/fe/be-java-extensions/fluss-client-patch/src/main/java/org/apache/fluss/utils/concurrent/ShutdownableThread.java
new file mode 100644
index 00000000000..8c1cf35f51c
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-client-patch/src/main/java/org/apache/fluss/utils/concurrent/ShutdownableThread.java
@@ -0,0 +1,165 @@
+// 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.fluss.utils.concurrent;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Doris's copy of the base class of fluss's long-lived worker threads: a 
thread whose work throws an
+ * {@link Error} ends here, where fluss's ends the process.
+ *
+ * <p>Fluss's version calls {@code System.exit(-1)} when {@link #doWork()} 
throws an {@code Error} - fluss
+ * took the class from Kafka, which exits only on its own {@code 
FatalExitError}, and widened that to
+ * every {@code Error}. In fluss-client the class runs {@code 
RemoteLogDownloader}'s thread, one per log
+ * scanner, which spends its life waiting for a remote log segment to fetch, 
and the waiting allocates. A
+ * scan that filled BE's JVM heap made that wait throw {@code 
OutOfMemoryError}, and the exit ran the C++
+ * global destructors under a BE that was still serving: BE aborted in its 
compaction thread on a
+ * destroyed mutex, as it did when {@link 
org.apache.fluss.utils.FatalExitExceptionHandler} exited. Here
+ * the thread logs the error, counts itself shut down as fluss's does, and 
stops - what fluss does on
+ * every other {@code Throwable}. Nothing restarts it, so the log scanner it 
served fetches no further
+ * remote segment: a range of that scanner that still needs one waits for it 
indefinitely.
+ *
+ * <p>Closing that scanner shuts the thread down and waits for it, on a scan 
thread, so nothing may leave
+ * that wait hanging: an {@code Error} thrown as the thread starts is caught 
like one from its work (fluss
+ * logs the start outside its {@code try}), and {@link #awaitShutdown()} also 
returns once the thread has
+ * ended, since an {@code OutOfMemoryError} can end it without running the 
code that would say so.
+ *
+ * <p>Found ahead of fluss-client's copy because the jar it ships in names it 
in
+ * {@code Doris-Shadows-Classes} (see this module's pom). Apart from {@link 
#run()} and
+ * {@link #awaitShutdown()} it is fluss's class member for member, so the 
fluss classes compiled against
+ * that one link to this one.
+ */
+public abstract class ShutdownableThread extends Thread {
+
+    protected final Logger log;
+
+    private final boolean isInterruptible;
+
+    private final CountDownLatch shutdownInitiated = new CountDownLatch(1);
+    private final CountDownLatch shutdownComplete = new CountDownLatch(1);
+
+    private volatile boolean isStarted = false;
+
+    public ShutdownableThread(String name) {
+        this(name, true);
+    }
+
+    public ShutdownableThread(String name, boolean isInterruptible) {
+        super(name);
+        this.isInterruptible = isInterruptible;
+        this.log = LoggerFactory.getLogger(getClass());
+        setDaemon(false);
+    }
+
+    public void shutdown() throws InterruptedException {
+        initiateShutdown();
+        awaitShutdown();
+    }
+
+    public boolean isShutdownInitiated() {
+        return shutdownInitiated.getCount() == 0;
+    }
+
+    /**
+     * Asks the thread to stop after its current unit of work, interrupting it 
if it was made
+     * interruptible; returns whether this call was the one that asked.
+     */
+    public boolean initiateShutdown() {
+        synchronized (this) {
+            if (isRunning()) {
+                log.info("Shutting down");
+                shutdownInitiated.countDown();
+                if (isInterruptible) {
+                    interrupt();
+                }
+                return true;
+            }
+            return false;
+        }
+    }
+
+    /**
+     * After calling {@link #initiateShutdown()}, waits for the thread to 
finish its work, or to have ended
+     * without saying so: fluss waits on the latch alone, and a thread can end 
without counting it down.
+     * Under a full heap the JVM may unwind a compiled frame without running 
its {@code catch} and
+     * {@code finally} blocks - when deoptimizing it cannot reallocate the 
frame's scalar-replaced objects,
+     * it throws "OutOfMemoryError: Java heap space: failed reallocation of 
scalar replaced objects" past
+     * them - and the close of the log scanner the thread served, on a scan 
thread, would never return.
+     */
+    public void awaitShutdown() throws InterruptedException {
+        if (!isShutdownInitiated()) {
+            throw new IllegalStateException("initiateShutdown() was not called 
before awaitShutdown()");
+        }
+        if (isStarted) {
+            while (!shutdownComplete.await(1, TimeUnit.SECONDS)) {
+                if (!isAlive()) {
+                    break;
+                }
+            }
+        }
+        log.info("Shutdown completed");
+    }
+
+    /** One unit of the thread's work, run over and over until the thread is 
shut down. */
+    public abstract void doWork() throws Exception;
+
+    @Override
+    public void run() {
+        isStarted = true;
+        try {
+            // Inside the try, where fluss logs it before: logging allocates, 
and an Error thrown there
+            // would end the thread without counting shutdownComplete down, so 
that awaitShutdown() - the
+            // close of the log scanner this thread serves, on a scan thread - 
would wait for ever.
+            log.info("Starting");
+            while (isRunning()) {
+                doWork();
+            }
+        } catch (Error e) {
+            shutdownInitiated.countDown();
+            shutdownComplete.countDown();
+            // Fluss's version calls System.exit(-1) here; see the class 
comment.
+            log.error("Fluss client thread '{}' died of an error. Fluss would 
stop the process here; inside BE "
+                    + "the thread is lost and the process carries on.", 
getName(), e);
+        } catch (Throwable e) {
+            if (isRunning()) {
+                log.error("Error due to", e);
+            }
+        } finally {
+            shutdownComplete.countDown();
+        }
+        log.info("Stopped");
+    }
+
+    /**
+     * Waits for {@code timeout}, or until shutdown is initiated, whichever 
comes first: a pause that a
+     * shutdown cuts short.
+     */
+    protected void pause(long timeout, TimeUnit unit) throws 
InterruptedException {
+        if (shutdownInitiated.await(timeout, unit)) {
+            log.trace("shutdownInitiated latch count reached zero. Shutdown 
called.");
+        }
+    }
+
+    public boolean isRunning() {
+        return !isShutdownInitiated();
+    }
+}
diff --git a/fe/be-java-extensions/fluss-scanner/pom.xml 
b/fe/be-java-extensions/fluss-scanner/pom.xml
index 1db5a0da8d9..de20e406b88 100644
--- a/fe/be-java-extensions/fluss-scanner/pom.xml
+++ b/fe/be-java-extensions/fluss-scanner/pom.xml
@@ -64,6 +64,28 @@ under the License.
             <version>${project.version}</version>
         </dependency>
 
+        <!--
+          Doris's copies of the fluss-client classes that would end BE when a 
client thread dies
+          (FatalExitExceptionHandler, ShutdownableThread; see that module). 
copy-dependencies puts
+          its jar in this plugin's directory, where its Doris-Shadows-Classes 
manifest entry makes
+          PluginRuntime search it ahead of fluss-client. Declared BEFORE 
fluss-client because
+          surefire's classpath follows declaration order, so the tests run on 
the copies the
+          deployed plugin runs on - FlussClientThreadDeathTest depends on 
that. Every transitive is
+          excluded: the module's own dependency is provided, so all that could 
come along are the
+          logging artifacts the fe parent pom declares on every module.
+        -->
+        <dependency>
+            <groupId>org.apache.doris</groupId>
+            <artifactId>fluss-client-patch</artifactId>
+            <version>${project.version}</version>
+            <exclusions>
+                <exclusion>
+                    <groupId>*</groupId>
+                    <artifactId>*</artifactId>
+                </exclusion>
+            </exclusions>
+        </dependency>
+
         <!--
           The whole fluss client, as one ~69MB unrelocated fat jar: 
fluss-common, fluss-rpc,
           frocksdbjni (with its natives, for the kv-snapshot path), a 
relocated Arrow
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussClientThreadDeathTest.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussClientThreadDeathTest.java
new file mode 100644
index 00000000000..6f0297c28a2
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussClientThreadDeathTest.java
@@ -0,0 +1,154 @@
+// 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.doris.fluss;
+
+import org.apache.fluss.utils.FatalExitExceptionHandler;
+import org.apache.fluss.utils.concurrent.ExecutorThreadFactory;
+import org.apache.fluss.utils.concurrent.FutureUtils;
+import org.apache.fluss.utils.concurrent.ShutdownableThread;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.slf4j.Logger;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Proxy;
+import java.time.Duration;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * A fluss client thread that dies must not take the process with it. Fluss's 
own handler for these
+ * threads, and the base class of its remote log download threads, call {@code 
System.exit}, and inside
+ * BE that runs the C++ global destructors under a BE that is still serving; 
the plugin ships classes
+ * of the same names that only log.
+ *
+ * <p>Should one of fluss's come back - a patched class dropped, or found 
after fluss-client's - these
+ * tests do not fail with an assertion: the JVM running them exits, and 
surefire reports a forked VM
+ * that terminated without saying goodbye.
+ */
+public class FlussClientThreadDeathTest {
+
+    @Test
+    public void clientThreadThatDiesOfAnUncaughtErrorLeavesTheProcessRunning() 
throws Exception {
+        // The factory behind fluss's admin refresh, remote file download and 
token renewal threads.
+        Thread thread = new 
ExecutorThreadFactory("doris-fatal-exit-test").newThread(() -> {
+            throw new OutOfMemoryError("simulated: Java heap space");
+        });
+        Assertions.assertSame(FatalExitExceptionHandler.INSTANCE, 
thread.getUncaughtExceptionHandler(),
+                "fluss no longer hands its threads this handler; the patch may 
not be needed any more");
+
+        thread.start();
+        thread.join(TimeUnit.SECONDS.toMillis(60));
+        Assertions.assertFalse(thread.isAlive(), "the thread was meant to 
die");
+    }
+
+    @Test
+    public void failedFutureFlussAssertsNeverFailsLeavesTheProcessRunning() {
+        CompletableFuture<Void> future = new CompletableFuture<>();
+        FutureUtils.assertNoException(future);
+        future.completeExceptionally(new IllegalStateException("simulated"));
+        Assertions.assertTrue(future.isCompletedExceptionally());
+    }
+
+    @Test
+    public void downloadThreadThatDiesOfAnErrorLeavesTheProcessRunning() 
throws Exception {
+        Assertions.assertSame(ShutdownableThread.class,
+                
Class.forName("org.apache.fluss.client.table.scanner.log.RemoteLogDownloader$DownloadRemoteLogThread")
+                        .getSuperclass(),
+                "fluss no longer runs its remote log download on this class; 
the patch may not be needed any more");
+        ShutdownableThread thread = new 
ShutdownableThread("doris-fatal-exit-test") {
+            @Override
+            public void doWork() {
+                throw new OutOfMemoryError("simulated: Java heap space");
+            }
+        };
+
+        thread.start();
+        thread.join(TimeUnit.SECONDS.toMillis(60));
+        Assertions.assertFalse(thread.isAlive(), "the thread was meant to 
die");
+        // What closing its log scanner does with it afterwards: must return, 
not wait on the dead thread.
+        Assertions.assertTimeoutPreemptively(Duration.ofSeconds(60), 
thread::shutdown);
+    }
+
+    /**
+     * The heap was full as the thread started, and logging that it starts is 
the first thing it
+     * allocates. The thread is lost all the same; the close of its log 
scanner, which shuts it down on a
+     * scan thread, must still return instead of waiting for a thread that 
never said it stopped.
+     */
+    @Test
+    public void downloadThreadThatDiesAsItStartsLetsItsScannerClose() throws 
Exception {
+        ShutdownableThread thread = new 
ShutdownableThread("doris-fatal-exit-test") {
+            @Override
+            public void doWork() throws InterruptedException {
+                pause(1, TimeUnit.DAYS);
+            }
+        };
+        Field log = ShutdownableThread.class.getDeclaredField("log");
+        log.setAccessible(true);
+        log.set(thread, runningOutOfHeapOn("Starting", (Logger) 
log.get(thread)));
+
+        thread.start();
+        thread.join(TimeUnit.SECONDS.toMillis(60));
+        Assertions.assertFalse(thread.isAlive(), "the thread was meant to 
die");
+        Assertions.assertTimeoutPreemptively(Duration.ofSeconds(60), 
thread::shutdown);
+    }
+
+    /**
+     * Under a full heap the JVM can end a thread without running its {@code 
finally} blocks: a compiled
+     * frame whose scalar-replaced objects cannot be reallocated on 
deoptimization is unwound past them. The
+     * thread then never counts itself shut down, and closing its log scanner 
must return all the same.
+     */
+    @Test
+    public void downloadThreadThatEndsWithoutSayingSoLetsItsScannerClose() 
throws Exception {
+        ShutdownableThread thread = new 
ShutdownableThread("doris-fatal-exit-test") {
+            @Override
+            public void doWork() {
+            }
+
+            @Override
+            public void run() {
+                // What the JVM leaves when it unwinds ShutdownableThread#run 
past its finally: a thread
+                // that started and ended, and a shutdownComplete nobody 
counted down.
+            }
+        };
+        Field started = ShutdownableThread.class.getDeclaredField("isStarted");
+        started.setAccessible(true);
+        started.set(thread, true);
+
+        thread.start();
+        thread.join(TimeUnit.SECONDS.toMillis(60));
+        Assertions.assertFalse(thread.isAlive(), "the thread was meant to 
end");
+        Assertions.assertTimeoutPreemptively(Duration.ofSeconds(60), 
thread::shutdown);
+    }
+
+    /** {@code logger}, except that logging {@code message} at info runs out 
of heap. */
+    private static Logger runningOutOfHeapOn(String message, Logger logger) {
+        return (Logger) Proxy.newProxyInstance(Logger.class.getClassLoader(), 
new Class<?>[] {Logger.class},
+                (proxy, method, args) -> {
+                    if (method.getName().equals("info") && args != null && 
message.equals(args[0])) {
+                        throw new OutOfMemoryError("simulated: Java heap 
space");
+                    }
+                    try {
+                        return method.invoke(logger, args);
+                    } catch (InvocationTargetException e) {
+                        throw e.getCause();
+                    }
+                });
+    }
+}
diff --git a/fe/be-java-extensions/paimon-common-patch/pom.xml 
b/fe/be-java-extensions/paimon-common-patch/pom.xml
new file mode 100644
index 00000000000..b2646255285
--- /dev/null
+++ b/fe/be-java-extensions/paimon-common-patch/pom.xml
@@ -0,0 +1,78 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+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.
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0";
+    xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance";
+    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
http://maven.apache.org/xsd/maven-4.0.0.xsd";>
+    <parent>
+        <artifactId>be-java-extensions</artifactId>
+        <groupId>org.apache.doris</groupId>
+        <version>${revision}</version>
+    </parent>
+    <modelVersion>4.0.0</modelVersion>
+
+    <artifactId>paimon-common-patch</artifactId>
+    <name>Doris BE paimon common patch</name>
+    <description>
+        Doris's copy of org.apache.paimon.utils.FatalExitExceptionHandler, the 
handler paimon gives
+        its pool threads, which ends the process when one of them dies (the 
class header says why).
+        Not a plugin: the paimon plugin declares this artifact, so its jar 
lands in the plugin
+        directory, and the jar's Doris-Shadows-Classes manifest entry is what 
makes PluginRuntime
+        search it ahead of paimon-common's copy and of the one 
paimon-hive-connector bundles.
+
+        A module of its own, not a second jar built by paimon-scanner, because 
of the Maven build
+        cache: a cache hit restores a module's own jar and nothing else, while 
copy-dependencies,
+        forced to run on a hit (fe/.mvn/maven-build-cache-config.xml), copies 
the jars of the
+        modules a plugin depends on. A second jar that paimon-scanner wrote 
into its own target/lib
+        was missing from every cache-restored build, and the plugin deployed 
without it. The
+        layout check (tools/be-java-plugins/check_plugin_layout.py) requires a 
shadowing jar to
+        hold nothing but the classes it names, which is all this module builds.
+    </description>
+
+    <properties>
+        <maven.compiler.source>8</maven.compiler.source>
+        <maven.compiler.target>8</maven.compiler.target>
+    </properties>
+
+    <dependencies>
+        <!-- The artifact this class replaces, provided by the plugin that 
deploys it. -->
+        <dependency>
+            <groupId>org.apache.paimon</groupId>
+            <artifactId>paimon-common</artifactId>
+            <scope>provided</scope>
+        </dependency>
+    </dependencies>
+
+    <build>
+        <plugins>
+            <plugin>
+                <groupId>org.apache.maven.plugins</groupId>
+                <artifactId>maven-jar-plugin</artifactId>
+                <configuration>
+                    <archive>
+                        <!-- Merged with the parent's 
Doris-Jni-Plugin-Api-Version entry, not replacing it. -->
+                        <manifestEntries>
+                            
<Doris-Shadows-Classes>org.apache.paimon.utils.FatalExitExceptionHandler</Doris-Shadows-Classes>
+                        </manifestEntries>
+                    </archive>
+                </configuration>
+            </plugin>
+        </plugins>
+    </build>
+</project>
diff --git 
a/fe/be-java-extensions/paimon-common-patch/src/main/java/org/apache/paimon/utils/FatalExitExceptionHandler.java
 
b/fe/be-java-extensions/paimon-common-patch/src/main/java/org/apache/paimon/utils/FatalExitExceptionHandler.java
new file mode 100644
index 00000000000..6393bf2eb7d
--- /dev/null
+++ 
b/fe/be-java-extensions/paimon-common-patch/src/main/java/org/apache/paimon/utils/FatalExitExceptionHandler.java
@@ -0,0 +1,57 @@
+// 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.paimon.utils;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Doris's copy of the handler paimon gives its threads for uncaught 
exceptions: this one logs them,
+ * where paimon's stops the process.
+ *
+ * <p>Paimon hands {@link #INSTANCE} to every thread its {@code 
ExecutorThreadFactory} makes unless the
+ * caller names another handler - among them the pool behind {@code 
AsyncRecordReader}, which reads
+ * each large data file of a merge read on a thread of its own, and the pools 
of
+ * {@code ParallelExecution}. Paimon's own version calls {@code 
System.exit(-17)}, which suits the Flink
+ * or Spark worker paimon is written for. Inside BE it is an {@code ::exit()} 
from a JVM thread: the C++
+ * global destructors run under a BE that is still serving, the failure {@code 
jvm_launcher.cpp}
+ * describes for the JVM's signal handlers, and BE aborts. The fluss client's 
handler of the same
+ * name did exactly that when concurrent reads exhausted the JVM heap, and 
paimon's merge reads are
+ * the readers that exhaust it most easily. A pool thread is cheap to lose: a 
task's own exception goes
+ * to whoever waits on its future, so this handler only sees what kills a 
thread between tasks - an
+ * {@code OutOfMemoryError} while an idle thread waits for work - and its pool 
starts another.
+ *
+ * <p>Found ahead of paimon-common's copy, and of the one 
paimon-hive-connector bundles, because the
+ * jar it ships in names it in {@code Doris-Shadows-Classes} (see this 
module's pom). It declares every
+ * member paimon's version does, so the paimon classes compiled against that 
one link to this one.
+ */
+public final class FatalExitExceptionHandler implements 
Thread.UncaughtExceptionHandler {
+
+    public static final FatalExitExceptionHandler INSTANCE = new 
FatalExitExceptionHandler();
+
+    /** The status paimon's version exits with. Nothing exits with it here; it 
stays for linkage. */
+    public static final int EXIT_CODE = -17;
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(FatalExitExceptionHandler.class);
+
+    @Override
+    public void uncaughtException(Thread thread, Throwable e) {
+        LOG.error("Paimon thread '{}' died of an uncaught exception. Paimon 
would stop the process"
+                + " here; inside BE the thread is lost and the process carries 
on.", thread.getName(), e);
+    }
+}
diff --git a/fe/be-java-extensions/paimon-scanner/pom.xml 
b/fe/be-java-extensions/paimon-scanner/pom.xml
index 01cfdb75cff..b69e15eb3c0 100644
--- a/fe/be-java-extensions/paimon-scanner/pom.xml
+++ b/fe/be-java-extensions/paimon-scanner/pom.xml
@@ -66,6 +66,28 @@ under the License.
             <version>${project.version}</version>
         </dependency>
 
+        <!--
+          Doris's copy of the paimon handler that would end BE when a pool 
thread dies
+          (FatalExitExceptionHandler; see that module). copy-dependencies puts 
its jar in this
+          plugin's directory, where its Doris-Shadows-Classes manifest entry 
makes PluginRuntime
+          search it ahead of paimon-common and of the copy 
paimon-hive-connector bundles. Declared
+          BEFORE the paimon artifacts because surefire's classpath follows 
declaration order, so the
+          tests run on the copy the deployed plugin runs on - 
PaimonThreadDeathTest depends on that.
+          Every transitive is excluded: the module's own dependency is 
provided, so all that could
+          come along are the logging artifacts the fe parent pom declares on 
every module.
+        -->
+        <dependency>
+            <groupId>org.apache.doris</groupId>
+            <artifactId>paimon-common-patch</artifactId>
+            <version>${project.version}</version>
+            <exclusions>
+                <exclusion>
+                    <groupId>*</groupId>
+                    <artifactId>*</artifactId>
+                </exclusion>
+            </exclusions>
+        </dependency>
+
         <dependency>
             <groupId>org.apache.paimon</groupId>
             <artifactId>paimon-core</artifactId>
diff --git 
a/fe/be-java-extensions/paimon-scanner/src/test/java/org/apache/doris/paimon/PaimonThreadDeathTest.java
 
b/fe/be-java-extensions/paimon-scanner/src/test/java/org/apache/doris/paimon/PaimonThreadDeathTest.java
new file mode 100644
index 00000000000..ee5a5fc5962
--- /dev/null
+++ 
b/fe/be-java-extensions/paimon-scanner/src/test/java/org/apache/doris/paimon/PaimonThreadDeathTest.java
@@ -0,0 +1,68 @@
+// 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.doris.paimon;
+
+import org.apache.paimon.utils.ExecutorThreadFactory;
+import org.apache.paimon.utils.FatalExitExceptionHandler;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * A paimon pool thread that dies must not take the process with it. Paimon's 
own handler for these
+ * threads calls {@code System.exit}, and inside BE that runs the C++ global 
destructors under a BE
+ * that is still serving; the plugin ships a handler of the same name that 
only logs.
+ *
+ * <p>Should paimon's handler come back - the patched class dropped, or found 
after paimon-common's -
+ * these tests do not fail with an assertion: the JVM running them exits, and 
surefire reports a
+ * forked VM that terminated without saying goodbye.
+ */
+public class PaimonThreadDeathTest {
+
+    @Test
+    public void threadThatDiesOfAnUncaughtErrorLeavesTheProcessRunning() 
throws Exception {
+        // The factory behind the pools of AsyncRecordReader and 
ParallelExecution.
+        Thread thread = new 
ExecutorThreadFactory("doris-fatal-exit-test").newThread(() -> {
+            throw new OutOfMemoryError("simulated: Java heap space");
+        });
+        Assertions.assertSame(FatalExitExceptionHandler.INSTANCE, 
thread.getUncaughtExceptionHandler(),
+                "paimon no longer hands its threads this handler; the patch 
may not be needed any more");
+
+        thread.start();
+        thread.join(TimeUnit.SECONDS.toMillis(60));
+        Assertions.assertFalse(thread.isAlive(), "the thread was meant to 
die");
+    }
+
+    @Test
+    public void poolKeepsServingAfterOneOfItsThreadsDies() throws Exception {
+        ExecutorService pool = Executors.newCachedThreadPool(new 
ExecutorThreadFactory("doris-fatal-exit-pool"));
+        try {
+            // Thrown outside any future, the way an OutOfMemoryError hits a 
worker waiting for work:
+            // it reaches the thread's uncaught exception handler and ends 
that worker.
+            pool.execute(() -> {
+                throw new OutOfMemoryError("simulated: Java heap space");
+            });
+            Assertions.assertEquals("served", pool.submit(() -> 
"served").get(60, TimeUnit.SECONDS));
+        } finally {
+            pool.shutdownNow();
+        }
+    }
+}
diff --git a/fe/be-java-extensions/pom.xml b/fe/be-java-extensions/pom.xml
index 529986b5f51..2176d691df3 100644
--- a/fe/be-java-extensions/pom.xml
+++ b/fe/be-java-extensions/pom.xml
@@ -24,6 +24,8 @@ under the License.
         <module>jni-spi</module>
         <module>jni-bootstrap</module>
         <module>plugin-toolkit</module>
+        <module>fluss-client-patch</module>
+        <module>paimon-common-patch</module>
         <module>iceberg-metadata-scanner</module>
         <module>hadoop-hudi-scanner</module>
         <module>hive-udf-shade</module>
diff --git a/fe/check/checkstyle/suppressions.xml 
b/fe/check/checkstyle/suppressions.xml
index 151b120d6ae..906a511b4be 100644
--- a/fe/check/checkstyle/suppressions.xml
+++ b/fe/check/checkstyle/suppressions.xml
@@ -52,6 +52,12 @@ under the License.
     <suppress 
files="org[\\/]apache[\\/]doris[\\/](?!nereids)[^\\/]+[\\/]|DorisFE\.java" 
checks="EmptyLineSeparator" id="forNereids" />
 
     <!-- exclude rules for special files -->
+    <!-- fluss-client-patch's replacements for fluss classes 
(FatalExitExceptionHandler, ShutdownableThread) have to
+         live in fluss's packages, which import-control.xml (rooted at 
org.apache.doris) does not cover; every other
+         check still applies -->
+    <suppress 
files="fluss-client-patch[\\/]src[\\/]main[\\/]java[\\/]org[\\/]apache[\\/]fluss[\\/]"
 checks="ImportControl" />
+    <!-- paimon-common-patch's replacement for paimon's 
FatalExitExceptionHandler, likewise in paimon's package -->
+    <suppress 
files="paimon-common-patch[\\/]src[\\/]main[\\/]java[\\/]org[\\/]apache[\\/]paimon[\\/]"
 checks="ImportControl" />
     <suppress 
files="org[\\/]apache[\\/]doris[\\/]load[\\/]loadv2[\\/]dpp[\\/]ColumnParser\.java"
 checks="OneTopLevelClass" />
     <suppress 
files="org[\\/]apache[\\/]doris[\\/]load[\\/]loadv2[\\/]dpp[\\/]SparkRDDAggregator\.java"
 checks="OneTopLevelClass" />
     <suppress 
files="org[\\/]apache[\\/]doris[\\/]catalog[\\/]FunctionSet\.java" 
checks="LineLength" />
diff --git a/run-fe-ut.sh b/run-fe-ut.sh
index a67d327ce50..9ffe1f5bcaf 100755
--- a/run-fe-ut.sh
+++ b/run-fe-ut.sh
@@ -219,7 +219,8 @@ FE_MODULES+=("fe-connector/fe-connector-fluss")
 # Listed one by one rather than as the aggregator: -pl on an aggregator 
selects that pom and none
 # of its children. Keep in sync with the module list in build.sh, which is the 
other complete
 # enumeration of this directory.
-for be_java_extension in jni-spi jni-bootstrap plugin-toolkit 
hive-apache-shade hive-udf-shade \
+for be_java_extension in jni-spi jni-bootstrap plugin-toolkit 
fluss-client-patch paimon-common-patch \
+    hive-apache-shade hive-udf-shade \
     hadoop-deps iceberg-metadata-scanner hadoop-hudi-scanner java-udf 
jdbc-scanner paimon-scanner \
     fluss-scanner max-compute-connector trino-connector-scanner java-writer; do
     FE_MODULES+=("be-java-extensions/${be_java_extension}")


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to