From 7e2108c8c479225347df2b92e31ac0ba7f767bf0 Mon Sep 17 00:00:00 2001 From: jaysunxiao Date: Sun, 15 Mar 2026 09:50:58 +0800 Subject: [PATCH] fix[command]: std thread returns without waiting for the execution to complete. --- .../java/com/zfoo/monitor/util/OSUtils.java | 33 +++++-------------- .../com/zfoo/protocol/util/ThreadUtils.java | 21 +++++++++--- 2 files changed, 25 insertions(+), 29 deletions(-) diff --git a/monitor/src/main/java/com/zfoo/monitor/util/OSUtils.java b/monitor/src/main/java/com/zfoo/monitor/util/OSUtils.java index bf820d55..28f54c10 100644 --- a/monitor/src/main/java/com/zfoo/monitor/util/OSUtils.java +++ b/monitor/src/main/java/com/zfoo/monitor/util/OSUtils.java @@ -33,10 +33,7 @@ import java.text.NumberFormat; import java.util.ArrayList; import java.util.HashMap; import java.util.List; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; /** @@ -275,7 +272,7 @@ public abstract class OSUtils { logger.info("execCommand [{}]", command); try { return doExecCommand(command, null, 5 * TimeUtils.MILLIS_PER_MINUTE); - } catch (IOException | InterruptedException e) { + } catch (IOException | InterruptedException | ExecutionException e) { throw new RuntimeException(e); } } @@ -290,37 +287,21 @@ public abstract class OSUtils { var wd = new File(workingDirectory); try { return doExecCommand(command, wd, timeoutMillis); - } catch (IOException | InterruptedException e) { + } catch (IOException | InterruptedException | ExecutionException e) { throw new RuntimeException(e); } } - private static String doExecCommand(String command, File wd, long timeoutMillis) throws IOException, InterruptedException { + private static String doExecCommand(String command, File wd, long timeoutMillis) throws IOException, InterruptedException, ExecutionException { var commandSplits = command.split(StringUtils.SPACE_REGEX); var process = new ProcessBuilder(commandSplits) .redirectErrorStream(true) .directory(wd) .start(); - var stdout = new StringBuilder(); - var stderr = new StringBuilder(); - // 异步读取输出,避免缓冲区阻塞 - executors.submit(ThreadUtils.safeRunnable(() -> { - try { - stdout.append(StringUtils.bytesToString(IOUtils.toByteArray(process.getInputStream()))); - } catch (IOException e) { - throw new RuntimeException(e); - } - })); - executors.submit(ThreadUtils.safeRunnable(() -> { - try { - stderr.append(StringUtils.bytesToString(IOUtils.toByteArray(process.getErrorStream()))); - } catch (IOException e) { - throw new RuntimeException(e); - } - - })); + var stdoutFuture = executors.submit(ThreadUtils.safeCallable(() -> StringUtils.bytesToString(IOUtils.toByteArray(process.getInputStream())))); + var stderrFuture = executors.submit(ThreadUtils.safeCallable(() -> StringUtils.bytesToString(IOUtils.toByteArray(process.getErrorStream())))); var finished = process.waitFor(timeoutMillis, TimeUnit.MILLISECONDS); if (!finished) { @@ -329,6 +310,8 @@ public abstract class OSUtils { } process.destroy(); + var stdout = stdoutFuture.get(); + var stderr = stderrFuture.get(); // 获取线程的退出值,0代表正常退出,非0代表异常中止 int exitValue = process.exitValue(); diff --git a/protocol/src/main/java/com/zfoo/protocol/util/ThreadUtils.java b/protocol/src/main/java/com/zfoo/protocol/util/ThreadUtils.java index 525f838a..d743dac5 100644 --- a/protocol/src/main/java/com/zfoo/protocol/util/ThreadUtils.java +++ b/protocol/src/main/java/com/zfoo/protocol/util/ThreadUtils.java @@ -18,10 +18,7 @@ import io.netty.util.concurrent.EventExecutorGroup; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.ForkJoinPool; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; /** * @author godotg @@ -108,6 +105,22 @@ public abstract class ThreadUtils { }; } + public static Callable safeCallable(Callable callable) { + return new Callable() { + @Override + public T call() { + try { + return callable.call(); + } catch (Exception e) { + logger.error("unknown exception", e); + } catch (Throwable t) { + logger.error("unknown error", t); + } + return null; + } + }; + } + // ----------------------------------------------------------------------------------------------------------------- // threadId -> (Thread, Executor) private static final CopyOnWriteHashMapLongObject> threadExecutorMap = new CopyOnWriteHashMapLongObject<>(Runtime.getRuntime().availableProcessors() * 8);