Spring Boot 3.0.x异步流式上传文件至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

