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

如何实现非字节数组型的Async HTTP Client异步输入流?

解决Async Http Client大文件流式传输避免OOM的问题

你说的没错,直接用getResponseBodyAsStream()确实会把整个文件加载到内存,对于大文件来说很容易触发OOM。你提到的两个思路都可行,我给你分别提供具体的实现示例,其中管道流的方案更高效,不需要落地磁盘。

方案1:写入临时文件后返回文件输入流

这个方案逻辑简单,把接收到的每块数据写入临时文件,最后返回文件的输入流给下游服务解析。需要注意临时文件的清理,避免磁盘空间浪费。

import org.asynchttpclient.*;
import java.io.*;
import java.nio.ByteBuffer;

public FileInputStream downloadToTempFile(String url, AsyncHttpClient client) throws Exception {
    // 创建临时文件,JVM退出时自动删除
    File tempFile = File.createTempFile("download-", ".tmp");
    tempFile.deleteOnExit();

    try (FileOutputStream fos = new FileOutputStream(tempFile)) {
        client.prepareGet(url).execute(new AsyncHandler<Void>() {
            @Override
            public State onBodyPartReceived(HttpResponseBodyPart bodyPart) throws Exception {
                ByteBuffer buffer = bodyPart.getBodyByteBuffer();
                // 将ByteBuffer写入文件输出流
                byte[] bytes = new byte[buffer.remaining()];
                buffer.get(bytes);
                fos.write(bytes);
                return State.CONTINUE;
            }

            @Override
            public State onHeadersReceived(HttpHeaders headers) throws Exception {
                return State.CONTINUE;
            }

            @Override
            public State onStatusReceived(HttpResponseStatus responseStatus) throws Exception {
                if (responseStatus.getStatusCode() != 200) {
                    return State.ABORT;
                }
                return State.CONTINUE;
            }

            @Override
            public Void onCompleted() throws Exception {
                fos.flush();
                return null;
            }

            @Override
            public void onThrowable(Throwable t) {
                // 发生异常时删除临时文件
                tempFile.delete();
            }
        }).get();

        return new FileInputStream(tempFile);
    }
}

方案2:使用PipedInputStream/PipedOutputStream(推荐)

这个方案不需要磁盘IO,直接在内存中通过管道流实现生产者-消费者模式:AsyncHandler作为生产者把数据写入PipedOutputStream,下游服务作为消费者从PipedInputStream读取数据。需要注意线程同步问题,因为AsyncHandler的回调是在异步线程中执行的,要避免死锁。

import org.asynchttpclient.*;
import java.io.*;
import java.nio.ByteBuffer;

public InputStream downloadWithPipedStream(String url, AsyncHttpClient client) throws Exception {
    PipedInputStream pipedIn = new PipedInputStream();
    PipedOutputStream pipedOut = new PipedOutputStream(pipedIn);

    client.prepareGet(url).execute(new AsyncHandler<Void>() {
        @Override
        public State onBodyPartReceived(HttpResponseBodyPart bodyPart) throws Exception {
            try {
                ByteBuffer buffer = bodyPart.getBodyByteBuffer();
                byte[] bytes = new byte[buffer.remaining()];
                buffer.get(bytes);
                // 写入管道输出流,注意如果管道满了会阻塞
                pipedOut.write(bytes);
            } catch (IOException e) {
                // 下游服务关闭输入流时会抛出异常,此时终止请求
                return State.ABORT;
            }
            return State.CONTINUE;
        }

        @Override
        public State onHeadersReceived(HttpHeaders headers) throws Exception {
            return State.CONTINUE;
        }

        @Override
        public State onStatusReceived(HttpResponseStatus responseStatus) throws Exception {
            if (responseStatus.getStatusCode() != 200) {
                // 状态码不对,关闭管道并终止
                pipedOut.close();
                return State.ABORT;
            }
            return State.CONTINUE;
        }

        @Override
        public Void onCompleted() throws Exception {
            // 请求完成后关闭输出流,让输入流知道没有更多数据
            pipedOut.close();
            return null;
        }

        @Override
        public void onThrowable(Throwable t) {
            try {
                pipedOut.close();
            } catch (IOException ignored) {}
        }
    });

    // 返回管道输入流给下游服务
    return pipedIn;
}

注意事项

  • 使用管道流时,下游服务必须尽快读取数据,否则AsyncHandler的写入操作会阻塞,影响请求处理。
  • 两种方案都要处理异常情况,比如请求失败时及时关闭流或删除临时文件,避免资源泄漏。
  • 你可以根据实际场景选择:如果下游服务解析速度慢,临时文件方案更稳妥;如果追求性能,管道流是更好的选择。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:22:16