基于Apache Async HttpClient实现Azure Storage流式响应的方案问询
实现Apache Async HttpClient响应体流式传输,避免内存缓冲
这个问题我之前也碰到过,核心是要让响应体的数据流直接从Async HttpClient的IO线程传递到客户端线程,而不是先缓冲到内存。下面是完整的实现方案:
核心思路
我们需要在自定义的AbstractAsyncResponseConsumer中创建一对管道流(PipedInputStream/PipedOutputStream):IO线程在onContentReceived回调中把收到的内容写入管道输出流,客户端则从管道输入流中读取数据,这样就能做到边接收边写入本地,完全避免大文件占用内存的问题。
自定义流式AsyncResponseConsumer实现
import org.apache.http.HttpResponse; import org.apache.http.HttpEntity; import org.apache.http.entity.ContentType; import org.apache.http.nio.ContentDecoder; import org.apache.http.nio.IOControl; import org.apache.http.nio.protocol.AbstractAsyncResponseConsumer; import org.apache.http.protocol.HttpContext; import java.io.IOException; import java.io.PipedInputStream; import java.io.PipedOutputStream; public class StreamingAsyncResponseConsumer extends AbstractAsyncResponseConsumer<HttpResponse> { private HttpResponse response; private PipedOutputStream pipedOut; private PipedInputStream pipedIn; private volatile boolean completed; public StreamingAsyncResponseConsumer() throws IOException { // 初始化管道流,设置8KB缓冲区平衡IO效率与内存占用 this.pipedIn = new PipedInputStream(8192); this.pipedOut = new PipedOutputStream(pipedIn); } @Override protected void onResponseReceived(HttpResponse response) throws HttpException, IOException { // 保存原始响应对象,后续替换entity为流式版本 this.response = response; } @Override protected void onEntityEnclosed(HttpEntity entity, ContentType contentType) throws IOException { // 跳过原始entity的处理,我们会用自己的管道流替换它 } @Override protected void onContentReceived(ContentDecoder decoder, IOControl ioctrl) throws IOException { try { byte[] buffer = new byte[8192]; int bytesRead; // 读取解码器中的数据,写入管道输出流 while ((bytesRead = decoder.read(buffer)) != -1) { pipedOut.write(buffer, 0, bytesRead); } // 通知HttpClient可以继续接收下一批数据 ioctrl.requestInput(); } catch (IOException e) { // 发生异常时关闭流,避免客户端线程永久阻塞 completed = true; pipedOut.close(); throw e; } } @Override protected HttpResponse buildResult(HttpContext context) throws Exception { // 用管道输入流创建自定义HttpEntity,替换原始响应的entity HttpEntity streamingEntity = new HttpEntity() { @Override public boolean isRepeatable() { return false; } @Override public boolean isChunked() { return true; } @Override public long getContentLength() { // Azure Storage下载响应多为chunked编码,无法提前获取长度 return -1; } @Override public ContentType getContentType() { return ContentType.getOrDefault(response.getEntity()); } @Override public java.io.InputStream getContent() throws IOException { return pipedIn; } @Override public void writeTo(java.io.OutputStream outputStream) throws IOException { throw new UnsupportedOperationException("客户端需自行处理流复制"); } @Override public boolean isStreaming() { return true; } @Override public void consumeContent() throws IOException { pipedIn.close(); } }; // 将流式entity设置到响应对象并返回 response.setEntity(streamingEntity); return response; } @Override protected void releaseResources() { // 释放所有资源,防止内存泄漏 completed = true; try { pipedOut.close(); } catch (IOException ignored) { } try { pipedIn.close(); } catch (IOException ignored) { } } }
修改请求调用代码
调整你的getFile方法,确保客户端能立即拿到包含流式entity的HttpResponse:
import org.apache.http.HttpRequestBase; import org.apache.http.HttpResponse; import org.apache.http.nio.protocol.HttpAsyncRequestProducer; import org.apache.http.nio.protocol.HttpAsyncMethods; import org.apache.http.impl.nio.client.CloseableHttpAsyncClient; import java.util.concurrent.CompletableFuture; public class AzureStorageClient { private final CloseableHttpAsyncClient httpclient; public AzureStorageClient(CloseableHttpAsyncClient httpclient) { this.httpclient = httpclient; } public CompletableFuture<HttpResponse> getFile(HttpRequestBase request) { CompletableFuture<HttpResponse> future = new CompletableFuture<>(); try { StreamingAsyncResponseConsumer consumer = new StreamingAsyncResponseConsumer(); HttpAsyncRequestProducer producer = HttpAsyncMethods.create(request); httpclient.execute(producer, consumer, new org.apache.http.concurrent.FutureCallback<HttpResponse>() { @Override public void completed(HttpResponse result) { future.complete(result); } @Override public void failed(Exception ex) { future.completeExceptionally(ex); consumer.releaseResources(); } @Override public void cancelled() { future.cancel(true); consumer.releaseResources(); } }); } catch (IOException e) { future.completeExceptionally(e); } return future; } }
客户端使用示例
客户端拿到CompletableFuture<HttpResponse>后,即可边接收边写入本地文件,无需等待整个响应完成:
azureStorageClient.getFile(request) .thenAccept(response -> { try (InputStream in = response.getEntity().getContent(); FileOutputStream out = new FileOutputStream("local-azure-file.bin")) { byte[] buffer = new byte[8192]; int bytesRead; // 流式复制数据,内存仅占用缓冲区大小 while ((bytesRead = in.read(buffer)) != -1) { out.write(buffer, 0, bytesRead); } } catch (IOException e) { e.printStackTrace(); } }) .exceptionally(ex -> { ex.printStackTrace(); return null; });
关键注意点
- 管道流缓冲区:设置合适的缓冲区大小(如8KB),可以平衡IO线程与客户端线程的阻塞频率,减少上下文切换。
- 背压控制:管道流会自动处理背压——如果客户端读取速度慢,IO线程写入管道时会被阻塞,避免内存溢出。
- 资源释放:在异常回调和
releaseResources中必须关闭管道流,防止内存泄漏和线程永久阻塞。
内容的提问来源于stack exchange,提问作者lostintranslation
相关产品推荐
相关产品推荐

