如何在AWS Lambda中用Java读取S3上的Parquet文件?
问题场景与代码
用户提供的处理SQS事件、读取S3中Parquet文件的代码:
private final S3Client s3Client = S3Client.builder().build(); @Override public Void handleRequest(SQSEvent event, Context context) { for (SQSMessage msg : event.getRecords()) { S3Event s3Event = gson.fromJson(msg.getBody(),S3Event.class); S3EventNotificationRecord s3EventRecord = s3Event.getRecords().get(0); String bucketName = s3EventRecord.getS3().getBucket().getName(); String objectKey = URLDecoder.decode(s3EventRecord.getS3().getObject().getKey(), StandardCharsets.UTF_8); GetObjectRequest getObjectRequest = GetObjectRequest.builder() .bucket(bucketName) .key(bucketName) // 此处存在bug:应使用objectKey而非bucketName .build(); InputStream inputStream = s3Client.getObject(getObjectRequest, ResponseTransformer.toBytes()).asInputStream(); // TODO: 读取Parquet文件中的单条记录 } return null; }
用户需求:从上述InputStream中读取Parquet文件的单条记录,且不想依赖Hadoop或Spark全量库。
可行解决方案
可以使用Apache Parquet官方的parquet-hadoop-core或parquet-avro核心库,只需引入必要模块,无需依赖完整的Hadoop分布式组件,就能直接从InputStream读取Parquet记录。
1. 引入依赖(Maven)
如果用Avro作为数据格式(Parquet常见搭配),可引入以下依赖,排除不必要的Hadoop客户端依赖以实现轻量化:
<dependency> <groupId>org.apache.parquet</groupId> <artifactId>parquet-avro</artifactId> <version>1.14.2</version> <exclusions> <exclusion> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>org.apache.parquet</groupId> <artifactId>parquet-hadoop</artifactId> <version>1.14.2</version> <exclusions> <exclusion> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> </exclusion> </exclusions> </dependency>
2. 读取Parquet记录的代码示例
方式一:使用预定义Avro类(类型安全)
假设Parquet文件对应已生成的Avro实体类(比如User),代码如下:
// 先修正GetObjectRequest的key错误 GetObjectRequest getObjectRequest = GetObjectRequest.builder() .bucket(bucketName) .key(objectKey) .build(); InputStream inputStream = s3Client.getObject(getObjectRequest, ResponseTransformer.toBytes()).asInputStream(); // 构建ParquetReader读取单条记录 try (ParquetReader<User> reader = AvroParquetReader.<User>builder(new Path("dummy")) .withReadSupport(new AvroReadSupport<>()) .withInputStream(inputStream) .build()) { User record; while ((record = reader.read()) != null) { // 处理单条记录,比如发送至目标队列 processAndSendToQueue(record); } } catch (IOException e) { context.getLogger().log("读取Parquet文件失败: " + e.getMessage()); }
方式二:使用GenericRecord处理任意Schema
如果没有预定义的Avro实体类,可通过GenericRecord读取任意Schema的Parquet文件:
try (ParquetReader<GenericRecord> reader = AvroParquetReader.<GenericRecord>builder(new Path("dummy")) .withReadSupport(new AvroReadSupport<>()) .withInputStream(inputStream) .build()) { GenericRecord record; while ((record = reader.read()) != null) { // 按需读取字段并处理 String targetFieldValue = record.get("targetField").toString(); // 执行发送队列逻辑 } } catch (IOException e) { context.getLogger().log("读取Parquet文件失败: " + e.getMessage()); }
关键说明
new Path("dummy")仅为占位符,ParquetReader实际会使用传入的InputStream读取数据,不会访问本地文件系统。- 排除Hadoop客户端依赖后,仅保留Parquet核心处理逻辑,不会引入Hadoop分布式集群相关代码,满足轻量化需求。
内容的提问来源于stack exchange,提问作者Ftwpker
相关产品推荐
相关产品推荐

