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

如何高效复制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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 15:05:16