AWS S3 SDK v2异步上传JSON文件出现前缀与数据损坏问题
问题:AWS S3 SDK v2异步上传JSON数据出现前缀与数据损坏
- 现象:使用
S3AsyncClient或TransferManager上传JSON流到S3后,文件开头出现随机字符串前缀,且JSON数据不完整/损坏;仅同步S3Client上传正常。 - 示例损坏文件内容:
A63 [{"id":1,"name":"apple","category":"fruit","price":0.5},{"id":2,"name":"banana","category":"fruit","price":0.3},{"id":3,"name":"carrot","category":"vegetable","price":0
- 额外现象:前缀随文件大小变化(14MB文件前缀为
800000,空文件为2);尝试临时文件上传、指定内容大小均无效。
相关代码
S3Service实现类
import software.amazon.awssdk.core.async.BlockingInputStreamAsyncRequestBody; import software.amazon.awssdk.core.sync.RequestBody; import software.amazon.awssdk.services.s3.S3AsyncClient; import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.PutObjectRequest; import software.amazon.awssdk.services.s3.model.PutObjectResponse; import software.amazon.awssdk.services.s3.model.S3Exception; import software.amazon.awssdk.services.s3.presigner.S3Presigner; import software.amazon.awssdk.services.s3.presigner.model.GetObjectPresignRequest; import software.amazon.awssdk.services.s3.presigner.model.PresignedGetObjectRequest; import software.amazon.awssdk.transfer.s3.S3TransferManager; import software.amazon.awssdk.transfer.s3.model.CompletedUpload; import software.amazon.awssdk.transfer.s3.model.Upload; public class S3ServiceImpl implements S3Service { @Override public CompletedUpload uploadStream(String key, InputStream stream) throws FileNotFoundException { BlockingInputStreamAsyncRequestBody body = AsyncRequestBody.forBlockingInputStream(null); // 'null' indicates a stream will be provided later. S3TransferManager transferManager = S3TransferManager.builder() .s3Client(s3AsyncClient) .build(); Upload upload = transferManager.upload(builder -> builder .requestBody(body) .putObjectRequest(req -> req.bucket(bucketName).key(key)) .build()); body.writeInputStream(stream); return upload.completionFuture().join(); } @Override public CompletableFuture<PutObjectResponse> uploadStream(String key, InputStream stream) { BlockingInputStreamAsyncRequestBody body = AsyncRequestBody.forBlockingInputStream(null); // 'null' indicates a stream will be provided later. CompletableFuture<PutObjectResponse> responseFuture = s3AsyncClient.putObject(r -> r.bucket(bucketName).key(key), body); body.writeInputStream(stream); return responseFuture; } }
导出逻辑
@Override @Async public void export(Query query) { try { PipedInputStream pipedInputStream = new PipedInputStream(); PipedOutputStream pipedOutputStream = new PipedOutputStream(pipedInputStream); executorService.submit(() -> exportToPipe(query, pipedOutputStream)); CompletedUpload completedUpload = s3Service.uploadStream(task.getFileName(), pipedInputStream); } catch (Exception e) { log.error("export ERROR {}", "", e); } } private void exportToPipe(Query query, PipedOutputStream pipedOutputStream) { try (pipedOutputStream; OutputStreamWriter writer = new OutputStreamWriter(new BufferedOutputStream(pipedOutputStream), StandardCharsets.UTF_8); Stream<Entity> entityStream= mongoOperations.stream(query, Entity.class)) { serializeEntityStream(entityStream, writer); } catch (IOException e) {} }
问题分析
核心错误在于BlockingInputStreamAsyncRequestBody的误用:
- 该类设计用于包装已就绪的阻塞输入流,而非后续动态写入的流。当调用
putObject或transferManager.upload时,异步客户端会立即尝试读取请求体,此时body.writeInputStream(stream)尚未执行,导致读取到流的初始异常数据(如管道流的内部标识),形成前缀。 - 流的读写时序混乱:异步上传启动后才写入流,导致SDK读取到不完整的数据流,最终造成JSON损坏。
解决方案
方案1:正确使用异步请求体(推荐)
放弃BlockingInputStreamAsyncRequestBody,改用AsyncRequestBody.fromInputStream直接包装输入流,同时修正读写时序:
修改S3Service异步上传方法
@Override public CompletedUpload uploadStream(String key, InputStream stream) { // 直接用输入流创建AsyncRequestBody,无需提前初始化空对象 AsyncRequestBody body = AsyncRequestBody.fromInputStream(stream); try (S3TransferManager transferManager = S3TransferManager.builder() .s3Client(s3AsyncClient) .build()) { Upload upload = transferManager.upload(builder -> builder .requestBody(body) .putObjectRequest(req -> req.bucket(bucketName).key(key)) .build()); return upload.completionFuture().join(); } } @Override public CompletableFuture<PutObjectResponse> uploadStream(String key, InputStream stream) { AsyncRequestBody body = AsyncRequestBody.fromInputStream(stream); return s3AsyncClient.putObject(r -> r.bucket(bucketName).key(key), body); }
修正导出逻辑的时序
确保写入线程先启动,再触发上传:
@Override @Async public void export(Query query) { try { PipedInputStream pipedInputStream = new PipedInputStream(); PipedOutputStream pipedOutputStream = new PipedOutputStream(pipedInputStream); // 先提交写入任务,确保写入线程启动 Future<?> writeFuture = executorService.submit(() -> exportToPipe(query, pipedOutputStream)); // 触发上传 CompletedUpload completedUpload = s3Service.uploadStream(task.getFileName(), pipedInputStream); // 等待写入完成,确保流正确关闭 writeFuture.get(); } catch (Exception e) { log.error("export ERROR {}", "", e); } }
方案2:同步客户端适配异步场景(简化版)
如果异步SDK的问题难以排查,可在异步业务方法中使用同步S3Client配合线程池,兼顾业务异步性与上传稳定性:
@Override public CompletableFuture<CompletedUpload> uploadStream(String key, InputStream stream) { return CompletableFuture.supplyAsync(() -> { try (S3Client s3Client = S3Client.create()) { s3Client.putObject(PutObjectRequest.builder() .bucket(bucketName) .key(key) .build(), RequestBody.fromInputStream(stream)); return CompletedUpload.builder().build(); } }, executorService); }
关键注意事项
- 所有流必须通过
try-with-resources自动关闭,避免流泄漏导致数据不完整。 - 不要用
BlockingInputStreamAsyncRequestBody处理动态写入的流,该类仅适用于已就绪的阻塞输入流。 - 若无法提前获取流的总大小,直接使用
AsyncRequestBody.fromInputStream(stream)即可,SDK会自动处理分块上传。
内容的提问来源于stack exchange,提问作者Med-KG
相关产品推荐
相关产品推荐

