morningman commented on code in PR #68710:
URL: https://github.com/apache/doris/pull/68710#discussion_r4178237847


##########
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) {

Review Comment:
   Confirmed, and it goes one layer further than this: failing the queued 
requests would not wake the reader either. `LogFetcher#pendRemoteFetches` 
registers `downloadFuture.onComplete(() -> 
logFetchBuffer.tryComplete(bucket))`, and `RemoteLogDownloadFuture#onComplete` 
is `logFileFuture.thenRun(callback)`, which skips an exceptional completion. 
Even if the callback ran, `RemotePendingFetch#toCompletedFetch` throws from 
`join()` inside `LogFetchBuffer#tryComplete` before the pending fetch is polled 
off, so it would still block the head of the bucket's queue.
   
   This is the residue the class comment describes (lines 38-39), and the trade 
this PR makes: before it, the same `Error` ran `System.exit(-1)` and took every 
query, load and compaction on the BE down with it; now one range of that log 
scanner waits, holding its scan thread and query context until BE restarts. I'm 
leaving it out of this PR:
   
   - Failing the pending reads cleanly has to happen inside fluss. From the 
plugin it means shadowing `RemoteLogDownloader` with its inner classes, 
`RemoteLogDownloadFuture`, and `LogFetchBuffer` or `RemotePendingFetch` as well 
-- three or four more `@Internal` classes copied member for member and pinned 
to 1.0.0.
   - The scanner has nothing to detect it by: the thread is named 
`DownloadRemoteLog-[<table path>]`, the same for every range of a table, so 
finding a range's own thread means reflecting through `LogScannerImpl`, 
`LogFetcher` and `RemoteLogDownloader`. A no-progress deadline in 
`BoundedLogRecords` would be a guess -- a remote segment can take minutes to 
download from slow storage -- and would not cover `PK_FULL`, which reads 
through fluss's own `KvSnapshotAndLogBatchScanner`.
   - Keeping the thread looping after the `Error` does not help: `fetchOnce` 
takes its prefetch permit before its `try`, so an `Error` from `poll` leaks it, 
and once the permits are gone the thread blocks in `acquire` for good.
   
   It will be reported to fluss together with the other ways an 
`OutOfMemoryError` leaves the client stuck (`ServerConnection` has no request 
timeout when a netty event loop dies; `RemoteLogDownloader#close` spins when it 
fails midway). The description no longer says every affected query fails: it 
lists what the error can leave behind, says why BE logs rather than exits, and 
asks for BE to be restarted after a JVM out-of-memory error, as after one 
anywhere else in BE's JVM.
   



##########
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")

Review Comment:
   Both premises hold, but this is how build.sh treats its helper modules, and 
I'd keep these two in line with the rest. `hive-udf-shade` (used only by 
`java-udf`) and `hive-apache-shade` (used only by `hadoop-hudi-scanner`) sit in 
the same unconditional block, and they are built, with their Hive dependencies 
resolved, when `--be-extension-ignore` drops their only consumer. The list is a 
complete enumeration on purpose (the comment above it says so), and 
`run-fe-ut.sh` follows the same one.
   
   What leaving a patch module in costs: an online build resolves fluss-client 
or paimon-common once and builds a jar of two or three classes that nothing 
deploys. It fails only offline, with a local repository that never had those 
artifacts. build.sh has no offline mode of its own (it passes `MVN_OPT` 
through), and the only caller of `--be-extension-ignore` in the repository, 
`build-for-release.sh`, ignores neither plugin.
   
   If `--be-extension-ignore` should drop helpers along with their consumers, 
it should do so for all of them -- `plugin-toolkit` only once fluss, jdbc and 
paimon are all ignored -- which is a build.sh change of its own rather than a 
special case for these two. The description now says these modules are built 
like the other helpers.
   



-- 
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]


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

Reply via email to