在 Ja va 中,直接使用 `BufferedReader.lines().forEach()` 来消费行数据,这其实是一个典型的“阻塞陷阱”。很多开发者初次接触时都会被它坑一把——因为 `forEach()` 会一口气消费整个流(Stream),而这个流是惰性求值的,它不会主动停止。这意味着,只要输入流没有关闭(比如进程没退出、socket 没断开),它就一直在那里等着,后续代码自然也就失去了执行机会。
你看到的“this never runs”输出,就是这种情况:`br.lines().forEach(...)` 在第一个流无数据时无限挂起,主线程被卡死,后面的逻辑根本轮不到执行。
所以,想要同时处理多个 InputStream 并执行其他任务,关键是理解一个基本事实:**I/O 操作本质上是阻塞的,你不可能在一个线程里通过“伪并发”绕过它。** 代码里那种“在单线程里用函数式流 API 一次性消费整个输入流”的做法,注定行不通。

**那么,正确的做法是什么?**
根本思路是:让每个 InputStream“各自为政”,彼此不干扰别人的执行路径。最直接、最可靠的方案,就是为每个流分配一个独立线程,用 `readLine()` 循环逐行读取,而不是用 `lines().forEach()` 一次性消费整个流。
具体来说,推荐使用 `ExecutorService` 来管理这些线程,配合标准的 `BufferedReader.readLine()` 循环。这样既能做到“边读边做其他事”,又能有效控制线程的生命周期,避免资源泄漏。
以下是改进后的完整示例(逻辑清晰,可直接参考):
```ja va
import ja va.io.*;
import ja va.nio.charset.StandardCharsets;
import ja va.util.concurrent.ExecutorService;
import ja va.util.concurrent.Executors;
public class ConcurrentInputStreamReader {
private static final ExecutorService executor = Executors.newCachedThreadPool();
public static void main(String[] args) {
// 假设 cmdsin 和 datain 是已打开的 InputStream(例如 Process.getInputStream())
InputStream cmdsin = ...; // e.g., process.getInputStream()
InputStream datain = ...; // e.g., another socket or pipe input
// 启动线程分别监听两个流
executor.submit(() -> readStream(cmdsin, "CMD"));
executor.submit(() -> readStream(datain, "DATA"));
// 主线程可自由执行其他逻辑(定时任务、状态检查、用户交互等)
while (!Thread.currentThread().isInterrupted()) {
System.out.println("Main thread is running other tasks...");
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
executor.shutdown();
}
private static void readStream(InputStream in, String prefix) {
try (BufferedReader reader = new BufferedReader(
new InputStreamReader(in, StandardCharsets.UTF_8))) {
String line;
while ((line = reader.readLine()) != null) {
System.out.println(prefix + ": " + line);
}
} catch (IOException e) {
System.err.println(prefix + " stream closed or error: " + e.getMessage());
}
}
}
```
几个关键要点
- `readLine()` 是逐行阻塞的,不是批量阻塞。每次只等待一行数据,这让你能更灵活地响应中断或配合超时控制。
- 使用 `ExecutorService` 统一管理线程,避免手动创建和销毁线程带来的麻烦,也防止资源泄漏。
- 需要特别提醒的是,`InputStream` 本身不原生支持超时。如果某个流长期没有数据,想要“及时放弃”,可以将其包装为 `ja va.nio.channels.Channels.newChannel(in)` 并配合 Selector(适用于 `ReadableByteChannel`)。不过,对于普通的 Process 或 Socket 输入流,更实用的做法是确保源头(如子进程)能在合适的时候正常关闭流。
- 另外,绝对不要像原代码那样,在循环内部重复创建 `InputStreamReader` / `BufferedReader`。这不仅效率极低,更危险的是,之前的流对象可能已经被消费,导致数据丢失。
说到底,Ja va 里并发读取多个 InputStream 的标准化解法,说穿了就是四个字:**一线一程**。再配合 `readLine` 循环,而不是在单线程里用函数式流 API 强撑“伪并发”。多线程不是权宜之计,而是 I/O 并发场景下合理且必要的设计模式。
本文转载于:https://www.php.cn/faq/2745414.html 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。