如何解决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
相关产品推荐
相关产品推荐

