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

如何解决Java版AWS SDK异步上传S3文件时的CancellationException

解决AWS SDK for Java流式上传S3时的CancellationException问题

问题背景

需要将压缩后生成的文件内容上传至AWS S3,由于压缩前无法预知文件大小,参照官方流式上传示例使用BlockingInputStreamAsyncRequestBody编写代码,但执行时抛出CancellationException: subscription has been cancelled错误。

环境信息:

  • AWS SDK版本:2.28.7
  • AWS CRT版本:0.33.2
  • Java版本:17

原错误代码:

BlockingInputStreamAsyncRequestBody body =
    AsyncRequestBody.forBlockingInputStream(null); // 'null' indicates a stream will be provided later.

CompletableFuture<PutObjectResponse> responseFuture =
    s33CrtAsyncClient.putObject(r -> r.bucket(bucketName).key(key), body);

// AsyncExampleUtils.randomString() returns a random string up to 100 characters.
String randomString = AsyncExampleUtils.randomString();
logger.info("random string to upload: {}: length={}", randomString, randomString.length());

// Provide the stream of data to be uploaded.
body.writeInputStream(new ByteArrayInputStream(randomString.getBytes()));

PutObjectResponse response = responseFuture.join(); // Wait for the response.

错误堆栈:

java.util.concurrent.CancellationException: subscription has been cancelled.
    at software.amazon.awssdk.utils.async.SimplePublisher.lambda$doProcessQueue$9(SimplePublisher.java:286)
    at software.amazon.awssdk.utils.async.SimplePublisher$FailureMessage.get(SimplePublisher.java:413)
    at software.amazon.awssdk.utils.async.SimplePublisher$FailureMessage.access$700(SimplePublisher.java:397)
    at software.amazon.awssdk.utils.async.SimplePublisher.doProcessQueue(SimplePublisher.java:260)
    at software.amazon.awssdk.utils.async.SimplePublisher.processEventQueue(SimplePublisher.java:224)
    at software.amazon.awssdk.utils.async.SimplePublisher.send(SimplePublisher.java:128)
    at software.amazon.awssdk.utils.async.InputStreamConsumingPublisher.doBlockingWrite(InputStreamConsumingPublisher.java:58)
    at software.amazon.awssdk.core.async.BlockingInputStreamAsyncRequestBody.writeInputStream(BlockingInputStreamAsyncRequestBody.java:95)

错误原因

CRT异步客户端调用putObject后会立即启动请求流程,尝试订阅请求体获取数据。但原代码先发起请求、再延迟提供InputStream,客户端等待超时后自动取消订阅,后续调用writeInputStream时就会触发取消异常。

解决方法

方法1:提前准备InputStream再发起请求

如果压缩后的内容可先写入内存生成InputStream,直接用该流创建AsyncRequestBody,避免延迟提供:

import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.CompletableFuture;
import java.util.zip.GZIPOutputStream;
import software.amazon.awssdk.core.async.AsyncRequestBody;
import software.amazon.awssdk.services.s3.model.PutObjectResponse;

// 模拟压缩逻辑:将原始数据压缩为InputStream
ByteArrayOutputStream baos = new ByteArrayOutputStream();
try (GZIPOutputStream gzipOut = new GZIPOutputStream(baos)) {
    // 写入原始数据到压缩流(替换为实际业务数据)
    gzipOut.write("原始数据内容".getBytes(StandardCharsets.UTF_8));
} catch (IOException e) {
    throw new RuntimeException("压缩数据失败", e);
}
ByteArrayInputStream compressedStream = new ByteArrayInputStream(baos.toByteArray());

// 直接用已准备好的流创建请求体
AsyncRequestBody body = AsyncRequestBody.forInputStream(compressedStream);

// 发起上传请求并等待结果
CompletableFuture<PutObjectResponse> responseFuture = 
    s33CrtAsyncClient.putObject(r -> r.bucket(bucketName).key(key), body);
PutObjectResponse response = responseFuture.join();

方法2:边压缩边上传(适合大文件)

如果原始数据量较大,不适合全部加载到内存,可通过自定义Publisher实现边压缩边上传:

import java.io.IOException;
import java.io.OutputStream;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.CompletableFuture;
import java.util.zip.GZIPOutputStream;
import org.reactivestreams.Subscriber;
import software.amazon.awssdk.core.async.AsyncRequestBody;
import software.amazon.awssdk.services.s3.model.PutObjectResponse;

AsyncRequestBody body = AsyncRequestBody.fromPublisher(subscriber -> {
    try (GZIPOutputStream gzipOut = new GZIPOutputStream(new OutputStream() {
        @Override
        public void write(byte[] b, int off, int len) throws IOException {
            // 将压缩后的字节发送给S3客户端
            subscriber.onNext(ByteBuffer.wrap(b, off, len));
        }

        @Override
        public void write(int b) throws IOException {
            subscriber.onNext(ByteBuffer.wrap(new byte[]{(byte) b}));
        }
    })) {
        // 写入原始数据到压缩流(可替换为从文件/网络读取的流式数据)
        gzipOut.write("大体积原始数据内容".getBytes(StandardCharsets.UTF_8));
    } catch (IOException e) {
        // 通知客户端上传失败
        subscriber.onError(e);
        return;
    }
    // 通知客户端数据发送完成
    subscriber.onComplete();
});

CompletableFuture<PutObjectResponse> responseFuture = 
    s33CrtAsyncClient.putObject(r -> r.bucket(bucketName).key(key), body);
PutObjectResponse response = responseFuture.join();

注意事项

  • 避免使用BlockingInputStreamAsyncRequestBody的延迟流方式,除非能确保在调用putObject前提供InputStream。
  • 大文件上传优先使用方法2的流式处理,避免内存溢出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 21:53:13