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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:41:22