如何在 Java 中使用 PipedInputStream 在两个本地线程之间建立直接的二进制字节通讯管道
在某些场景下,我们需要在两个线程之间快速、简单地交换二进制数据,既不想通过网络Socket,也不想依赖文件系统。这时候,Ja va 提供的 PipedInputStream 和 PipedOutputStream 就是一套非常轻巧的解决方案——它们就像一根内存中的水管,让一个线程往一端倒数据,另一个
PipedInputStream 和 PipedOutputStream 就是一套非常轻巧的解决方案——它们就像一根内存中的水管,让一个线程往一端倒数据,另一个线程从另一端接着。
先明确几个核心要点:这套机制是单向、阻塞、基于字节的。也就是说,数据在管道里只能朝一个方向流动;写满或读空时,对应的线程会自动进入阻塞等待状态;缓冲区默认只有1024字节,当然你可以通过构造方法调整大小。最重要的是,它只适用于同一个JVM内的线程通信,无法跨进程或跨JVM。

说白了,这就是一个在生产者和消费者之间传递原始字节流的工具。用的时候需要特别注意:输入流和输出流必须配对使用,而且在使用之前必须完成连接。
核心步骤:创建并连接管道流
最稳妥的做法是显式调用 connect() 方法。原因很简单——如果连接顺序没控制好,很容易抛出 IOException 异常,报错内容往往是“Pipe not connected”。避免这种尴尬其实不复杂:
- 先分别创建
PipedInputStream和PipedOutputStream实例 - 然后调用
pipedReader.connect(pipeWriter)(注意方向:输入流连接输出流) - 确保连接发生在任何读或写操作之前
如果你觉得显式连接不够优雅,也可以使用带参数的构造方法在创建时自动完成绑定,但对于初学者,我还是推荐先用 connect(),逻辑更清晰。有钱有闲,也可以提前指定缓冲区容量,比如 new PipedInputStream(8192),别等满了才改。
典型生产者-消费者线程结构
这套管道的经典用法就是实现一个最简单的生产者-消费者模型。一个线程往 PipedOutputStream 里写字节,另一个线程从 PipedInputStream 中往外读。关键点在于:读写操作是自动同步的——写满缓冲区时写线程阻塞,读空时读线程阻塞。同步靠阻塞来实现,不需要显式加锁。这么设计,简单、直接、有效。
- 生产者线程:拿到
PipedOutputStream后,调用write(byte[])或write(int) - 消费者线程:拿到
PipedInputStream后,调用read(byte[])或read() - 默认缓冲区是1024字节,想要更高吞吐的话,直接指定更大的尺寸,比如
new PipedInputStream(8192)
重要注意事项与常见陷阱
尽管这套工具看起来挺美好,实际用起来还是有不少坑要小心绕开:
- 无法跨JVM或跨进程使用:这条是硬性限制,别想着用它来做RPC或微服务通信。
- 不支持双向通信:一对管道只能实现单向数据传输。如果需要双向通信,必须搞两套:A→B 和 B→A。
- 异常处理必须到位:如果写端提前关闭,读端的
read()会返回 -1(类似文件读完)。但如果写端因为异常挂掉但没有关闭流,读端就会一直阻塞下去,造成线程“假死”。 - 小心死锁:千万不要在同一个线程中对同一对管道既读又写。那样的话,缓冲区满了或空了就会互相等待,彻底锁死。
完整可运行示例
废话不多说,直接看代码。下面这个例子展示两个线程通过管道传递5个整数,每个整数以4字节的二进制形式传输:
PipedInputStream pis = new PipedInputStream();
PipedOutputStream pos = new PipedOutputStream();
pis.connect(pos); // 先连接,再干活
Thread writer = new Thread(() -> {
try (DataOutputStream dos = new DataOutputStream(pos)) {
for (int i = 1; i <= 5; i++) {
dos.writeInt(i * 10);
System.out.println("写入: " + (i * 10));
Thread.sleep(100);
}
} catch (Exception e) {
e.printStackTrace();
}
});
Thread reader = new Thread(() -> {
try (DataInputStream dis = new DataInputStream(pis)) {
for (int i = 0; i < 5; i++) {
int val = dis.readInt();
System.out.println("读取: " + val);
}
} catch (Exception e) {
e.printStackTrace();
}
});
writer.start();
reader.start();
运行后会依次打印出“写入: 10”、“读取: 10”、“写入: 20”、“读取: 20”……整个过程完全同步,体现了字节级的阻塞行为。代码中用到了 try-with-resources,确保流能够及时关闭、资源不会泄漏。
说真的,这种一次性、内存级的管道,在特定场景下特别好用。但你也看到了,使用门槛不低——连接、阻塞、异常、死锁……每个环节都不能马虎。只要掌握好了,它就是轻量级线程协作里的一把利刃。当然,如果你需要跨JVM、跨网络,还是老老实实用Socket或消息队列吧。


































