【问题标题】:Java example of using ExecutorService and PipedReader/PipedWriter (or PipedInputStream/PipedOutputStream) for consumer-producer为消费者生产者使用 ExecutorService 和 PipedReader/PipedWriter(或 PipedInputStream/PipedOutputStream)的 Java 示例
【发布时间】:2012-05-01 21:01:05
【问题描述】:

我正在寻找一个简单的 Java 生产者 - 消费者实现并且不想重新发明轮子

我找不到同时使用新并发包和管道类的示例

是否有同时使用 PipedInputStream 和新的 Java 并发包的示例?

有没有更好的方法不使用管道类来完成这样的任务?

【问题讨论】:

  • 您到底想达到什么目的?你提出了一个非常广泛的问题。如果我使用管道流,我不介意只启动一个线程。
  • 您是否只想为消费者和生产者创建Runnables 并将它们提交给ExecutorService
  • 任务只是从数据库中读取,并以非阻塞/异步/缓冲的方式写入文件,问题中提到的工具正是我认为适合该工作的工具,如果有更简单/不同的方式,我会很高兴听到
  • @trutheality - 是的,差不多

标签: java producer-consumer executorservice java.util.concurrent


【解决方案1】:

对于您的任务,当您从数据库中读取数据时,只使用一个线程并使用BufferedOutputStream 写入文件可能就足够了。

如果您想更好地控制缓冲区大小和写入文件的块大小,您可以执行以下操作:

class Producer implements Runnable {

    private final OutputStream out;
    private final SomeDBClass db;

    public Producer( OutputStream out, SomeDBClass db ){
        this.out = out;
        this.db = db;
    }

    public void run(){
        // If you're writing to a text file you might want to wrap
        // out in a Writer instead of using `write` directly.
        while( db has more data ){
            out.write( the data );
        }
        out.flush();
        out.close();
    }
}

class Consumer implements Runnable {

    private final InputStream in;
    private final OutputStream out;
    public static final int CHUNKSIZE=512;

    public Consumer( InputStream in, OutputStream out ){
        this.out = out;
        this.in = in;
    }

    public void run(){
        byte[] chunk = new byte[CHUNKSIZE];

        for( int bytesRead; -1 != (bytesRead = in.read(chunk,0,CHUNKSIZE) );;){
            out.write(chunk, 0, bytesRead);
        }
        out.close();
    }
}

在调用代码中:

FileOutputStream toFile = // Open the stream to a file
SomeDBClass db = // Set up the db connection
PipedInputStream pi = new PipedInputStream(); // Optionally specify a size
PipedOutputStream po = new PipedOutputStream( pi );

ExecutorService exec = Executors.newFixedThreadPool(2);
exec.submit( new Producer( po, db ) );
exec.submit( new Consumer( pi, toFile ) );
exec.shutdown();
  • 同时捕获任何可能引发的异常。

请注意,如果这就是您所做的一切,那么使用ExecutorService 没有任何优势。当您有很多任务(太多而无法同时在线程中启动所有任务)时,执行程序很有用。这里只有两个线程必须同时运行,所以直接调用Thread#start 开销会更小。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-10-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多