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

将S3Client升级为S3AsyncClient时如何将InputStream转为Publisher<ByteBuffer>

AWS S3AsyncClient 异步上传InputStream解决方案

你可以选择两种方案实现InputStream到AsyncRequestBody所需的Publisher<ByteBuffer>转换:

方案1:直接使用高版本SDK内置方法

AWS Java SDK 2.18及更高版本已经内置了AsyncRequestBody.fromInputStream()方法,和同步接口的用法完全一致,直接替换即可:

return asyncClient.putObject(myObjectRequestBuild, AsyncRequestBody.fromInputStream(inputStream, contentSize));

如果你的SDK版本较低无法升级,可使用方案2手动实现转换逻辑。

方案2:基于Reactor手动构造Publisher

AWS SDK 2.x内部默认依赖Reactor响应式库,不需要额外引入依赖,你可以用Flux快速封装InputStream为Publisher<ByteBuffer>:

完整实现代码

import reactor.core.publisher.Flux;
import java.nio.ByteBuffer;

// 原有逻辑不变
URL url = new URL(fileUrl);
String[] fileNameArray = url.getFile().split("\\.");
var uniqueFileName = prepareFileName(fileNameArray[fileNameArray.length -1]);
URLConnection connection = url.openConnection();
long contentSize = connection.getContentLengthLong();
InputStream inputStream = connection.getInputStream();

// 转换InputStream为Publisher<ByteBuffer>
Flux<ByteBuffer> bufferPublisher = Flux.generate(
    // 初始化8KB缓冲区,可根据文件大小调整
    () -> new byte[1024 * 8],
    (buffer, sink) -> {
        try {
            int readLen = inputStream.read(buffer);
            if (readLen == -1) {
                // 流读取完成,结束发布
                sink.complete();
            } else {
                // 发布读取到的字节缓冲
                sink.next(ByteBuffer.wrap(buffer, 0, readLen));
            }
            return buffer;
        } catch (Exception e) {
            // 读取异常向下传递
            sink.error(e);
            return buffer;
        }
    }
// 结束后自动关闭输入流,避免资源泄漏
).doFinally(signal -> {
    try {
        inputStream.close();
    } catch (Exception ignore) {}
});

return asyncClient.putObject(myObjectRequestBuild, AsyncRequestBody.fromPublisher(bufferPublisher));

注意事项

  • 缓冲区大小可根据你的业务场景调整,大文件上传可以适当调大缓冲区,减少IO次数提升效率
  • 如果你不想依赖Reactor,也可以手动实现JDK的Flow.Publisher接口,但需要自己处理背压、异常、订阅等逻辑,成本较高,不推荐

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 03:06:02