将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
相关产品推荐
相关产品推荐

