三公机器人

牛牛机器人,三公撑船机器人,微信牛牛机器人

三公机器人 在 Windows 操作系统下使用 JDK 8 进行开发

在 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 包中的并发容器来实现线程间数据传递。


Powered By Z-BlogPHP 1.7.3

三公机器人,牛牛机器人,三公撑船机器人,微信牛牛机器人