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

基于Armeria实现高并发大文件流式代理的技术问询

解决Armeria代理服务器的请求堆积与阻塞问题

首先,你的核心问题在于阻塞式的PipedStream搭配默认ForkJoinPool的supplyAsync(),和Armeria的异步非阻塞模型完全不兼容,导致高并发下线程耗尽、请求堆积。下面我会给出两种解决方案,优先推荐纯Armeria实现的方案,更贴合框架原生特性。

一、现有AsyncHttpClient代码的优化方案

针对你当前的代码,主要优化三个核心点:

  1. 替换默认线程池:放弃CompletableFuture.supplyAsync()的默认ForkJoinPool,改用Armeria提供的EventLoopGroup,避免通用线程池在高并发下被阻塞IO耗尽。
  2. 移除阻塞的PipedStream:改用AsyncHttpClient的异步流处理,避免阻塞IO导致线程挂起。
  3. 异步化所有阻塞调用:把getResponse()这类同步等待操作转成异步回调,彻底避免线程阻塞。

优化后的代码示例:

public HttpResponse doLogic(URL getUrl, String requestId, String key, Lock lock, String token) {
    // 复用Armeria的EventLoop线程池,适配异步非阻塞场景
    Executor executor = ServerConfig.defaultEventLoopGroup().next();
    
    CompletableFuture<HttpResponse> r = CompletableFuture.supplyAsync(() -> {
        AsyncHttpClient asyncHttpClient = ...; // 务必全局复用客户端实例,不要每次创建
        try {
            // 异步获取源端响应,避免同步阻塞
            return asyncHttpClient.prepareGet(getUrl.toString())
                    .execute()
                    .toCompletableFuture()
                    .thenApply(sourceResp -> {
                        if (sourceResp.getStatusCode() != 200) {
                            log.warn("[requestId: {}] s3 came back with status code: {}", requestId, sourceResp.getStatusCode());
                            return CompletableFuture.completedFuture(unlockingResponse(200, lock, key));
                        }
                        // 转换并处理响应头
                        HttpHeaders headers = translateHeaders(sourceResp.getHeaders());
                        if (headers.get("etag") != null && headers.get("etag").contains("-")) {
                            headers.remove("etag");
                        }
                        headers.add("x-auth-token", token);
                        // 异步发起PUT请求,直接用源端响应流作为body
                        return asyncHttpClient.preparePut(Config.getString("swift.proxy.base.url") + "v1/" + key)
                                .setHeaders(headers)
                                .setBody(sourceResp.getResponseBodyAsStream())
                                .execute()
                                .toCompletableFuture()
                                .thenApply(targetResp -> {
                                    if (targetResp.getStatusCode() != 201 && targetResp.getStatusCode() != 202) {
                                        log.error("[requestId: {}] swift returned {}", requestId, targetResp.getStatusCode());
                                        return unlockingResponse(targetResp.getStatusCode(), lock, key);
                                    }
                                    log.info("[requestId: {}] File: {} rollbacked!", requestId, key);
                                    return unlockingResponse(200, lock, key);
                                })
                                .exceptionally(e -> {
                                    e.printStackTrace();
                                    return unlockingResponse(500, lock, key);
                                });
                    })
                    .exceptionally(e -> {
                        e.printStackTrace();
                        return CompletableFuture.completedFuture(unlockingResponse(500, lock, key));
                    })
                    .join();
        } catch (Exception e) {
            e.printStackTrace();
            return unlockingResponse(500, lock, key);
        }
    }, executor);
    
    return HttpResponse.from(r);
}

二、纯Armeria实现方案(强烈推荐)

Armeria本身提供了异步非阻塞的HttpClient,完全适配其HttpResponse模型,不需要依赖第三方AsyncHttpClient,能实现真正的零拷贝无存储流式转发,完美支持高并发场景:

核心优势:

  • 零拷贝流式传输:源端的响应流直接转发到目标端,不需要中间存储或阻塞IO转换
  • 原生适配异步模型:所有操作都在Armeria的EventLoop上执行,避免线程阻塞和上下文切换
  • 内置连接池优化:HttpClient自带连接池,自动处理高并发下的连接复用和限流

纯Armeria代码示例:

// 全局复用HttpClient实例,不要每次请求创建,避免资源浪费
private static final HttpClient sourceS3Client = HttpClient.of("https://your-source-s3-endpoint");
private static final HttpClient targetSwiftClient = HttpClient.of(Config.getString("swift.proxy.base.url"));

public HttpResponse doLogic(URL getUrl, String requestId, String key, Lock lock, String token) {
    // 1. 异步发起源端GET请求
    return sourceS3Client.get(getUrl.toString())
            .thenCompose(sourceResp -> {
                // 2. 处理源端非200响应
                if (sourceResp.status().code() != 200) {
                    log.warn("[requestId: {}] s3 came back with status code: {}", requestId, sourceResp.status().code());
                    return HttpResponse.of(200)
                            .toCompletionStage()
                            .thenApply(resp -> unlockingResponse(200, lock, key));
                }
                // 3. 转换并清理响应头
                HttpHeaders headers = HttpHeaders.of(sourceResp.headers());
                if (headers.get("etag") != null && headers.get("etag").contains("-")) {
                    headers.remove("etag");
                }
                headers.add("x-auth-token", token);
                // 4. 流式转发到目标端:直接用源端的StreamMessage作为PUT请求体(零拷贝)
                return targetSwiftClient.put("/v1/" + key)
                        .headers(headers)
                        .body(sourceResp.content())
                        .execute()
                        .thenCompose(targetResp -> {
                            // 5. 处理目标端响应状态
                            if (targetResp.status().code() != 201 && targetResp.status().code() != 202) {
                                log.error("[requestId: {}] swift returned {}", requestId, targetResp.status().code());
                                return HttpResponse.of(targetResp.status().code())
                                        .toCompletionStage()
                                        .thenApply(resp -> unlockingResponse(targetResp.status().code(), lock, key));
                            }
                            log.info("[requestId: {}] File: {} rollbacked!", requestId, key);
                            return HttpResponse.of(200)
                                    .toCompletionStage()
                                    .thenApply(resp -> unlockingResponse(200, lock, key));
                        })
                        .exceptionally(e -> {
                            e.printStackTrace();
                            return unlockingResponse(500, lock, key);
                        });
            })
            .exceptionally(e -> {
                e.printStackTrace();
                return unlockingResponse(500, lock, key);
            })
            .toHttpResponse();
}

// 辅助方法:释放锁并返回对应HttpResponse
private HttpResponse unlockingResponse(int statusCode, Lock lock, String key) {
    try {
        lock.unlock();
    } catch (Exception e) {
        log.error("Failed to unlock key: {}", key, e);
    }
    return HttpResponse.of(statusCode);
}

关键注意点:

  • 务必复用HttpClient实例:全局单例复用连接池,避免每次请求创建新客户端导致的资源耗尽
  • 用thenCompose链式处理异步操作:不要用同步等待(如get()),保持异步流的连贯性
  • 异常处理全覆盖:每个异步步骤都要处理异常,避免未捕获的异常导致请求堆积
  • 锁的释放要可靠:确保无论成功还是失败,锁都能被正确释放,避免死锁

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:18:44