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

AWS SDK2 Java S3 Select示例:如何获取结果字节数据

AWS SDK2 Java 实现S3 Select并获取结果流字节

问题说明

需要使用AWS SDK2 Java执行S3 Select操作,提取查询结果的字节数据。已有AWS SDK V1的实现参考,但在SDK V2中编写了部分代码后,不清楚如何获取完整的结果流字节。

AWS SDK V1 参考实现

InputStream resultInputStream = result.getPayload().getRecordsInputStream(
    new SelectObjectContentEventVisitor() {
        @Override
        public void visit(SelectObjectContentEvent.StatsEvent event)
        {
            System.out.println(
                "Received Stats, Bytes Scanned: " + event.getDetails().getBytesScanned()
                +  " Bytes Processed: " + event.getDetails().getBytesProcessed());
        }

        /*
         * An End Event informs that the request has finished successfully.
         */
        @Override
        public void visit(SelectObjectContentEvent.EndEvent event)
        {
            isResultComplete.set(true);
            System.out.println("Received End Event. Result is complete.");
        }
    }
);

用户的AWS SDK V2 部分代码

public byte[] getQueryResults() {
    logger.info("V2 query");

    S3AsyncClient s3Client = null;
    s3Client = S3AsyncClient.builder()
            .region(Region.US_WEST_2)
            .build();


    String fileObjKeyName = "upload/" + filePath;

try{
    logger.info("Filepath: " + fileObjKeyName);

        ListObjectsV2Request listObjects = ListObjectsV2Request
                .builder()
                .bucket(Constants.bucketName)
                .build();
          ......


           InputSerialization inputSerialization = InputSerialization.builder().
                    json(JSONInput.builder().type(JSONType.LINES).build()).build();           
           OutputSerialization outputSerialization = OutputSerialization.builder().
                        json(JSONOutput.builder()
                                .build()

                        ).build();



            SelectObjectContentRequest selectObjectContentRequest = SelectObjectContentRequest.builder()
                    .bucket(Constants.bucketName)
                    .key(partFilename)
                    .expression(query)
                    .expressionType(ExpressionType.SQL)
                    .inputSerialization(inputSerialization)
                    .outputSerialization(outputSerialization)
                    .scanRange(ScanRange.builder().start(0L).end(Constants.limitBytes).build())
                    .build();


                final DataHandler handler = new DataHandler();

                CompletableFuture future = s3Client.selectObjectContent(selectObjectContentRequest, handler);

//hold it till we get a end event
                EndEvent endEvent = (EndEvent) handler.receivedEvents.stream()
                        .filter(e -> e.sdkEventType() == SelectObjectContentEventStream.EventType.END)
                        .findFirst()
                        .orElse(null);

//现在,我该如何从这里获取响应字节?
////////---> 问题:如何获取ResultStream字节????

                return <bytes>
}
// 处理器代码
private static  class DataHandler implements SelectObjectContentResponseHandler {
private SelectObjectContentResponse response;
private List receivedEvents = new ArrayList<>();
private Throwable exception;

    @Override
    public void responseReceived(SelectObjectContentResponse response) {
        this.response = response;
    }

    @Override
    public void onEventStream(SdkPublisher<SelectObjectContentEventStream> publisher) {
        publisher.subscribe(receivedEvents::add);
    }

    @Override
    public void exceptionOccurred(Throwable throwable) {
        exception = throwable;
    }

    @Override
    public void complete() {
    }
}  

解决方案:完整实现代码

修改处理器逻辑来收集RecordsEvent中的字节数据,并正确等待异步操作完成,最终合并所有结果字节:

import software.amazon.awssdk.core.SdkBytes;
import software.amazon.awssdk.services.s3.model.*;
import software.amazon.awssdk.services.s3.model.SelectObjectContentEventStream.EventType;
import software.amazon.awssdk.services.s3.async.S3AsyncClient;
import java.io.ByteArrayOutputStream;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

public byte[] getQueryResults() {
    logger.info("V2 query");

    // 复用S3客户端,避免重复创建
    try (S3AsyncClient s3Client = S3AsyncClient.builder()
            .region(Region.US_WEST_2)
            .build()) {

        String fileObjKeyName = "upload/" + filePath;
        logger.info("Filepath: " + fileObjKeyName);

        // 配置输入输出序列化规则
        InputSerialization inputSerialization = InputSerialization.builder()
                .json(JSONInput.builder().type(JSONType.LINES).build())
                .build();
        OutputSerialization outputSerialization = OutputSerialization.builder()
                .json(JSONOutput.builder().build())
                .build();

        // 构建Select请求
        SelectObjectContentRequest selectRequest = SelectObjectContentRequest.builder()
                .bucket(Constants.bucketName)
                .key(partFilename)
                .expression(query)
                .expressionType(ExpressionType.SQL)
                .inputSerialization(inputSerialization)
                .outputSerialization(outputSerialization)
                .scanRange(ScanRange.builder().start(0L).end(Constants.limitBytes).build())
                .build();

        // 使用自定义处理器收集结果
        final ResultCollectingHandler handler = new ResultCollectingHandler();
        CompletableFuture<Void> future = s3Client.selectObjectContent(selectRequest, handler);

        // 等待异步请求完成
        future.join();

        // 检查请求是否异常
        if (handler.getException() != null) {
            throw new RuntimeException("S3 Select执行失败", handler.getException());
        }

        // 返回合并后的结果字节
        return handler.getResultBytes();

    } catch (Exception e) {
        logger.error("S3 Select操作出错", e);
        throw new RuntimeException(e);
    }
}

// 自定义处理器:负责收集结果字节、处理各类事件
private static class ResultCollectingHandler implements SelectObjectContentResponseHandler {
    private final ByteArrayOutputStream resultStream = new ByteArrayOutputStream();
    private final CountDownLatch endLatch = new CountDownLatch(1);
    private SelectObjectContentResponse response;
    private Throwable exception;

    @Override
    public void responseReceived(SelectObjectContentResponse response) {
        this.response = response;
    }

    @Override
    public void onEventStream(SdkPublisher<SelectObjectContentEventStream> publisher) {
        publisher.subscribe(event -> {
            switch (event.sdkEventType()) {
                case RECORDS:
                    // 提取结果数据字节并写入输出流
                    RecordsEvent recordsEvent = (RecordsEvent) event;
                    resultStream.write(recordsEvent.payload().asByteArray());
                    break;
                case STATS:
                    // 打印统计信息,与V1逻辑一致
                    StatsEvent statsEvent = (StatsEvent) event;
                    System.out.println("Received Stats, Bytes Scanned: " + statsEvent.details().bytesScanned()
                            + " Bytes Processed: " + statsEvent.details().bytesProcessed());
                    break;
                case END:
                    // 收到结束事件,释放等待锁
                    endLatch.countDown();
                    System.out.println("Received End Event. Result is complete.");
                    break;
                default:
                    // 忽略其他类型事件
                    break;
            }
        });
    }

    @Override
    public void exceptionOccurred(Throwable throwable) {
        this.exception = throwable;
        endLatch.countDown(); // 异常时也释放锁,避免永久等待
    }

    @Override
    public void complete() {
        try {
            // 等待结束事件或异常,设置超时时间
            endLatch.await(5, TimeUnit.MINUTES);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            this.exception = e;
        }
    }

    // 获取最终合并的结果字节数组
    public byte[] getResultBytes() {
        return resultStream.toByteArray();
    }

    // 获取异常信息
    public Throwable getException() {
        return exception;
    }
}

关键说明

  1. 结果字节收集:在事件流中监听RECORDS类型事件,将每个事件的payload字节写入ByteArrayOutputStream,最终合并为完整的byte数组返回。
  2. 异步等待机制:通过CompletableFuture.join()等待异步请求完成,同时用CountDownLatch确保等待到END事件或异常,避免数据未完全收集就返回。
  3. 事件兼容处理:保留了Stats事件的打印逻辑,和V1的行为保持一致;同时处理异常场景,确保错误能被及时捕获。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 12:30:46