如何高效复制PipedOutputStream?解决Quarkus API线程阻塞问题
我想最佳实践且内存高效地把一个PipedOutputStream复制到另一个PipedOutputStream,这可行吗?
我有个方法会往PipedOutputStream output1写数据,写完后想把这些数据复制到另一个PipedOutputStream output2,该怎么实现?
我的现有代码如下:
private InputStream converterMethod(final InputStream inputDocument) { try { final ExecutorService executorService = Executors.newWorkStealingPool(); final PipedOutputStream output1 = new PipedOutputStream(); final CurrentHandler handler = new CurrentHandler(output1); final PipedOutputStream output2 = new PipedOutputStream(); final InputStream finalInputStream = new PipedInputStream(output2); executorService.execute( () -> { try { new ConverterApplication().convert(inputDocument, handler); output1.close(); } catch (Exception e) { e.printStackTrace(); } }); executorService.execute(() ->{ try { //Planning to copy the output1 to output2 in the best way possible. } catch (Exception e) { e.printStackTrace(); } }); return finalInputStream; } catch (Exception e) { throw new MyCustomException("Exception occurred during"); } }
我试过只保留一个PipedOutputStream,本地运行正常,但通过Quarkus RestAPI调用时出现如下错误:
Thread Thread[vert.x-eventloop-thread-0,5,main] has been blocked for 620239 ms, time limit is 2000 ms: io.vertx.core.VertxException: Thread blocked at java.base@17.0.3/java.lang.Object.wait(Native Method) at java.base@17.0.3/java.io.PipedInputStream.read(PipedInputStream.java:326) at java.base@17.0.3/java.io.PipedInputStream.read(PipedInputStream.java:377) at java.base@17.0.3/java.io.InputStream.read(InputStream.java:218) at org.jboss.resteasy.reactive.common.providers.serialisers.InputStreamMessageBodyHandler.writeTo(InputStreamMessageBodyHandler.java:39) at org.jboss.resteasy.reactive.server.providers.serialisers.ServerInputStreamMessageBodyHandler.writeResponse(ServerInputStreamMessageBodyHandler.java:46) at org.jboss.resteasy.reactive.server.providers.serialisers.ServerInputStreamMessageBodyHandler.writeResponse(ServerInputStreamMessageBodyHandler.java:16)
对应的简化代码如下:
private InputStream converterMethod(final InputStream inputDocument) { try { final ExecutorService executorService = Executors.newWorkStealingPool(); final PipedOutputStream output1 = new PipedOutputStream(); final CurrentHandler handler = new CurrentHandler(output1); final InputStream finalInputStream = new PipedInputStream(output1); executorService.execute( () -> { try { new ConverterApplication().convert(inputDocument, handler); output1.close(); } catch (Exception e) { e.printStackTrace(); } }); return finalInputStream; } catch (Exception e) { throw new MyCustomException("Exception occurred during"); } }
我的核心需求是:如何从ConverterApplication中获取最终数据作为InputStream返回。
问题根源分析
Quarkus基于Vert.x事件循环,要求事件线程不能被阻塞。简化方案中,PipedInputStream的read()方法会在无数据时阻塞线程,而RESTEasy Reactive在事件循环线程中处理InputStream的写入操作,直接触发了阻塞超时。
正确实现方式
不需要额外复制PipedOutputStream,核心是确保PipedInputStream的读取操作不在事件循环线程中执行,或通过更适配的异步流处理避免阻塞。以下是两种可靠方案:
方案1:管道流+虚拟线程(内存高效)
利用Java虚拟线程(Java 19+)执行转换任务,同时包装返回的InputStream,确保读取操作不会阻塞事件循环:
private InputStream converterMethod(final InputStream inputDocument) { try { // 创建管道流对:转换线程写outputPipe,外部从inputPipe读数据 final PipedInputStream inputPipe = new PipedInputStream(); final PipedOutputStream outputPipe = new PipedOutputStream(inputPipe); final CurrentHandler handler = new CurrentHandler(outputPipe); // 用虚拟线程执行转换任务,避免占用平台线程资源 CompletableFuture.runAsync(() -> { try { new ConverterApplication().convert(inputDocument, handler); outputPipe.close(); } catch (Exception e) { // 异常时关闭管道,避免永久阻塞 try { outputPipe.close(); } catch (IOException ex) { ex.addSuppressed(e); } throw new CompletionException(e); } }, Executors.newVirtualThreadPerTaskExecutor()); // 包装InputStream,确保读取逻辑不绑定事件循环线程 return new InputStream() { @Override public int read() throws IOException { return inputPipe.read(); } @Override public int read(byte[] b, int off, int len) throws IOException { return inputPipe.read(b, off, len); } @Override public void close() throws IOException { inputPipe.close(); outputPipe.close(); } }; } catch (Exception e) { throw new MyCustomException("Exception occurred during conversion"); } }
方案2:Quarkus异步响应类型(推荐)
直接返回Uni<InputStream>适配Quarkus异步模型,从根源避免事件循环阻塞:
import io.smallrye.mutiny.Uni; import java.io.ByteArrayOutputStream; import java.io.InputStream; // 修改方法返回类型为Uni<InputStream> public Uni<InputStream> converterMethod(final InputStream inputDocument) { return Uni.createFrom().item(() -> { ByteArrayOutputStream baos = new ByteArrayOutputStream(); CurrentHandler handler = new CurrentHandler(baos); new ConverterApplication().convert(inputDocument, handler); return new ByteArrayInputStream(baos.toByteArray()); }).runSubscriptionOn(Executors.newVirtualThreadPerTaskExecutor()); }
注意:如果转换后数据量极大,ByteArrayOutputStream会占用过多内存,此时优先选择方案1的管道流+虚拟线程组合。
关于PipedOutputStream复制的说明
直接复制两个PipedOutputStream是可行的,但需要额外线程做桥接,核心代码如下(需放在第二个executor任务中):
// 连接output1到输入流,再将数据复制到output2 PipedInputStream inputFromOutput1 = new PipedInputStream(output1); byte[] buffer = new byte[8192]; // 8KB缓冲区,平衡内存占用与复制效率 int bytesRead; while ((bytesRead = inputFromOutput1.read(buffer)) != -1) { output2.write(buffer, 0, bytesRead); } // 关闭相关流 output2.close(); inputFromOutput1.close();
但这种方式多了一层线程和管道开销,不如直接调整管道流的使用方式高效,且同样要确保复制操作在非事件线程执行。
内容的提问来源于stack exchange,提问作者BATMAN_2008

