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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:23:47