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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 21:05:24