基于Armeria实现高并发大文件流式代理的技术问询
解决Armeria代理服务器的请求堆积与阻塞问题
首先,你的核心问题在于阻塞式的PipedStream搭配默认ForkJoinPool的supplyAsync(),和Armeria的异步非阻塞模型完全不兼容,导致高并发下线程耗尽、请求堆积。下面我会给出两种解决方案,优先推荐纯Armeria实现的方案,更贴合框架原生特性。
一、现有AsyncHttpClient代码的优化方案
针对你当前的代码,主要优化三个核心点:
- 替换默认线程池:放弃
CompletableFuture.supplyAsync()的默认ForkJoinPool,改用Armeria提供的EventLoopGroup,避免通用线程池在高并发下被阻塞IO耗尽。 - 移除阻塞的PipedStream:改用AsyncHttpClient的异步流处理,避免阻塞IO导致线程挂起。
- 异步化所有阻塞调用:把
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
相关产品推荐
相关产品推荐

