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

AWS SDK v2上传大文件时CompletableFuture.join()无法返回问题

大文件上传时CompletableFuture.join()无限挂起问题排查

问题场景

因AWS SDK v1停止支持,升级至SDK v2后替换为异步客户端及新版传输管理器,通过CompletableFuture.join()同步调用上传逻辑。小文件上传正常,但5GB以上大文件上传时出现异常:Spring Boot服务的http-nio请求线程打印“Before request”后无响应,日志显示传输管理器已发送completeMultipartUploadRequest,但CompletableFuture始终未完成,join()调用导致线程无限挂起。

相关代码:

public void putObject(PutObjectRequest putObjectRequest, InputStream inputStream) {
    ExecutorService executorService = null;
    try {
        executorService = Executors.newSingleThreadExecutor();
        transferManager.upload(UploadRequest.builder()
                        .putObjectRequest(putObjectRequest)
                        .requestBody(
                                AsyncRequestBody.fromInputStream(
                                        AsyncRequestBodyFromInputStreamConfiguration.builder()
                                                .executor(executorService)
                                                .contentLength(putObjectRequest.contentLength())
                                                .inputStream(inputStream)
                                                .build()
                                )
                        )
                        .build()
                )
                .completionFuture()
                .join();
    } catch (Exception ex) {
        throwErrorOnUploadFail(ex);
    } finally {
        if (Objects.nonNull(executorService)) {
            executorService.shutdown();
        }
    }
}

问题原因

  1. 线程池过早关闭:finally块中调用executorService.shutdown()不会等待已提交的任务完成,仅停止接受新任务。大文件采用分片上传,AsyncRequestBody.fromInputStream依赖该线程池处理流的读取,主线程调用join()时,若分片读取任务尚未完成,线程池已被关闭,导致任务无法继续推进,CompletableFuture永远无法触发完成,主线程无限阻塞。
  2. 单线程池瓶颈:newSingleThreadExecutor()是单线程池,大文件分片上传需要处理多个分片的流读取任务,单线程无法及时处理所有任务,加剧了线程池关闭后的任务停滞问题。
  3. IO线程阻塞风险:Spring Boot的http-nio线程是处理请求的核心非阻塞线程,调用join()会直接阻塞该线程,导致服务无法处理其他请求,同时也放大了线程池异常带来的影响。

解决方案

1. 延迟线程池关闭,确保任务完成

将finally块中的shutdown()替换为awaitTermination(),等待所有任务执行完毕后再关闭线程池:

public void putObject(PutObjectRequest putObjectRequest, InputStream inputStream) {
    ExecutorService executorService = null;
    try {
        executorService = Executors.newFixedThreadPool(4); // 改用多线程池适配分片任务
        CompletableFuture<Void> uploadFuture = transferManager.upload(UploadRequest.builder()
                        .putObjectRequest(putObjectRequest)
                        .requestBody(
                                AsyncRequestBody.fromInputStream(
                                        AsyncRequestBodyFromInputStreamConfiguration.builder()
                                                .executor(executorService)
                                                .contentLength(putObjectRequest.contentLength())
                                                .inputStream(inputStream)
                                                .build()
                                )
                        )
                        .build()
                )
                .completionFuture();
        
        uploadFuture.join();
    } catch (Exception ex) {
        throwErrorOnUploadFail(ex);
    } finally {
        if (Objects.nonNull(executorService)) {
            executorService.shutdown();
            try {
                // 根据大文件上传时长设置足够的等待时间
                if (!executorService.awaitTermination(30, TimeUnit.MINUTES)) {
                    executorService.shutdownNow();
                }
            } catch (InterruptedException e) {
                executorService.shutdownNow();
                Thread.currentThread().interrupt();
            }
        }
    }
}

2. 使用全局复用的线程池

避免每次上传都创建销毁线程池,改用全局线程池(如Spring管理的TaskExecutor),在服务关闭时统一销毁:

// 注入Spring全局线程池,或自行创建全局线程池
@Autowired
private TaskExecutor uploadTaskExecutor;

public void putObject(PutObjectRequest putObjectRequest, InputStream inputStream) {
    try {
        transferManager.upload(UploadRequest.builder()
                        .putObjectRequest(putObjectRequest)
                        .requestBody(
                                AsyncRequestBody.fromInputStream(
                                        AsyncRequestBodyFromInputStreamConfiguration.builder()
                                                .executor(uploadTaskExecutor)
                                                .contentLength(putObjectRequest.contentLength())
                                                .inputStream(inputStream)
                                                .build()
                                )
                        )
                        .build()
                )
                .completionFuture()
                .join();
    } catch (Exception ex) {
        throwErrorOnUploadFail(ex);
    }
}

// 服务关闭时销毁线程池
@PreDestroy
public void shutdownExecutor() {
    if (uploadTaskExecutor instanceof ExecutorService) {
        ExecutorService executor = (ExecutorService) uploadTaskExecutor;
        executor.shutdown();
        try {
            if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
                executor.shutdownNow();
            }
        } catch (InterruptedException e) {
            executor.shutdownNow();
        }
    }
}

3. 避免阻塞IO线程

将上传逻辑提交到独立线程池执行,返回CompletableFuture给上层,让请求线程保持非阻塞:

public CompletableFuture<Void> putObjectAsync(PutObjectRequest putObjectRequest, InputStream inputStream) {
    return CompletableFuture.runAsync(() -> {
        try {
            transferManager.upload(UploadRequest.builder()
                            .putObjectRequest(putObjectRequest)
                            .requestBody(
                                    AsyncRequestBody.fromInputStream(
                                            AsyncRequestBodyFromInputStreamConfiguration.builder()
                                                    .executor(uploadTaskExecutor)
                                                    .contentLength(putObjectRequest.contentLength())
                                                    .inputStream(inputStream)
                                                    .build()
                                    )
                            )
                            .build()
                    )
                    .completionFuture()
                    .join();
        } catch (Exception ex) {
            throwErrorOnUploadFail(ex);
        }
    }, uploadTaskExecutor);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 01:50:20