0%

操作Parquet

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

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

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

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

添加依赖

1
2
3
4
5
6
7
8
9
10
<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 文件的核心步骤是三件事:① 创建输入格式 → ② 创建记录读取器 → ③ 遍历解析每条记录

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
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。所以你要自己接管整个过程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
@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 的”行对象”,根据字段名和位置取数据:

1
2
3
4
5
6
7
8
9
10
11
// 取字符串字段
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 做默认值:

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

几个常见坑

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

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

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

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

坑2:MapReduce 跑起来报 NoSuchMethodError

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

坑3:空字段处理

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

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

坑4:性能问题

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

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

欢迎关注我的其它发布渠道