展开讲 · 一
SequenceFile 的一个例子
假设 HDFS 上有一百万张小图片。每个文件都要占 NameNode 的内存,每个文件还会单独起一个读取任务。常见的做法是把它们打包进一个 SequenceFile:key 是文件名,value 是图片的字节。
# PySpark:把很多小文件打包成一个 SequenceFile
# key 是文件名,value 是文件内容
pairs = sc.parallelize([
("cat-001.jpg", bytearray(open("cat-001.jpg", "rb").read())),
("cat-002.jpg", bytearray(open("cat-002.jpg", "rb").read())),
])
pairs.saveAsSequenceFile("hdfs:///datasets/cats.seq")
# 读回来:得到 (文件名, 内容) 的键值对
sc.sequenceFile("hdfs:///datasets/cats.seq").take(1)
写出来的文件,里面大致是这样:
SEQ\x06 魔数和版本
org.apache.hadoop.io.Text key 的类
org.apache.hadoop.io.BytesWritable value 的类
压缩: 否 块压缩: 否 压缩标志
<16 字节 sync marker> 文件头结束
[记录长度][key 长度] "cat-001.jpg" <图片字节>
[记录长度][key 长度] "cat-002.jpg" <图片字节>
<16 字节 sync marker> 每隔一段出现一次
[记录长度][key 长度] "cat-003.jpg" <图片字节>
文件头里写的是 Java 的类名,这就是它难以被其他语言读取的原因。sync marker 让一个任务可以从文件中间开始,找到下一条完整的记录。
动手练习 · 约 5 分钟
亲手看一眼 Parquet 的文件尾
需要 Python 和 pyarrow(pip install pyarrow)。把下面的代码存成 demo.py 再运行。
import os
import pyarrow as pa
import pyarrow.csv as pacsv
import pyarrow.parquet as pq
# 1. 造一张 100 万行的表:一列递增的整数,一列只有 4 种取值,一列小数
n = 1_000_000
table = pa.table({
"user_id": pa.array(range(n), pa.int64()),
"country": pa.array(["US", "CN", "IN", "BR"] * (n // 4)),
"amount": pa.array([(i * 7919) % 100_000 / 100 for i in range(n)]),
})
# 2. 分别写成 CSV 和 Parquet,比一比大小
pacsv.write_csv(table, "demo.csv")
pq.write_table(table, "demo.parquet", row_group_size=250_000)
print("CSV ", os.path.getsize("demo.csv") // 1024, "KB")
print("Parquet", os.path.getsize("demo.parquet") // 1024, "KB")
# 3. 读文件尾(footer):不读数据,就能知道文件的结构
meta = pq.ParquetFile("demo.parquet").metadata
print("row group 数:", meta.num_row_groups, " 总行数:", meta.num_rows)
# 4. 看每一列在第一个 row group 里占多少字节、用了什么编码
for i in range(meta.num_columns):
col = meta.row_group(0).column(i)
print(col.path_in_schema, col.total_compressed_size // 1024, "KB", col.encodings)
# 5. 看每个 row group 里 user_id 的最小值和最大值
for g in range(meta.num_row_groups):
stats = meta.row_group(g).column(0).statistics
print("row group", g, "user_id:", stats.min, "到", stats.max)
# 6. 只读一列,并且只要 user_id >= 900000 的行
result = pq.read_table(
"demo.parquet", columns=["amount"], filters=[("user_id", ">=", 900_000)]
)
print("读到", result.num_rows, "行")
在 pyarrow 19.0 上的实际输出:
CSV 18221 KB
Parquet 8991 KB
row group 数: 4 总行数: 1000000
user_id 1249 KB ('PLAIN', 'RLE', 'RLE_DICTIONARY')
country 2 KB ('PLAIN', 'RLE', 'RLE_DICTIONARY')
amount 994 KB ('PLAIN', 'RLE', 'RLE_DICTIONARY')
row group 0 user_id: 0 到 249999
row group 1 user_id: 250000 到 499999
row group 2 user_id: 500000 到 749999
row group 3 user_id: 750000 到 999999
读到 100000 行
- country 列在第一个 row group 里有 25 万个值,只占 3,049 字节(输出里按 KB 取整,显示成 2 KB)。它只有 4 种取值,字典编码把每个值变成 2 bit 的编号,25 万个值约 62 KB;这些编号按固定顺序循环,Snappy 再把它压到约 3 KB。这里游程编码没起作用,因为相邻的值都不相同。
- 想看游程编码的效果,把 country 这一列排序后再写一次:相同的值连在一起,这一列会小到几十个字节。
- 第 5 步打印的最小值和最大值存在文件尾里。第 6 步的过滤条件是 user_id >= 900000,前三个 row group 的最大值都不到 900000,所以整组被跳过,只读了第四组。
- 自己改一改:把 row_group_size 改成 50000 再跑,看 row group 数和文件大小怎么变;把 country 换成每行都不同的字符串,看那一列变成多大。
速查表
考前十分钟过一遍
- CSV / JSON
- 文本,通用,没有类型,整体压缩后不能切分。
- SequenceFile
- 二进制键值对,sync marker 让它可切分,绑死 Java。
- Thrift / Protobuf
- 不是文件格式。字段编号带来 schema 演进。
- Avro
- 行存,schema 在文件头。写入侧主流:Kafka、Iceberg manifest。
- Dremel
- repetition level 和 definition level,嵌套数据按列存。
- RCFile
- Record Columnar:先分行组,组内按列。不认识类型,没有索引。
- ORC
- stripe 加索引加文件尾,认识类型。Hive 体系。
- Parquet
- row group、column chunk、page 三层,文件尾存统计。事实标准。
- Arrow
- 内存里的列式布局,零拷贝。和 Parquet 搭配。
- Iceberg / Delta / Hudi
- 表格式:原子提交、历史版本、改表结构。
- Lance
- 没有 row group,随机读快,加列不重写,自带索引。
- Nimble / Vortex / F3
- 宽表、可扩展编码、GPU。都还新。
留言
说说你的看法
有想法或问题都可以写在这里。留言会立刻显示;想收到回复通知再填邮箱。还没有留言,你可以写第一条。