PipedReader和PipedWriter的源码分析和使用方法详细分析(windows操作系统,JDK8)
允许一批多个线程向PipedWriter写入数据另一批多个线程从PipedReader读取数据。但是同一批多个线程相互之间会存在竞争比如同一批向PipedWriter写入数据的线程会存在竞争同一批从PipedReader读取数据的线程也会存在竞争。因此PipedWriter和PipedReader中的线程安全需要通过synchronized关键字和wait()/notifyAll()机制实现。不建议在一个线程中同时使用PipedWriter和PipedReader因为这样可能会导致这个线程陷入死锁状态。PipedWriter和PipedReader之间的通信本质上也是一个生产者-消费者模型其中PipedWriter作为生产者PipedReader作为消费者。两者通过一个循环缓冲区char[]数组进行数据交换PipedWriter将数据缓存在PipedReader的数组当中等待PipedReader的读取。PipedWriter和PipedReader的UML关系图如下所示image一、PipedWriter生产者源码——向PipedReader消费者中的缓冲区char[]数组写入字符数据的字符输出流生产者package java.io;public class PipedWriter extends Writer {//与这个PipedWriter生产者相关联的 PipedReader 消费者private PipedReader sink;//标记当前这个PipedWriter对象是否关闭true表示关闭false表示开启private boolean closed false;//构造函数 public PipedWriter(PipedReader snk) throws IOException { connect(snk); } //构造函数 public PipedWriter() { } //线程同步函数用来改变将要关联的PipedReader 消费者中一些变量的值 public synchronized void connect(PipedReader snk) throws IOException { if (snk null) { throw new NullPointerException();//如果将要关联的PipedReader 消费者为null抛出NullPointerException } else if (sink ! null || snk.connected) { //如果与这个PipedWriter生产者相关联的 PipedReader 消费者!null或者将要关联的PipedReader 消费者的boolean connected变量为true则抛出IOException throw new IOException(Already connected); } else if (snk.closedByReader || closed) { //如果将要关联的PipedReader 消费者的boolean closedByReader变量为true或者当前这个PipedWriter对象已经关闭则抛出IOException throw new IOException(Pipe closed); } sink snk;//将这个PipedWriter生产者与这个PipedReader 消费者相关联 snk.in -1;//改变PipedReader 消费者中的变量int in-1 snk.out 0;//改变PipedReader 消费者中的变量int out0 snk.connected true;//改变PipedReader 消费者中的变量boolean connectedtrue表示该PipedReader 消费者已经关联了某个PipedWriter生产者了 } //向与这个PipedWriter生产者相关联的 PipedReader 消费者的缓冲区char[]数组写入1个字符 public void write(int c) throws IOException { if (sink null) { //如果与这个PipedWriter生产者相关联的 PipedReader 消费者 null抛出IOException throw new IOException(Pipe not connected); } sink.receive(c);//最终调用的是这个相关联的 PipedReader 消费者的receive(int b)函数 } //向与这个PipedWriter生产者相关联的 PipedReader 消费者的缓冲区char[]数组写入char[]数组cbuf的[off,offlen)左闭右开不包括offlen索引位置的字符 public void write(char cbuf[], int off, int len) throws IOException { if (sink null) { //如果与这个PipedWriter生产者相关联的 PipedReader消费者 null抛出IOException throw new IOException(Pipe not connected); } else if ((off | len | (off len) | (cbuf.length - (off len))) 0) {//char[]数组cbuf的[off,offlen)左闭右开索引位置是否有越界的检查 throw new IndexOutOfBoundsException();//越界的话抛出一个IndexOutOfBoundsException } //最终调用的是这个相关联的 PipedReader 消费者的receive(byte b[], int off, int len)函数 sink.receive(cbuf, off, len); } //线程同步函数使用notifyAll()函数唤醒所有与这个PipedWriter生产者相关联的 PipedReader 消费者线程这个消费者可以绑定1~n个线程 public synchronized void flush() throws IOException { if (sink ! null) { if (sink.closedByReader || closed) { throw new IOException(Pipe closed); } synchronized (sink) { sink.notifyAll(); } } } //关闭这个PipedWriter生产者这个PipedWriter生产者不能再向与它相关联的PipedReader消费者中的缓冲区char[]数组写入字符数据 public void close() throws IOException { closed true; if (sink ! null) { sink.receivedLast(); } }}二、PipedReader消费者源码——从自己的缓冲区char[]数组读取字符数据的字符输入流消费者package java.io;public class PipedReader extends Reader {//标记符true表示与这个 PipedReader 消费者相关联的PipedWriter生产者已经关闭反之反之boolean closedByWriter false;//标记符true表示当前这个 PipedReader 消费者已经关闭了反之反之boolean closedByReader false;//标记符true表示与这个 PipedReader 消费者相关联的PipedWriter生产者已经持有了这个PipedReader 消费者对象或者叫已经连接上了反之反之boolean connected false;Thread readSide;//当前消费的线程 Thread writeSide;//当前生产者的线程 //默认的PipedReader 消费者的缓冲区char[]数组的长度 private static final int DEFAULT_PIPE_SIZE 1024; //PipedReader 消费者的缓冲区char[]数组 char buffer[]; //缓冲区char[]数组的写指针 int in -1; //缓冲区char[]数组的读指针 int out 0; //构造函数 public PipedReader(PipedWriter src) throws IOException { this(src, DEFAULT_PIPE_SIZE);//缓冲区char[]数组的长度使用默认值1024 } //构造函数 public PipedReader(PipedWriter src, int pipeSize) throws IOException { initPipe(pipeSize);//缓冲区char[]数组的长度使用指定的长度 //最终还是调用PipedWriter生产者的connect()函数并把自身对象this传递进去然后在PipedWriter生产者的connect()函数中改变自己的3个变量int in-1、int out0、boolean connectedtrue connect(src); } //构造函数缓冲区char[]数组的长度使用默认值1024 public PipedReader() { initPipe(DEFAULT_PIPE_SIZE); } //构造函数缓冲区char[]数组的长度使用指定的长度 public PipedReader(int pipeSize) { initPipe(pipeSize); } //初始化缓冲区char[]数组 private void initPipe(int pipeSize) { if (pipeSize 0) { throw new IllegalArgumentException(Pipe size 0); } buffer new char[pipeSize]; } public void connect(PipedWriter src) throws IOException { src.connect(this); //最终还是调用PipedWriter生产者的connect()函数并把自身对象this传递进去然后在PipedWriter生产者的connect()函数中改变自己的3个变量int in-1、int out0、boolean connectedtrue } //线程同步函数该函数只被PipedWriter生产者的write(int b)函数调用 synchronized void receive(int c) throws IOException { //检查PipedReader 消费者的状态 if (!connected) { throw new IOException(Pipe not connected); } else if (closedByWriter || closedByReader) { throw new IOException(Pipe closed); } else if (readSide ! null !readSide.isAlive()) { throw new IOException(Read end dead); } writeSide Thread.currentThread();//当前执行该函数的线程就是生产者线程 while (in out) { //如果缓冲区char[]数组的读指针缓冲区char[]数组的写指针唤醒所有消费者线程自己这个生产者线程调用wait(1000)函数 if ((readSide ! null) !readSide.isAlive()) { throw new IOException(Pipe broken); } /* full: kick any waiting readers */ notifyAll(); try { wait(1000); } catch (InterruptedException ex) { throw new java.io.InterruptedIOException(); } } if (in 0) { //缓冲区char[]数组的写指针0时设置缓冲区char[]数组的写指针0缓冲区char[]数组的读指针0 in 0; out 0; } buffer[in] (char) c;//向缓冲区的写指针位置写入1个字节 if (in buffer.length) { in 0;//如果缓冲区满了设置缓冲区的写指针0 } } //线程同步函数该函数只被PipedWriter生产者的write(char cbuf[], int off, int len)函数调用 synchronized void receive(char c[], int off, int len) throws IOException { //如果缓冲区足够大可以装下len个字符的话不停地向缓冲区char[]数组中顺序写入char[]数组c的[off,offlen)左闭右开不包括offlen索引位置的字符 while (--len 0) { receive(c[off]); } } //关闭与这个 PipedReader 消费者相关联的PipedWriter生产者 synchronized void receivedLast() { closedByWriter true; notifyAll();//唤醒所有消费者线程 } public synchronized int read() throws IOException { if (!connected) {//检查标记符connected如果为false抛出IOException throw new IOException(Pipe not connected); } else if (closedByReader) {//检查标记符closedByReader如果为true抛出IOException throw new IOException(Pipe closed); } else if (writeSide ! null !writeSide.isAlive() !closedByWriter (in 0)) { //检查当前这个PipedReader 消费者对象中引用的生产者线程和生产者线程的状态如果和标记符closedByWriter还有缓冲区char[]数组的写指针in不能对应的话抛出一个IOException throw new IOException(Write end dead); } readSide Thread.currentThread();//当前执行该函数的线程就是消费者线程 int trials 2;//这是一个多次检测的策略变量防止生产者线程没有关闭了与这个 PipedReader 消费者相关联的PipedWriter生产者时便抛出IOException //in-1的情况有3种 //①、生产者线程还没有向缓冲区char[]数组中写任何字符 //②、消费者线程从缓冲区char[]数组中读完字符char数据以后读指针out写指针in那么当前消费者线程会设置写指针in-1 //③、消费者线程执行PipedReader 的close()函数后关闭了这个PipedReader消费者 while (in 0) { if (closedByWriter) { /* closed by writer, return EOF */ return -1; } if ((writeSide ! null) (!writeSide.isAlive()) (--trials 0)) { //多个消费者线程从缓冲区char[]数组中读的时候并且前一个消费者线程已经把缓冲区char[]数组中写入的字符读完了并且前一个线程设置了写指针in-1生产者线程也关闭了与这个 PipedReader 消费者相关联的PipedWriter生产者时抛出一个IOException throw new IOException(Pipe broken); } /* might be a writer waiting */ notifyAll();//此处的目的是为了唤醒所有生产者线程 try { wait(1000); } catch (InterruptedException ex) { throw new java.io.InterruptedIOException(); } } int ret buffer[out];//获取缓冲区char[]数组中读指针out索引位置的字符并且将读指针out1 if (out buffer.length) { out 0;//如果读指针out缓冲区char[]数组的长度设置读指针out0 } if (in out) { /* now empty */ in -1;//如果消费者线程从缓冲区char[]数组中读完字符char数据以后读指针out写指针in那么当前消费者线程会设置写指针in-1 } return ret; } //线程同步函数如果缓冲区char[]数组中有足够多的字符的话数量len消费者线程每次从缓冲区char[]数组中读取len个字符放到char[]数组cbuf的[off, offlen)索引位置左闭右开不包括offlen //如果缓冲区char[]数组中字符的数量len个比如有in写指针-out读指针个消费者线程每次从缓冲区char[]数组中读取in-out个字符放到char[]数组cbuf的[off, offin-out)索引位置左闭右开不包括offin-out public synchronized int read(char cbuf[], int off, int len) throws IOException { if (!connected) {//检查标记符connected如果为false抛出IOException throw new IOException(Pipe not connected); } else if (closedByReader) {//检查标记符closedByReader如果为true抛出IOException throw new IOException(Pipe closed); } else if (writeSide ! null !writeSide.isAlive() !closedByWriter (in 0)) { //检查当前这个PipedReader 消费者对象中引用的生产者线程和生产者线程的状态如果和标记符closedByWriter还有缓冲区char[]数组的写指针in不能对应的话抛出一个IOException throw new IOException(Write end dead); } if ((off 0) || (off cbuf.length) || (len 0) || ((off len) cbuf.length) || ((off len) 0)) {//char[]数组cbuf的[off,offlen)左闭右开索引位置是否有越界的检查 throw new IndexOutOfBoundsException();//越界的话抛出一个IndexOutOfBoundsException } else if (len 0) { return 0;//如果len0返回0 } /* possibly wait on the first character */ int c read();//先调用read()函数试探性从缓冲区char[]数组中读1个字符 if (c 0) { return -1;//如果试探性的从缓冲区char[]数组中都读不到1个字符返回-1 } cbuf[off] (char)c;//把试探性从缓冲区char[]数组中读到的第1个字符放到char[]数组b的off索引位置 int rlen 1;//累计从缓冲区char[]数组中读到的所有字符数量 while ((in 0) (--len 0)) { //从缓冲区char[]数组中向char[]数组cbuf中读取字符 cbuf[off rlen] buffer[out]; rlen;//累计从缓冲区char[]数组中读到的所有字符数量 if (out buffer.length) { out 0;//如果读指针out缓冲区char[]数组的长度设置读指针out0 } if (in out) { /* now empty */ in -1;//如果消费者线程从缓冲区char[]数组中读完字符char数据以后读指针out写指针in那么当前消费者线程会设置写指针in-1 } } return rlen;//返回累计从缓冲区char[]数组中读到的所有字符数量 } public synchronized boolean ready() throws IOException { if (!connected) { throw new IOException(Pipe not connected); } else if (closedByReader) { throw new IOException(Pipe closed); } else if (writeSide ! null !writeSide.isAlive() !closedByWriter (in 0)) { throw new IOException(Write end dead); } if (in 0) { return false; } else { return true; } } //关闭这个PipedReader消费者其实就是设置标记符closedByReadertrue 设置写指针in-1 public void close() throws IOException { in -1; closedByReader true; }}三、1个线程向PipedWriter生产者写字符数据1个线程从PipedReader消费者读取字符数据的过程3.1、非循环直接写和非循环直接读整个过程和字节管道流PipedInputStream消费者和PipedOutputStream生产者的1个线程向PipedOutputStream生产者写字节数据1个线程从PipedInputStream消费者读取字节数据的过程相同唯一不同的是字符管道流PipedWriter生产者和PipedReader消费者写入和读取的缓冲区是字符数组char[] buffer字节管道流PipedInputStream消费者和PipedOutputStream生产者写入和读取字节数组缓冲区byte[] buffer详细过程请参考我的另一篇博客9、PipedInputStream和PipedOutputStream的源码分析和使用方法详细分析3.2、加锁循环写和非加锁循环读到byte[]数组b中再处理同3.1。四、多个大于等于2个线程向PipedWriter生产者写字符数据多个线程从PipedReader消费者读取字符数据的过程4.1、循环直接写和循环直接读package com.chelong.pipe.PipedReaderAndPipedWriter;import java.io.IOException;import java.io.PipedReader;import java.io.PipedWriter;/**Created by chelong on 2026/2/12*/public class PipeTest {public static void main(String[] args) throws IOException, InterruptedException {final PipedWriter writer new PipedWriter();final PipedReader reader new PipedReader(writer);Thread consumer1 new Thread(new Runnable() { Override public void run() { try { char[] readChars new char[11]; while (reader.read(readChars) ! -1) {//read()函数是阻塞的 System.out.println(Thread.currentThread().getName() read new String(readChars)); } } catch (IOException e) { } } }, consumer1-Thread); Thread consumer2 new Thread(new Runnable() { Override public void run() { try { char[] readChars new char[11]; while (reader.read(readChars) ! -1) {//read()函数是阻塞的 System.out.println(Thread.currentThread().getName() read new String(readChars)); } } catch (IOException e) { } } }, consumer2-Thread); Thread producer1 new Thread(new Runnable() { Override public void run() { try { for (int i 1; i 3; i) { writer.write(Hello Pipe i); System.out.println(Thread.currentThread().getName() write Hello Pipe i); } writer.close(); } catch (IOException e) { e.printStackTrace(); } } }, producer1-Thread); Thread producer2 new Thread(new Runnable() { Override public void run() { try { for (int i 4; i 6; i) { writer.write(Hello Pipe i); System.out.println(Thread.currentThread().getName() write Hello Pipe i); } writer.close(); } catch (IOException e) { e.printStackTrace(); } } }, producer2-Thread); consumer1.start(); consumer2.start(); producer1.start(); producer2.start();}}程序运行结果如下所示imagemain线程构造PipedWriter生产者和PipedReader消费者的过程如下image向PipedWriter生产者写字符数据的生产者线程的执行过程如下image从PipedReader消费者读取字符数据的消费者线程的执行过程如下image4.1.1、循环直接写和循环直接读时多个大于等于2个生产者线程和多个大于等于2个消费者线程处理数据的过程当使用者执行4.1中的代码时多个大于等于2个个生产者线程和多个大于等于2个消费者线程处理数据的过程如下①、main线程初始化一个缓冲区char[]数组长度为1024默认值然后producer1-Thread简称生产者线程1和producer2-Thread简称生产者线程2都会根据使用者创建的字符串在堆中创建了自己要写入缓冲区char[]数组中的字符数组生产者线程1使用的字符串为Hello Pipe1生产者线程2使用的字符串为Hello Pipe2如下所示image②、然后producer1-Thread简称生产者线程1和producer2-Thread简称生产者线程2会同时调用PipedReader.class::recevice(char c[], int off, int len)函数由于这个函数是一个线程同步的函数所以同一时刻只有获取到锁的线程才能进入后续Monitor监控的代码片段本示例中是生产者线程1先获取到了锁另外一个没有获取到锁的线程只能等待获取到锁的生产者线程1释放掉锁如下所示image③、获取到锁的生产者生产者线程1将自己堆中引用的字符数组填充到缓冲区char[]数组当生产者线程填充完缓冲区之后写指针变量int in11读指针变量int out0Thread writeSide 当前这个生产者线程Thread对象生产者线程会把自己线程栈中修改的变量最终刷新到堆中PipedReader对象中以确保其它消费者线程的线程栈从堆中读取这3个变量时这3个变量已经为修改后的值如下所示image④、然后consumer1-Thread简称消费者线程1和consumer2-Thread简称消费者线程2会同时调用PipedReader.class::read(char cbuf[], int off, int len)函数由于这个函数也是一个线程同步的函数所以同一时刻只有获取到锁的线程才能进入后续Monitor监控的代码片段本示例中是消费者线程1先获取到了锁另外一个没有获取到锁的线程只能等待获取到锁的消费者线程1释放掉锁如下所示image⑤、消费者线程读缓冲区char[]数组的过程中会不断地执行out读指针以读取缓冲区char[]数组中的可用字节并返回直到out读指针in写指针或者使用者创建的字符数组读满了如果往使用者创建的字符数组中读取字符的过程中out读指针in写指针会修改in写指针-1但是这里在多核cpu上执行的时候生产者线程2此时会继续将自己堆中引用的字符数组填充到缓冲区char[]数组中这个过程会修改in写指针使得最后一次读取缓冲区char[]数组之前的所有读取都不会让out读指针和写指针相遇也就是out读指针in写指针并且每次同步执行PipedReader.class::read(char cbuf[], int off, int len)函数时都会更新Thread readSide 当前这个消费者线程Thread对象消费者线程也会把自己线程栈中修改的变量最终刷新到堆中PipedReader对象中以确保其它消费者线程的线程栈从堆中读取这3个变量时这3个变量已经为修改后的值如下所示image⑥、获取了锁的producer2-Thread简称生产者线程2将自己堆中的字符数组填充到缓冲区char[]数组获取了锁的consumer1-Thread简称消费者线程1也会从缓冲区char[]数组的out读指针和in写指针之间不断地读取字符到自己的字符数组中如下所示imageimage⑦、获取了锁的producer1-Thread简称生产者线程1将自己堆中的字符数组填充到缓冲区char[]数组获取了锁的consumer1-Thread简称消费者线程1也会从缓冲区char[]数组的out读指针和in写指针之间不断地读取字符到自己的字符数组中如下所示imageimage⑧、获取了锁的producer2-Thread简称生产者线程2将自己堆中的字符数组填充到缓冲区char[]数组获取了锁的consumer1-Thread简称消费者线程1也会从缓冲区char[]数组的out读指针和in写指针之间不断地读取字符到自己的字符数组中如下所示imageimage⑨、获取了锁的producer1-Thread简称生产者线程1将自己堆中的字符数组填充到缓冲区char[]数组获取了锁的consumer1-Thread简称消费者线程1也会从缓冲区char[]数组的out读指针和in写指针之间不断地读取字符到自己的字符数组中如下所示imageimage⑩、获取了锁的producer2-Thread简称生产者线程2将自己堆中的字符数组填充到缓冲区char[]数组获取了锁的consumer2-Thread简称消费者线程2也会从缓冲区char[]数组的out读指针和in写指针之间不断地读取字符到自己的字符数组中如下所示imageimage
