在 Windows 操作系统下使用 JDK 8 进行开发时,PipedReader 和 PipedWriter 是 Java IO 体系中用于
实现线程间字符数据通信的重要组件。它们构成了一个单向的、基于内存的字符管道,允许一个线程向管
道写入字符,另一个线程从管道读取字符,而无需借助临时文件或网络套接字。
以下是对这两个类的详细使用方法、核心机制及源码层面的深度分析。
一、 核心概念与使用场景
定义:PipedWriter 是管道字符输出流,PipedReader 是管道字符输入流。二者必须配套使用,形成“生产者-消费者”模型。
通信方向:单向通信。数据从 PipedWriter 流入,从 PipedReader 流出。
适用场景:
两个线程之间需要传递大量字符数据(如文本处理、日志中转)。
希望解耦生产者和消费者的处理逻辑,利用管道的缓冲区平衡读写速度差异。
局限性:
仅支持一对一通信(一个写端对应一个读端)。
读写操作是阻塞的,若缓冲区满或空,线程会挂起。
在现代高并发场景中,通常推荐使用 java.util.concurrent.BlockingQueue 替代,因为管道流的异常处理较复杂且容易死锁。
二、 基本使用方法
使用 PipedReader 和 PipedWriter 的关键步骤是建立连接。连接可以通过构造函数完成,也可以通过 connect() 方法显式完成。
1. 标准代码示例
java
import java.io.IOException;
import java.io.PipedReader;
import java.io.PipedWriter;
public class PipeExample {
public static void main(String[] args) {
// 1. 创建管道流对象
PipedReader reader = new PipedReader();
PipedWriter writer = new PipedWriter();
try {
// 2. 建立连接 (必须在读写操作开始前完成)
// 方式A: 显式连接
reader.connect(writer);
// 方式B: 也可以在构造时连接: new PipedWriter(reader);
// 3. 启动写线程 (生产者)
Thread writerThread = new Thread(() -> {
try {
System.out.println("Writer: Starting to write...");
String message = "Hello, PipedReader in JDK8!";
for (char c : message.toCharArray()) {
writer.write(c);
// 模拟慢速写入,观察阻塞效果
Thread.sleep(100);
}
writer.close(); // 重要:关闭写端,通知读端数据结束
System.out.println("Writer: Finished and closed.");
} catch (IOException | InterruptedException e) {
e.printStackTrace();
}
});
// 4. 启动读线程 (消费者)
Thread readerThread = new Thread(() -> {
try {
System.out.println("Reader: Starting to read...");
int ch;
// read() 是阻塞方法,直到有数据可读或流关闭
while ((ch = reader.read()) != -1) {
System.out.print((char) ch);
}
System.out.println("\nReader: End of stream reached (-1).");
reader.close();
} catch (IOException e) {
e.printStackTrace();
}
});
// 5. 启动线程
// 注意:建议先启动读线程或同时启动,避免写线程填满缓冲区后阻塞且无读者消费
readerThread.start();
writerThread.start();
// 等待线程结束
writerThread.join();
readerThread.join();
} catch (IOException | InterruptedException e) {
e.printStackTrace();
}
}
}
2. 关键注意事项
连接唯一性:connect() 只能调用一次。重复连接或连接已连接的流会抛出 IOException: Already connected。
关闭顺序:写端必须调用 close(),否则读端的 read() 会一直阻塞等待新数据,无法返回 -1 退出循环。
死锁风险:如果在同一个线程中先执行写操作填满缓冲区,再执行读操作,会导致死锁。因为写操作阻塞等待读操作释放空间,而读操作尚未开始。务必在不同线程中执行读写。
异常传播:如果写线程异常退出而未正常关闭,读线程在下一次 read() 时会抛出 IOException: Pipe closed。
三、 源码深度分析 (JDK 8)
理解源码有助于掌握其底层同步机制和缓冲区管理。
1. 内部结构
PipedWriter 和 PipedReader 通过共享一个缓冲区进行通信。
PipedReader:持有缓冲区 char[] buffer,以及指针 in (写指针/下一个写入位置) 和 out (读指针/下一个读取位置)。它还维护一个 boolean connected 状态。
PipedWriter:持有一个指向关联 PipedReader 的引用 private PipedReader sink。它本身不存储数据,所有写操作都委托给 sink。
2. 连接机制 (connect)
在 PipedWriter.connect(PipedReader snk) 中:
java
public synchronized void connect(PipedReader snk) throws IOException {
if (snk == null) {
throw new NullPointerException();
} else if (sink != null || snk.connected) {
throw new IOException("Already connected");
} else if (snk.closedByReader || closed) {
throw new IOException("Pipe closed");
}
sink = snk;
snk.in = -1; // 初始化写指针
snk.out = 0; // 初始化读指针
snk.connected = true; // 标记已连接
}
使用 synchronized 保证线程安全。
检查状态防止重复连接或连接已关闭的流。
初始化 PipedReader 的内部指针,准备接收数据。
3. 写入机制 (write)
PipedWriter.write(int c) 实际上调用了 sink.receive(c):
java
public void write(int c) throws IOException {
if (sink == null) {
throw new IOException("Pipe not connected");
}
sink.receive(c);
}
核心逻辑在 PipedReader.receive(int c) 中:
java
protected synchronized void receive(int c) throws IOException {
if (!connected) {
throw new IOException("Pipe not connected");
} else if (closedByWriter || closedByReader) {
throw new IOException("Pipe closed");
}
// 如果缓冲区满 (in == out 且 in != -1 表示非初始状态,或者环形缓冲区满)
// JDK8 中判断满的逻辑较为复杂,简而言之:当 (in + 1) % buffer.length == out 时视为满
while (in == out && in != -1) {
// 缓冲区满,写线程等待
notifyAll(); // 唤醒可能正在等待读取的线程(虽然通常是读线程等写,但这里为了对称)
try {
wait(); // 写线程进入等待状态,释放锁,直到被读线程唤醒
} catch (InterruptedException ex) {
throw new java.io.InterruptedIOException();
}
}
if (in == -1) {
in = 0; // 初始状态,设置写指针为0
}
buffer[in] = (char) c; // 写入数据
in = (in + 1) % buffer.length; // 移动写指针,环形覆盖
notifyAll(); // 唤醒正在等待数据的读线程
}
同步控制:整个 receive 方法是 synchronized 的,确保对缓冲区和指针的操作原子性。
阻塞策略:当缓冲区满时,写线程调用 wait() 进入等待池,释放锁。
环形缓冲区:使用取模运算 (in + 1) % buffer.length 实现环形队列,复用内存空间。
通知机制:写入数据后,调用 notifyAll() 唤醒因缓冲区为空而等待的读线程。
4. 读取机制 (read)
PipedReader.read() 的核心逻辑:
java
public synchronized int read() throws IOException {
if (!connected) {
throw new IOException("Pipe not connected");
} else if (closedByWriter && in == -1) {
// 写端关闭且无剩余数据
return -1;
}
// 如果缓冲区空 (in == -1 或 in == out)
while (in == -1 || in == out) {
if (closedByWriter) {
return -1; // 写端已关闭,返回结束标志
}
notifyAll(); // 唤醒可能等待写入的写线程
try {
wait(); // 读线程等待,直到有数据写入或写端关闭
} catch (InterruptedException ex) {
throw new java.io.InterruptedIOException();
}
}
char c = buffer[out]; // 读取数据
out = (out + 1) % buffer.length; // 移动读指针
notifyAll(); // 唤醒因缓冲区满而等待的写线程
return c;
}
阻塞策略:当缓冲区为空时,读线程调用 wait() 等待。
结束判断:如果 closedByWriter 为真且缓冲区无数据,返回 -1。
通知机制:读取数据后,腾出了空间,调用 notifyAll() 唤醒等待写入的写线程。
5. 缓冲区大小
默认缓冲区大小为 1024 个字符。可以在构造 PipedReader 时指定:
java
public PipedReader(PipedWriter src, int pipeSize) throws IOException
较大的缓冲区可以减少写线程阻塞的频率,但会增加内存占用。
四、 常见问题与最佳实践
Broken Pipe 异常:
如果读线程关闭了 PipedReader,而写线程继续写入,写线程会抛出 IOException: Pipe closed。
反之,如果写线程异常终止未关闭,读线程下次读取时会抛出 IOException: Pipe closed。
对策:始终在 finally 块中关闭流,并妥善处理 IOException。
性能考量:
由于每次读写都涉及 synchronized 锁和可能的 wait/notify 上下文切换,PipedWriter/PipedReader 的性能低于直接内存操作或 BlockingQueue。
对于高频、小数据量的通信,开销较大。
替代方案推荐:
ArrayBlockingQueue<Character> 或 LinkedBlockingQueue<String>:更灵活,支持多生产者多消费者,API 更友好,非阻塞选项丰富。
Exchanger:适用于两个线程交换数据块的场景。
BlockingQueue + StringBuilder:如果需要传递字符串,建议将字符组装成字符串后放入队列,减少同步次数。
五、 总结
在 JDK 8 Windows 环境下,PipedReader 和 PipedWriter 提供了标准的线程间字符管道通信机制。其核心在于共享环形缓冲区和基于监视器锁(Monitor)的等待/通知机制。虽然功能强大且符合 IO 流的设计哲学,但由于其严格的同步阻塞特性和一对一流的限制,在现代 Java 并发编程中,除非必须兼容旧的 IO 流接口,否则通常建议优先使用 java.util.concurrent 包中的并发容器来实现线程间数据传递。