使用AWS Java SDK V2调用S3Select仅获部分JSON响应问题排查
问题描述
- 场景:Spring Boot应用中使用AWS Java SDK V2实现S3Select,查询S3存储桶中的Parquet文件
- 问题:仅能获取部分查询结果,控制台输出的JSON内容仅65k字符
- 已尝试操作:取消Eclipse控制台偏好设置中的“Limit console output”选项,问题仍未解决
核心问题分析
你的代码逻辑存在关键错误:S3Select的查询结果会通过多个RecordsEvent分片返回,但你只调用了findFirst()获取第一个事件的内容,自然只能拿到部分数据,和控制台输出限制无关。
修复后的完整代码
import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider; import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; import software.amazon.awssdk.core.async.SdkPublisher; import software.amazon.awssdk.regions.Region; import software.amazon.awssdk.services.s3.S3AsyncClient; import software.amazon.awssdk.services.s3.model.CompressionType; import software.amazon.awssdk.services.s3.model.EndEvent; import software.amazon.awssdk.services.s3.model.ExpressionType; import software.amazon.awssdk.services.s3.model.InputSerialization; import software.amazon.awssdk.services.s3.model.JSONOutput; import software.amazon.awssdk.services.s3.model.OutputSerialization; import software.amazon.awssdk.services.s3.model.ParquetInput; import software.amazon.awssdk.services.s3.model.RecordsEvent; import software.amazon.awssdk.services.s3.model.SelectObjectContentEventStream; import software.amazon.awssdk.services.s3.model.SelectObjectContentEventStream.EventType; import software.amazon.awssdk.services.s3.model.SelectObjectContentRequest; import software.amazon.awssdk.services.s3.model.SelectObjectContentResponse; import software.amazon.awssdk.services.s3.model.SelectObjectContentResponseHandler; public class ParquetSelect { private static final String BUCKET_NAME = "<bucket-name>"; private static final String KEY = "<object-key>"; private static final String QUERY = "select * from S3Object s"; public static S3AsyncClient s3; public static void selectObjectContent() { Handler handler = new Handler(); SelectQueryWithHandler(handler).join(); // 遍历所有RecordsEvent,拼接完整结果 StringBuilder fullResult = new StringBuilder(); for (SelectObjectContentEventStream event : handler.receivedEvents) { if (event.sdkEventType() == EventType.RECORDS) { RecordsEvent recordsEvent = (RecordsEvent) event; fullResult.append(recordsEvent.payload().asUtf8String()); } } System.out.println(fullResult.toString()); } private static CompletableFuture<Void> SelectQueryWithHandler(SelectObjectContentResponseHandler handler) { InputSerialization inputSerialization = InputSerialization.builder() .parquet(ParquetInput.builder().build()) .compressionType(CompressionType.NONE) .build(); OutputSerialization outputSerialization = OutputSerialization.builder() .json(JSONOutput.builder().build()) .build(); SelectObjectContentRequest select = SelectObjectContentRequest.builder() .bucket(BUCKET_NAME) .key(KEY) .expression(QUERY) .expressionType(ExpressionType.SQL) .inputSerialization(inputSerialization) .outputSerialization(outputSerialization) .build(); return s3.selectObjectContent(select, handler); } private static class Handler implements SelectObjectContentResponseHandler { private SelectObjectContentResponse response; private List<SelectObjectContentEventStream> 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() { } } }
额外注意事项
- 确保
S3AsyncClient已正确初始化(代码中s3变量未初始化,需补充类似s3 = S3AsyncClient.builder().region(Region.CN_NORTH_1).credentialsProvider(StaticCredentialsProvider.create(AwsBasicCredentials.create("accessKey", "secretKey"))).build();的逻辑) - 如果查询结果量极大,建议不要直接拼接成字符串输出到控制台,改为写入文件或分块处理,避免内存溢出
内容的提问来源于stack exchange,提问作者Barani
相关产品推荐
相关产品推荐

