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

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的误用:

  1. 该类设计用于包装已就绪的阻塞输入流,而非后续动态写入的流。当调用putObject或transferManager.upload时,异步客户端会立即尝试读取请求体,此时body.writeInputStream(stream)尚未执行,导致读取到流的初始异常数据(如管道流的内部标识),形成前缀。
  2. 流的读写时序混乱:异步上传启动后才写入流,导致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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 16:47:10