如何基于RxJava与Vert.x将S3Object转换为Observable
解决方案
你需要完成两个核心操作来实现需求:一是把阻塞的S3同步调用适配成异步模型(避免阻塞Vert.x事件循环),二是读取S3Object的内容并包装成Observable<String>。下面分两种场景给出实现方案:
推荐方案:使用AWS异步S3客户端
AWS提供了原生异步的S3客户端,完美适配RxJava和Vert.x的异步编程模型,无需手动处理线程池,是最优选择:
import io.reactivex.rxjava3.core.Observable; import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; import software.amazon.awssdk.regions.Region; import software.amazon.awssdk.services.s3.S3AsyncClient; import software.amazon.awssdk.services.s3.model.GetObjectRequest; import software.amazon.awssdk.core.ResponseBytes; import java.nio.charset.StandardCharsets; public Observable<String> getCustomerFroms3(String orderId) { // 初始化异步S3客户端 S3AsyncClient s3AsyncClient = S3AsyncClient.builder() .credentialsProvider(StaticCredentialsProvider.create( AwsBasicCredentials.create( "你的AccessKey", "你的SecretKey" ) )) .region(Region.AP_SOUTH_1) .build(); // 构建S3对象获取请求 GetObjectRequest getObjectRequest = GetObjectRequest.builder() .bucket("database8") .key(orderId) .build(); // 将异步调用结果转换为Observable<String> return Observable.fromCompletionStage( s3AsyncClient.getObject(getObjectRequest, ResponseBytes::asByteArray) ).map(responseBytes -> { // 将字节数组转为UTF-8编码的JSON字符串 String jsonContent = new String(responseBytes.asByteArray(), StandardCharsets.UTF_8); // 关闭响应流,避免资源泄漏 responseBytes.response().close(); return jsonContent; }); }
兼容方案:使用旧版同步S3客户端(不推荐)
如果因历史原因必须使用旧版同步客户端,一定要把阻塞操作放到Vert.x的worker线程池,避免阻塞事件循环:
import io.reactivex.rxjava3.core.Observable; import io.vertx.core.Vertx; import com.amazonaws.auth.BasicAWSCredentials; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.regions.Regions; import com.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.AmazonS3ClientBuilder; import com.amazonaws.services.s3.model.S3Object; import java.io.BufferedReader; import java.io.InputStreamReader; import java.nio.charset.StandardCharsets; import java.util.stream.Collectors; public Observable<String> getCustomerFroms3(String orderId) { Vertx vertx = Vertx.currentContext().owner(); // 获取当前Vert.x实例 return Observable.create(emitter -> { // 将阻塞操作提交到Vert.x worker线程池 vertx.executeBlocking(future -> { AmazonS3 s3client = AmazonS3ClientBuilder.standard() .withCredentials(new AWSStaticCredentialsProvider( new BasicAWSCredentials("你的AccessKey", "你的SecretKey") )) .withRegion(Regions.AP_SOUTH_1) .build(); try { S3Object s3Object = s3client.getObject("database8", orderId); // 读取S3Object的输入流转为字符串 try (BufferedReader reader = new BufferedReader( new InputStreamReader(s3Object.getObjectContent(), StandardCharsets.UTF_8) )) { String jsonContent = reader.lines().collect(Collectors.joining(System.lineSeparator())); future.complete(jsonContent); } finally { // 必须关闭S3Object,释放底层资源 s3Object.close(); } } catch (Exception e) { future.fail(e); } }, false, result -> { if (result.succeeded()) { emitter.onNext(result.result().toString()); emitter.onComplete(); } else { emitter.onError(result.cause()); } }); }); }
关键注意事项
- 资源释放:无论哪种方案,都必须关闭S3的响应流/对象,否则会导致连接泄漏
- 异步适配:Vert.x中绝对不能在事件循环线程执行阻塞操作,必须使用异步客户端或worker线程池
- 编码一致性:明确使用UTF-8编码读取文件,避免因系统默认编码导致JSON乱码
内容的提问来源于stack exchange,提问作者Vivek Sharma
相关产品推荐
相关产品推荐

