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 out=0
    snk.connected = true;//改变PipedReader (消费者)中的变量boolean connected=true,表示该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,off+len)(左闭右开,不包括off+len)索引位置的字符
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,off+len)(左闭右开)索引位置是否有越界的检查
        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 out=0、boolean connected=true
    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 out=0、boolean connected=true
}

//线程同步函数:该函数只被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,off+len)(左闭右开,不包括off+len)索引位置的字符
    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)索引位置的字符,并且将读指针(out)+1
    if (out >= buffer.length) {
        out = 0;//如果读指针(out)>=缓冲区(char[]数组)的长度,设置读指针(out)=0
    }
    if (in == out) {
        /* now empty */
        in = -1;//如果消费者线程从缓冲区(char[]数组)中读完字符(char)数据以后读指针(out)=写指针(in),那么,当前消费者线程会设置写指针(in)=-1
    }
    return ret;
}

//线程同步函数:如果缓冲区(char[]数组)中有足够多的字符的话(数量>len),消费者线程每次从缓冲区(char[]数组)中读取len个字符放到char[]数组cbuf的[off, off+len)索引位置(左闭右开,不包括off+len)
//如果缓冲区(char[]数组)中字符的数量<len个(比如有in(写指针)-out(读指针)个),消费者线程每次从缓冲区(char[]数组)中读取(in-out)个字符放到char[]数组cbuf的[off, off+in-out)索引位置(左闭右开,不包括off+in-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,off+len)(左闭右开)索引位置是否有越界的检查
        throw new IndexOutOfBoundsException();//越界的话,抛出一个IndexOutOfBoundsException
    } else if (len == 0) {
        return 0;//如果len==0,返回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[]数组)的长度,设置读指针(out)=0
        }
        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(消费者),其实就是设置标记符closedByReader=true, 设置写指针(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。

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐