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; } }
关键说明
- 结果字节收集:在事件流中监听
RECORDS类型事件,将每个事件的payload字节写入ByteArrayOutputStream,最终合并为完整的byte数组返回。 - 异步等待机制:通过
CompletableFuture.join()等待异步请求完成,同时用CountDownLatch确保等待到END事件或异常,避免数据未完全收集就返回。 - 事件兼容处理:保留了Stats事件的打印逻辑,和V1的行为保持一致;同时处理异常场景,确保错误能被及时捕获。
内容的提问来源于stack exchange,提问作者user1805280
相关产品推荐
相关产品推荐

