【问题标题】:Exception propagation within PipedInputStream and PipedOutputStreamPipedInputStream 和 PipedOutputStream 中的异常传播
【发布时间】:2015-11-13 05:23:27
【问题描述】:

我有一个数据生产者,它在单独的线程中运行并将生成的数据推送到连接到PipedInputStream 的PipedOutputStream。此输入流的引用通过公共 API 公开,以便任何客户端都可以使用它。 PipedInputStream 包含一个有限的缓冲区,如果已满,则阻塞数据生产者。基本上,当客户端从输入流中读取数据时,数据生产者会生成新数据。

问题在于数据生产者可能会失败并抛出异常。但是由于消费者在单独的线程中运行,因此没有很好的方法将异常发送给客户端。

我所做的是捕获该异常并关闭输入流。这将导致IOException 在客户端出现消息“管道已关闭”,但我真的很想向客户说明这背后的真正原因。

这是我的 API 的粗略代码:

public InputStream getData() {
    final PipedInputStream inputStream = new PipedInputStream(config.getPipeBufferSize());
    final PipedOutputStream outputStream = new PipedOutputStream(inputStream);

    Thread thread = new Thread(() -> {
        try {
          // Start producing the data and push it into output stream.
          // The production my fail and throw an Exception with the reason
        } catch (Exception e) {
            try {
                // What to do here?
                outputStream.close();
                inputStream.close();
            } catch (IOException e1) {
            }
        }
    });
    thread.start();

    return inputStream;
}

我有两个想法可以解决这个问题:

  1. 将异常存储在父对象中并通过 API 将其公开给客户端。 IE。如果读取失败并返回 IOException,客户端可以向 API 询问原因。
  2. 扩展/重新实现管道流,以便我可以将原因传递给close() 方法。然后流抛出的IOException 可以包含该原因作为消息。

有更好的想法吗?

【问题讨论】:

    标签: java multithreading


    【解决方案1】:

    巧合的是,我刚刚编写了类似的代码来允许对流进行 GZip 压缩。您不需要扩展 PipedInputStream,只需 FilterInputStream 即可完成并返回一个包装的版本,例如

    final PipedInputStream in = new PipedInputStream();
    final InputStreamWithFinalExceptionCheck inWithException = new InputStreamWithFinalExceptionCheck(in);
    final PipedOutputStream out = new PipedOutputStream(in);
    Thread thread = new Thread(() -> {
        try {
          // Start producing the data and push it into output stream.
          // The production my fail and throw an Exception with the reason
        } catch (final IOException e) {
            inWithException.fail(e);
        } finally {
            inWithException.countDown();
        }
    });
    thread.start();
    return inWithException;
    

    然后 InputStreamWithFinalExceptionCheck 只是

    private static final class InputStreamWithFinalExceptionCheck extends FilterInputStream {
        private final AtomicReference<IOException> exception = new AtomicReference<>(null);
        private final CountDownLatch complete = new CountDownLatch(1);
    
        public InputStreamWithFinalExceptionCheck(final InputStream stream) {
            super(stream);
        }
    
        @Override
        public void close() throws IOException {
            try {
                complete.await();
                final IOException e = exception.get();
                if (e != null) {
                    throw e;
                }
            } catch (final InterruptedException e) {
                throw new IOException("Interrupted while waiting for synchronised closure");
            } finally {
                stream.close();
            }
        }
    
        public void fail(final IOException e) {
            exception.set(Preconditions.checkNotNull(e));
        }
    
        public void countDown() {complete.countDown();}
    }
    

    【讨论】:

    • 如果你扩展FilterInputStream,你可以让它变得更简单。它为您完成了大部分工作。
    • 该解决方案对我有用,但值得注意的是,两个流都必须明确关闭(上面的代码中没有,但我想这是一种礼貌的暗示:)) .在我的例子中,数据是由多个并行作业产生的,最后一个关闭输出流。如果失败,则必须在 finally 块中完成,否则会抛出“Write end dead”。在输入流中,close() 是引发异常的原因。我在 JUnit 测试中使用IOUtils.toString(),它没有正确关闭流,因此不会引发异常。
    【解决方案2】:

    这是我的实现,取自上面接受的答案 https://stackoverflow.com/a/33698661/5165540 ,我不使用 CountDownLatch complete.await() 因为如果在作者完成写入完整内容之前 InputStream 突然关闭,它会导致死锁。 我仍然设置在使用 PipedOutpuStream 时捕获的异常,并在生成线程中创建 PipedOutputStream,使用 try-finally-resource 模式确保它被关闭,在供应商中等待直到 2 个流被管道传输。

    Supplier<InputStream> streamSupplier = new Supplier<InputStream>() {
            @Override
            public InputStream get() {
                final AtomicReference<IOException> osException = new AtomicReference<>();
                final CountDownLatch piped = new CountDownLatch(1);
    
                final PipedInputStream is = new PipedInputStream();
    
                FilterInputStream fis = new FilterInputStream(is) {
                    @Override
                    public void close() throws IOException {
                        try {
                            IOException e = osException.get();
                            if (e != null) {
                                //Exception thrown by the write will bubble up to InputStream reader
                                throw new IOException("IOException in writer", e);
                            }
                        } finally {
                            super.close();
                        }
                    };
                };
    
                Thread t = new Thread(() -> {
                        try (PipedOutputStream os = new PipedOutputStream(is)) {
                            piped.countDown();
                            writeIozToStream(os, projectFile, dataFolder);
                        } catch (final IOException e) {
                            osException.set(e);
                        }
                });
                t.start();
    
                try {
                    piped.await();
                } catch (InterruptedException e) {
                    t.cancel();
                    Thread.currentThread().interrupt();
                }
    
                return fis;
            }
        };
    

    调用代码类似于

    try (InputStream is = streamSupplier.getInputStream()) {
         //Read stream in full 
    }
    

    因此,当 InputStream 关​​闭时,这将在 PipedOutputStream 中发出信号,最终导致“管道关闭”IOException,此时将被忽略。

    如果我在 FilterInputStreamclose() 中保留 complete.await() 行,我可能会遇到死锁(PipedInputStream 试图关闭,在 complete.await() 上等待,而 PipedOutputStream 在 PipedInputStreamawaitSpace 上永远等待)

    【讨论】:

    • 为了避免死锁等待超时也可以使用。在您的情况下,可以在 writeIozToStream 引发异常之前关闭 InputStream 并且不会检查此异常。
    猜你喜欢
    • 2023-04-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-03-18
    • 2017-10-28
    • 2018-10-06
    • 1970-01-01
    相关资源
    最近更新 更多