在 Java NIO 里,Pipe 是一个位于同一 JVM 内部的单向数据通道,由一个可写的 SinkChannel 和一个可读的 SourceChannel 组成。它非常适合用来在线程之间传递字节流,而且底层通常基于内存循环缓冲实现,不需要走网络协议栈,因此性能很高。

一、Pipe 的基本结构
Pipe 通过 Pipe.open() 创建,得到的两个通道分别负责写入和读取:
- SinkChannel:只能写,生产者线程往里面塞数据
- SourceChannel:只能读,消费者线程从里面取数据
由于是单向的,如果需要双向通信,就要创建两个 Pipe。另外,Pipe 本身线程安全程度有限,一般建议一个写线程、一个读线程,避免多写或多读造成混乱。
二、简单使用示例
下面演示一个生产者线程写数据、消费者线程读数据的完整例子:
import java.nio.ByteBuffer;
import java.nio.channels.Pipe;
public class PipeDemo {
public static void main(String[] args) throws Exception {
// 打开一个 Pipe
Pipe pipe = Pipe.open();
// 获取写通道和读通道
Pipe.SinkChannel sink = pipe.sink();
Pipe.SourceChannel source = pipe.source();
// 生产者线程
Thread writer = new Thread(() -> {
try {
ByteBuffer buf = ByteBuffer.allocate(1024);
for (int i = 0; i < 5; i++) {
buf.clear();
String msg = "msg-" + i;
buf.put(msg.getBytes());
buf.flip();
// 写入 sink 通道
while (buf.hasRemaining()) {
sink.write(buf);
}
}
// 写完关闭 sink
sink.close();
} catch (Exception e) {
e.printStackTrace();
}
});
// 消费者线程
Thread reader = new Thread(() -> {
try {
ByteBuffer buf = ByteBuffer.allocate(1024);
int len;
// 从 source 通道读取
while ((len = source.read(buf)) > 0) {
buf.flip();
byte[] data = new byte[buf.remaining()];
buf.get(data);
System.out.println("read: " + new String(data));
buf.clear();
}
source.close();
} catch (Exception e) {
e.printStackTrace();
}
});
writer.start();
reader.start();
writer.join();
reader.join();
}
}
三、和 BlockingQueue 的对比
很多开发者习惯用 BlockingQueue 做线程间通信,那 Pipe 有什么不同?
| 对比项 | Pipe | BlockingQueue |
|---|---|---|
| 通信模型 | 流式字节 | 对象或消息块 |
| 底层实现 | 内存循环缓冲 | 数组或链表加锁 |
| 适用场景 | 连续字节流 | 离散任务或数据块 |
如果你的数据是连续字节流,比如序列化后的对象、日志流,Pipe 会更自然。如果是离散任务,BlockingQueue 更简单。
四、常见注意事项
1. 缓冲区满的问题
当消费者读得慢,SinkChannel 的写入会阻塞(或在非阻塞模式下返回 0)。生产方需要处理好这种情况,避免假死。
2. 正确关闭通道
写完数据后一定要关闭 SinkChannel,否则 SourceChannel.read 会一直阻塞,以为还有数据要来。
3. 不要多写多读
单个 Pipe 最好只配对一个写线程和一个读线程。多写需要自己加同步,多读可能丢数据。
五、小结
利用 NIO 的 Pipe 可以在同一个 JVM 内部以内存级方式实现高效的单向流通信。它比传统加锁队列更贴近流式处理,也更容易和 Channel 体系里的其他组件组合。只要注意关闭通道和控制读写线程数量,就能在生产者消费者场景中稳定发挥作用。