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]
