18、PipedReader和PipedWriter的源码分析和使用方法详细分析(windows操作系统,JDK8)
18、PipedReader和PipedWriter的源码分析和使用方法详细分析(windows操作系统,JDK8)
大家好,我是你们的技术博主。今天我们来聊聊Java I/O中一个有趣但容易被忽略的兄弟组合:PipedReader和PipedWriter。它们就像管道工一样,在字符流的世界里搭建起一条数据通道,让线程之间能够高效地“传纸条”。本文会从源码分析到实战应用,带你彻底搞懂这对好搭档。## 什么是PipedReader和PipedWriter?在Java中,PipedReader和PipedWriter是字符流(Reader/Writer)的管道实现。它们允许一个线程向管道写入字符,另一个线程从管道读取字符,实现线程间的数据交换。想象一下,你有一个生产者和一个消费者,生产者往管道里放数据,消费者从管道里取数据,双方不需要直接接触,管道就是它们的桥梁。与字节流的PipedInputStream和PipedOutputStream类似,但这对兄弟专门处理字符数据,避免了编码转换问题。它们常用于多线程协作,比如一个线程生成日志,另一个线程读取并写入文件。## 源码分析:揭开管道的神秘面纱在JDK8中,PipedReader和PipedWriter的源码位于java.io包中。我们先看看它们的核心设计。### PipedWriter的源码关键点PipedWriter的构造函数需要绑定一个PipedReader作为接收端。核心方法是write(int c)和write(char[] cbuf, int off, int len)。它的内部维护了一个缓冲区(实际上是通过PipedReader的缓冲区实现的)。java// 简化版PipedWriter源码public class PipedWriter extends Writer { private PipedReader sink; // 目标PipedReader public PipedWriter(PipedReader snk) throws IOException { connect(snk); } public void connect(PipedReader snk) throws IOException { if (snk == null) { throw new NullPointerException(); } else if (sink != null || snk.connected) { throw new IOException("Already connected"); } sink = snk; snk.connected = true; } public void write(int c) throws IOException { if (sink == null) { throw new IOException("Pipe not connected"); } sink.receive(c); // 委托给PipedReader的receive方法 } // 其他方法类似...}关键点:PipedWriter不直接管理缓冲区,而是通过PipedReader的receive()方法写入数据。这体现了“写入端”和“读取端”的紧密耦合。### PipedReader的源码关键点PipedReader是管道的接收端,它维护了一个循环缓冲区(char buffer[]),以及读写指针(in和out)。当缓冲区满时,写入线程会阻塞;当缓冲区空时,读取线程会阻塞。这是通过wait()和notifyAll()实现的。java// 简化版PipedReader源码public class PipedReader extends Reader { boolean closedByWriter = false; boolean closedByReader = false; boolean connected = false; Thread readSide; // 当前读取线程 Thread writeSide; // 当前写入线程 private static final int DEFAULT_PIPE_SIZE = 1024; char buffer[]; // 循环缓冲区 int in = -1; // 写入指针 int out = 0; // 读取指针 public PipedReader(PipedWriter src) throws IOException { this(src, DEFAULT_PIPE_SIZE); } public PipedReader(PipedWriter src, int pipeSize) throws IOException { initPipe(pipeSize); connect(src); } private void initPipe(int pipeSize) { if (pipeSize <= 0) throw new IllegalArgumentException("Pipe size <= 0"); buffer = new char[pipeSize]; } // 核心方法:从Writer接收字符 synchronized void receive(int c) throws IOException { if (!connected) throw new IOException("Pipe not connected"); if (closedByWriter || closedByReader) throw new IOException("Pipe closed"); if (readSide != null && !readSide.isAlive()) throw new IOException("Read end dead"); writeSide = Thread.currentThread(); while (in == out) { // 缓冲区满时阻塞 if (readSide != null && !readSide.isAlive()) throw new IOException("Pipe broken"); notifyAll(); // 通知读取线程 try { wait(1000); } catch (InterruptedException e) { throw new InterruptedIOException(); } } if (in < 0) in = 0; buffer[in++] = (char) c; if (in >= buffer.length) in = 0; } // 核心方法:读取字符 public synchronized int read() throws IOException { if (!connected) throw new IOException("Pipe not connected"); if (closedByReader) throw new IOException("Pipe closed"); if (writeSide != null && !writeSide.isAlive() && closedByWriter) throw new IOException("Write end dead"); readSide = Thread.currentThread(); int trials = 2; while (in < 0) { // 缓冲区空时阻塞 if (closedByWriter) return -1; // 写入端关闭,返回-1 if (writeSide != null && !writeSide.isAlive() && --trials < 0) throw new IOException("Write end dead"); notifyAll(); try { wait(1000); } catch (InterruptedException e) { throw new InterruptedIOException(); } } int ret = buffer[out++]; if (out >= buffer.length) out = 0; if (out == in) { // 缓冲区空时重置 in = -1; out = 0; } return ret; }}关键点:- 循环缓冲区:in和out指针在数组范围内循环,节省空间。- 线程安全:所有方法都使用synchronized,并在缓冲区满/空时通过wait()和notifyAll()进行线程间协调。- 死线程检测:如果写入线程死亡,读取线程会收到异常或返回-1,防止死锁。## 使用方法:实战演练理论讲完,我们来看两个可运行的例子,感受一下管道的魅力。### 示例1:基础管道通信这个例子创建了两个线程:一个写入数据,一个读取数据。写入线程发送"Hello, PipedReader!",读取线程将其打印。javaimport java.io.*;public class PipeExample1 { public static void main(String[] args) throws IOException { // 创建PipedReader和PipedWriter,并连接 PipedReader reader = new PipedReader(); PipedWriter writer = new PipedWriter(reader); // 或者 reader.connect(writer) // 写入线程 Thread writerThread = new Thread(() -> { try (writer) { // 使用try-with-resources自动关闭 writer.write("Hello, PipedReader!"); System.out.println("写入线程:数据已写入"); } catch (IOException e) { e.printStackTrace(); } }); // 读取线程 Thread readerThread = new Thread(() -> { try (reader) { int data; StringBuilder sb = new StringBuilder(); while ((data = reader.read()) != -1) { sb.append((char) data); } System.out.println("读取线程:收到数据 - " + sb.toString()); } catch (IOException e) { e.printStackTrace(); } }); // 启动线程 writerThread.start(); readerThread.start(); }}运行结果:写入线程:数据已写入读取线程:收到数据 - Hello, PipedReader!说明:- PipedWriter和PipedReader通过构造函数连接。- 使用try-with-resources自动关闭流,避免资源泄漏。- 读取线程循环读取直到返回-1(表示写入端关闭)。### 示例2:多线程数据交换模拟这个例子模拟生产者-消费者模式:生产者生成1到10的数字,消费者读取并累加。注意缓冲区大小设为5,演示阻塞效果。javaimport java.io.*;public class PipeExample2 { public static void main(String[] args) throws IOException, InterruptedException { // 创建缓冲区大小为5的管道 PipedReader reader = new PipedReader(5); PipedWriter writer = new PipedWriter(reader); // 生产者线程:写入1到10 Thread producer = new Thread(() -> { try (writer) { for (int i = 1; i <= 10; i++) { writer.write(i); // 写入单个字符(实际上是ASCII码) System.out.println("生产者写入: " + i + " (字符: " + (char)(i + '0') + ")"); Thread.sleep(200); // 模拟生产间隔 } System.out.println("生产者完成"); } catch (IOException | InterruptedException e) { e.printStackTrace(); } }); // 消费者线程:读取并累加 Thread consumer = new Thread(() -> { try (reader) { int sum = 0; int data; while ((data = reader.read()) != -1) { // 注意:这里写入的是int值1-10,不是字符'1' sum += data; System.out.println("消费者读取: " + data + ",当前累加: " + sum); Thread.sleep(300); // 模拟消费间隔 } System.out.println("最终累加和: " + sum); } catch (IOException | InterruptedException e) { e.printStackTrace(); } }); // 启动线程 producer.start(); consumer.start(); // 等待线程结束 producer.join(); consumer.join(); System.out.println("主线程结束"); }}运行结果示例:生产者写入: 1 (字符: 1)消费者读取: 1,当前累加: 1生产者写入: 2 (字符: 2)生产者写入: 3 (字符: 3)消费者读取: 2,当前累加: 3生产者写入: 4 (字符: 4)生产者写入: 5 (字符: 5)消费者读取: 3,当前累加: 6...(继续直到10)最终累加和: 55说明:- 缓冲区大小为5,当生产者写入速度超过消费者时,会阻塞等待。- writer.write(i)写入的是int值(1-10),不是字符’1’。如果需要写入字符,应使用writer.write(String.valueOf(i))。- 由于生产者比消费者快(200ms vs 300ms),可以看到缓冲区满时生产者阻塞。## 注意事项和最佳实践1. 连接顺序:必须先连接再使用,否则抛出IOException。2. 线程安全:PipedReader和PipedWriter是线程安全的,但需要确保读写在不同线程中执行。如果在同一线程中读写,可能导致死锁(因为write()和read()都持有锁)。3. 缓冲区大小:默认大小为1024字符,可以根据数据量调整。太小会导致频繁阻塞,太大浪费内存。4. 异常处理:当写入线程意外死亡时,读取线程会收到IOException,需要妥善处理。5. 关闭流:使用try-with-resources自动关闭,或确保在finally中关闭。关闭写入端后,读取端会收到-1。## 总结PipedReader和PipedWriter是Java I/O库中一对精巧的线程间通信工具。它们通过循环缓冲区和synchronized+wait/notifyAll机制,实现了字符数据的安全传递。源码分析让我们看到:PipedWriter将数据委托给PipedReader的缓冲区,而PipedReader负责阻塞和唤醒逻辑。在实际开发中,不建议在生产环境中大量使用管道,因为它的性能和扩展性有限(例如,无法跨JVM通信)。但在学习和简单多线程场景下,它比直接使用BlockingQueue更贴近I/O流的概念。理解它的原理,有助于你掌握Java的线程同步和I/O设计模式。希望这篇文章能帮你彻底搞懂这对管道兄弟!如果你有任何问题,欢迎在评论区交流。下期我们聊聊BufferedReader和BufferedWriter的缓冲机制,敬请期待!
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐


所有评论(0)