You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Java的PipedOutputStream和PipedInputStream实现类生产者消费者问题示例

完整的Java PipedOutputStream/PipedInputStream 生产者消费者示例

我来帮你完善这个基于管道流的生产者消费者示例!PipedOutputStream 和 PipedInputStream 本身就是为线程间的字节流通信设计的,天生适配生产者消费者模型,下面是完整可运行的实现:

import java.io.IOException;
import java.io.PipedInputStream;
import java.io.PipedOutputStream;
import java.util.concurrent.ThreadLocalRandom;

public class PipeProducerConsumer {

    public static void main(String[] args) {
        try {
            // 关联管道输出流和输入流,也可以用构造函数直接绑定:new PipedInputStream(pipedSrc)
            final PipedOutputStream pipedSrc = new PipedOutputStream();
            final PipedInputStream pipedSnk = new PipedInputStream(pipedSrc);

            // 随机生成5-10个要生产的数字
            int putNumbers = ThreadLocalRandom.current().nextInt(5, 10);
            System.out.printf("生产者将生产 %d 个数字%n", putNumbers);

            // 消费者线程:从管道读取数字
            Runnable consumer = () -> {
                try {
                    int receivedNum;
                    while ((receivedNum = pipedSnk.read()) != -1) {
                        System.out.printf("消费者读取到数字:%d%n", receivedNum);
                        // 模拟消费耗时
                        Thread.sleep(ThreadLocalRandom.current().nextInt(500, 1000));
                    }
                    System.out.println("消费者:管道已关闭,停止读取");
                } catch (IOException | InterruptedException e) {
                    Thread.currentThread().interrupt();
                    System.err.println("消费者线程异常:" + e.getMessage());
                } finally {
                    try {
                        pipedSnk.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            };

            // 生产者线程:往管道写入数字
            Runnable producer = () -> {
                try {
                    for (int i = 0; i < putNumbers; i++) {
                        int num = ThreadLocalRandom.current().nextInt(1, 100);
                        System.out.printf("生产者写入数字:%d%n", num);
                        pipedSrc.write(num);
                        // 模拟生产耗时
                        Thread.sleep(ThreadLocalRandom.current().nextInt(300, 800));
                    }
                    System.out.println("生产者:所有数字已写入,关闭管道输出流");
                } catch (IOException | InterruptedException e) {
                    Thread.currentThread().interrupt();
                    System.err.println("生产者线程异常:" + e.getMessage());
                } finally {
                    try {
                        pipedSrc.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            };

            // 启动两个线程
            new Thread(consumer, "Consumer-Thread").start();
            new Thread(producer, "Producer-Thread").start();

        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

关键细节说明:

  • 管道关联:可以通过connect()方法绑定输入输出流,也可以直接使用PipedInputStream(PipedOutputStream src)构造函数一步完成绑定,两种方式效果一致。
  • 阻塞特性:管道流自带阻塞机制:当管道为空时,消费者的read()方法会阻塞等待数据;当管道满(默认缓冲区大小为1024字节)时,生产者的write()方法会阻塞,不需要额外加锁或同步工具,天然适配生产者消费者模型。
  • 资源关闭:必须在finally块中关闭流,避免资源泄漏;当生产者关闭PipedOutputStream时,消费者的read()方法会返回-1,以此判断数据读取完毕。
  • 异常处理:线程中的异常需要捕获并处理,同时要设置线程中断状态,避免线程资源泄漏。

内容的提问来源于stack exchange,提问作者user7339033

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 08:26:29