如何实现非字节数组型的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
相关产品推荐
相关产品推荐

