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

Spring Boot 3.0.x异步流式上传文件至S3失败求助

解决Spring Boot 3流式上传至S3失败的问题

针对你用Spring Boot 3的PartEvent+DataBuffer流式上传到S3失败、数据流未启动且FilePartEvent重复触发的问题,结合你的代码和约束(不能缓存全文件到内存/磁盘),给出以下具体解决思路:

1. 修复DataBuffer释放时机问题

当前代码在map中直接释放DataBuffer,但AWS的AsyncRequestBody可能还未读取对应的ByteBuffer,导致数据丢失、上传中断。需确保DataBuffer在ByteBuffer被使用后再释放:

修改S3AttachmentService.save方法:

public Mono<Boolean> save(String filename, long contentLength, Flux<DataBuffer> contentBuffers) {
    Flux<ByteBuffer> buffers = contentBuffers.flatMap(db -> {
        ByteBuffer buf = db.toByteBuffer();
        // 在Mono完成后释放DataBuffer,确保S3已读取数据
        return Mono.just(buf)
                .doOnTerminate(() -> DataBufferUtils.release(db));
    });

    return s3Service.uploadBuffers(attachmentConfigProperties.getBucketName(), filename, contentLength, buffers);
}

2. 确保contentLength的准确性

你当前使用accumulator.attachmentSize()作为S3上传的contentLength,但这个值可能并非实际文件的完整大小(因为流式处理时文件还未完全接收),导致S3接收到的字节数与声明值不匹配,触发上传失败进而重试FilePartEvent。

  • 如果无法提前获取准确的文件大小,直接移除contentLength参数,让AWS SDK自动处理流式上传:
// 修改S3AttachmentService.save方法,去掉contentLength参数
public Mono<Boolean> save(String filename, Flux<DataBuffer> contentBuffers) {
    Flux<ByteBuffer> buffers = contentBuffers.flatMap(db -> {
        ByteBuffer buf = db.toByteBuffer();
        return Mono.just(buf)
                .doOnTerminate(() -> DataBufferUtils.release(db));
    });

    return s3Service.uploadBuffers(attachmentConfigProperties.getBucketName(), filename, buffers);
}

// 修改S3Service.uploadBuffers方法,移除contentLength相关配置
public Mono<Boolean> uploadBuffers(String bucketName, String objectKey, Flux<ByteBuffer> buffers) {
    return Mono.create(sink -> {
        CompletableFuture<PutObjectResponse> future = s3AsyncClient.putObject(
                PutObjectRequest.builder()
                        .bucket(bucketName)
                        .key(objectKey)
                        .build(),
                AsyncRequestBody.fromPublisher(buffers)
        );
        future.whenComplete((response, error) -> {
            if (error != null) {
                sink.error(error);
            } else {
                sink.success(checkResult(response));
            }
        });
        // 处理Reactor取消逻辑,避免资源泄漏
        sink.onCancel(future::cancel);
    }).doOnError(this::handleError);
}

同时对应修改handlePart中的调用:

else if (event instanceof FilePartEvent fileEvent) {
    log.debug("!!! FilePartEvent!");
    String filename = accumulator.generateFilename(fileEvent.filename());
    // 不再传递contentLength
    return partEvents.map(PartEvent::content)
            .transform(contentBuffers -> attachmentService.save(filename, contentBuffers))
            .then(Mono.empty());
}

3. 适配Reactor与AWS SDK的线程模型

直接用Mono.fromFuture可能无法正确绑定Reactor的订阅上下文,导致数据流未被触发。改用Mono.create手动处理Future的完成/错误/取消逻辑,确保数据流能被AWS SDK正确订阅(如上述uploadBuffers代码所示)。

4. 避免PartEvent流未被完全消费

在handlePart处理FilePartEvent时,确保partEvents流被正确订阅消费。使用transform操作符将流传递给save方法,而不是直接调用save后返回空Mono,保证整个文件的分块流都被处理。

5. 排查重试触发根源

FilePartEvent重复触发通常是因为上传失败后Spring WebFlux的重试机制导致的。解决上述核心问题后,上传成功就不会再触发重试。如果仍有重试,需检查是否有其他错误(如S3权限、网络问题),可在handleError中打印详细错误日志定位。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 15:25:24