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

如何读取AWS S3中的Parquet文件并转换为Java对象?

从S3 Parquet文件流转换为Java对象的实现方案

以下是两种常用的实现方式,基于Apache Parquet生态工具完成转换:


方式一:使用Parquet-Jackson直接映射到普通Java对象

这种方式适合Parquet结构与普通Java实体类匹配的场景,无需额外定义Schema。

1. 添加Maven依赖

<dependencies>
    <dependency>
        <groupId>org.apache.parquet</groupId>
        <artifactId>parquet-jackson</artifactId>
        <version>1.14.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.parquet</groupId>
        <artifactId>parquet-hadoop</artifactId>
        <version>1.14.0</version>
    </dependency>
</dependencies>

2. 定义匹配Parquet结构的Java实体类

确保字段名(或通过@JsonProperty映射)、数据类型与Parquet文件列完全匹配,且提供无参构造函数:

import com.fasterxml.jackson.annotation.JsonProperty;

public class Recommendation {
    @JsonProperty("user_id")
    private String userId;
    @JsonProperty("item_id")
    private String itemId;
    private double score;

    public Recommendation() {}

    // Getter & Setter
    public String getUserId() { return userId; }
    public void setUserId(String userId) { this.userId = userId; }
    public String getItemId() { return itemId; }
    public void setItemId(String itemId) { this.itemId = itemId; }
    public double getScore() { return score; }
    public void setScore(double score) { this.score = score; }
}

3. 读取S3流并转换为Java对象

利用Parquet-Jackson的Reader解析输入流:

import org.apache.parquet.hadoop.ParquetReader;
import org.apache.parquet.jackson.ParquetJacksonReader;
import org.apache.parquet.io.InputFile;
import org.apache.parquet.hadoop.util.HadoopInputFile;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.conf.Configuration;

// 从已获取的S3Object中获取输入流
SeekableInputStream inputStream = object.getObjectContent();

try {
    // 构建Parquet输入文件对象
    Configuration conf = new Configuration();
    InputFile inputFile = HadoopInputFile.fromStream(inputStream, new Path("s3://temp/recommendations.parquet"), conf);

    // 创建Reader并逐行读取转换
    try (ParquetReader<Recommendation> reader = ParquetJacksonReader.builder(Recommendation.class, inputFile).build()) {
        Recommendation rec;
        while ((rec = reader.read()) != null) {
            // 处理转换后的Java对象,比如业务逻辑处理
            System.out.printf("用户ID: %s, 物品ID: %s, 评分: %.2f%n", rec.getUserId(), rec.getItemId(), rec.getScore());
        }
    }
} catch (Exception e) {
    e.printStackTrace();
} finally {
    // 确保输入流关闭
    inputStream.close();
}

方式二:基于Avro Schema转换(适合有预定义Schema的场景)

如果Parquet文件是基于Avro Schema生成的,可使用Parquet-Avro工具链完成转换。

1. 添加Maven依赖

<dependencies>
    <dependency>
        <groupId>org.apache.parquet</groupId>
        <artifactId>parquet-avro</artifactId>
        <version>1.14.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.parquet</groupId>
        <artifactId>parquet-hadoop</artifactId>
        <version>1.14.0</version>
    </dependency>
</dependencies>

2. 生成Avro对应的Java类

通过Avro的Maven插件,根据Avro Schema文件(.avsc)生成Java实体类,示例插件配置:

<build>
    <plugins>
        <plugin>
            <groupId>org.apache.avro</groupId>
            <artifactId>avro-maven-plugin</artifactId>
            <version>1.11.0</version>
            <executions>
                <execution>
                    <phase>generate-sources</phase>
                    <goals>
                        <goal>schema</goal>
                    </goals>
                    <configuration>
                        <sourceDirectory>${project.basedir}/src/main/resources/avro</sourceDirectory>
                        <outputDirectory>${project.basedir}/src/main/java</outputDirectory>
                    </configuration>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>

3. 读取并转换

import org.apache.parquet.avro.AvroParquetReader;
import org.apache.parquet.hadoop.ParquetReader;
import org.apache.parquet.io.InputFile;
import org.apache.parquet.hadoop.util.HadoopInputFile;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.conf.Configuration;

SeekableInputStream inputStream = object.getObjectContent();

try {
    Configuration conf = new Configuration();
    InputFile inputFile = HadoopInputFile.fromStream(inputStream, new Path("s3://temp/recommendations.parquet"), conf);

    // 假设生成的Avro Java类为RecommendationAvro
    try (ParquetReader<RecommendationAvro> reader = AvroParquetReader.<RecommendationAvro>builder(inputFile).build()) {
        RecommendationAvro rec;
        while ((rec = reader.read()) != null) {
            // 处理Avro实体对象
        }
    }
} catch (Exception e) {
    e.printStackTrace();
} finally {
    inputStream.close();
}

关键注意事项

  • 确保Parquet依赖版本统一,避免版本冲突
  • Java实体类必须与Parquet列的名称、数据类型严格匹配,字段名不匹配时用@JsonProperty映射
  • 优先使用try-with-resources自动关闭Reader和流资源,避免内存泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 10:36:01