操作Parquet

Hadoop 操作 Parquet 文件:怎么读、怎么写、要注意什么

Parquet 是 Hadoop 生态里最常用的列式存储格式。跟 Text 格式相比,它的核心优势:

  • 压缩率高:列式存储 + 编码压缩,文件体积比 Text 小得多
  • 查询快:只读需要的列,不用整行加载(谓词下推 + 列裁剪)
  • Schema 自带:文件里包含字段名、类型,读的时候不用额外定义

在 Hive、Spark、Impala 里直接用 SQL 就能查,但如果你要自己写 MapReduce 读 Parquet,就得用专门的 API。

添加依赖

<dependency>
    <groupId>org.apache.parquet</groupId>
    <artifactId>parquet-column</artifactId>
    <version>1.8.1</version>
</dependency>
<dependency>
    <groupId>org.apache.parquet</groupId>
    <artifactId>parquet-hadoop</artifactId>
    <version>1.8.1</version>
</dependency>

注意: Parquet 依赖的 Guava 版本如果跟 Hadoop 的不一样,会报 NoSuchMethodError。如果遇到,参考前面问题集锦的处理方式——统一两边 Guava 版本。

读取流程三步走

读取 Parquet 文件的核心步骤是三件事:① 创建输入格式 → ② 创建记录读取器 → ③ 遍历解析每条记录

public class ParquetReaderMapper extends Mapper<Void, Group, NullWritable, Text> {
    
    private Gson gson = new Gson();
    private Text outValue = new Text();
    
    @Override
    protected void map(Void key, Group value, Context context) 
            throws IOException, InterruptedException {
        // 第三步:解析每条记录
        outValue.set(groupToJson(value));
        context.write(NullWritable.get(), outValue);
    }
    
    private String groupToJson(Group group) {
        Map<String, Object> map = new HashMap<>();
        // 按字段名取数,类型要匹配
        map.put("name", group.getString("name", 0));
        map.put("age", group.getInteger("age", 0));
        return gson.toJson(map);
    }
}

但这里有个关键点: 如果你只是这么写 map() 方法,它是不会执行的。Parquet 的 RecordReader 需要手动初始化,所以需要重写 run() 方法来手动控制读取流程。

重写 run() 方法:手动控制读取

Mapper 的默认 run() 方法会调用 map(),但不适用于 Parquet 的 RecordReader。所以你要自己接管整个过程:

@Override
public void run(Context context) throws IOException, InterruptedException {
    setup(context);
    
    // 第一步:创建 ParquetInputFormat,指定 GroupReadSupport
    ParquetInputFormat<Group> inputFormat = new ParquetInputFormat<>(GroupReadSupport.class);
    
    try {
        while (context.nextKeyValue()) {
            InputSplit split = context.getInputSplit();
            
            if (split instanceof FileSplit) {
                // 第二步:创建 RecordReader 并初始化
                TaskAttemptContext attemptContext = new TaskAttemptContextImpl(
                    context.getConfiguration(),
                    context.getTaskAttemptID()
                );
                
                try (RecordReader<Void, Group> reader = inputFormat.createRecordReader(
                        split, attemptContext)) {
                    reader.initialize(split, attemptContext);
                    
                    // 第三步:遍历读取每条记录,调用 map() 处理
                    while (reader.nextKeyValue()) {
                        map(reader.getCurrentKey(), reader.getCurrentValue(), context);
                    }
                }
            }
        }
    } finally {
        cleanup(context);
    }
}

从 Group 对象里取数据

Group 是 Parquet 的”行对象”,根据字段名和位置取数据:

// 取字符串字段
String name = group.getString("name", 0);

// 取整数字段
Integer age = group.getInteger("age", 0);

// 取长整型
Long timestamp = group.getLong("timestamp", 0);

// 取布尔型
Boolean active = group.getBoolean("is_active", 0);

注意: 字段名必须跟 Parquet 文件里写的一模一样,大小写敏感。类型必须匹配,getString() 拿字符串字段,getInteger() 拿 int 字段,拿错会报错。

如果字段可能为空,用 null 做默认值:

String name = group.getString("name", null);

几个常见坑

坑1:字段名/类型不匹配

group.getString("age", 0) 如果 age 在 Parquet schema 里是 int 类型,会报错。字段名写错也一样。

解决办法:parquet-tools 先看一眼 schema:

parquet-tools schema /path/to/file.parquet

坑2:MapReduce 跑起来报 NoSuchMethodError

基本上是 Guava 版本冲突。把 Hadoop 和 Parquet 的 Guava 版本统一即可。

坑3:空字段处理

Parquet 里某个字段为空,直接用 getXxx() 会报 NullPointerException。用带默认值的重载方法:

String maybeNull = group.getString("field", null);

坑4:性能问题

大 Parquet 文件逐行解析转 JSON,如果字段多、行数大,性能会比较差。可以:

  • 只读需要的列(通过 ParquetInputFormatsetReadSupport 配置读哪些列)
  • 输出用二进制格式而不是 JSON
  • 结合 Spark 做批量处理