上一篇讲了 RPC。本篇讲 RPC 和数据存储都依赖的两个底层机制——序列化和压缩。

序列化与压缩常被合并介绍。这两个机制其实是不同维度——序列化解决"对象怎么变字节流"(结构化),压缩解决"字节流怎么变更小的字节流"(空间优化)。两者经常组合使用——序列化后的字节流再压缩存储或传输。准确的说法是:Hadoop 提供多种序列化格式(Writable / Protobuf / Avro),与多种压缩格式(Gzip / BZip2 / Snappy / LZ4 / Zstd)正交组合,按场景选择最优组合。

本篇只抓一个问题:Writable 与 Java Serializable 的差别、Protobuf 与 Avro 各自的取舍、Snappy vs Zstd vs Gzip 的压缩比/速度权衡、压缩在 Hadoop 各个环节的具体应用位置。

序列化的本质

序列化(serialization)解决"对象 → 字节流"和"字节流 → 对象"的双向转换。分布式系统的所有跨进程数据交换都依赖序列化——RPC 参数、HDFS 文件内容、Shuffle 中间数据、作业配置传递。

序列化机制的设计目标可以拆成五个维度:

1
2
3
4
5
1. 速度(CPU 开销) - 序列化和反序列化的 CPU 占用
2. 紧凑度(大小) - 输出字节流的大小,影响网络带宽和存储
3. 兼容性 - 字段变化时旧版本能否反序列化新版本
4. 跨语言 - 是否支持 Java 之外的语言
5. 可拆分性 - 反序列化能否从字节流任意位置开始(Hadoop MapReduce 关键需求)

没有一种序列化在所有维度都最优——这是典型的工程权衡。Hadoop 支持多种序列化格式让用户按场景选择。

Writable:Hadoop 原生序列化

Writable 是 Hadoop 自己的序列化机制,所有 Hadoop 内部数据结构(Text、IntWritable、LongWritable、ArrayWritable 等)都实现 Writable 接口:

1
2
3
4
public interface Writable {
void write(DataOutput out) throws IOException;
void readFields(DataInput in) throws IOException;
}

实现 Writable 的类自己定义怎么序列化到 DataOutput、怎么从 DataInput 反序列化。这种"自己实现"的方式让序列化速度极快(直接调用方法,无反射开销),同时让字节流紧凑(没有 Java Serializable 的类名等冗余信息)。

Writable 与 Java Serializable 的对比:

1
2
3
4
5
6
7
8
9
                       Java Serializable              Writable
───────────────── ─────────────────────────── ──────────────────────
序列化方式 反射 + 默认字段写入 用户自己实现 write/readFields
字节流大小 大(含类名、字段元数据) 小(只含数据字节)
CPU 开销 高(反射 + 元数据处理) 低(直接调用方法)
跨语言 不支持(Java 专属) 不支持(Java 专属)
版本兼容 弱(serialVersionUID 严格匹配) 弱(用户自己处理)
拆分性 不支持 取决于实现
典型场景 Java RMI、单进程 Hadoop MapReduce(内置)

Writable 的设计动机是 Java Serializable 不适合 Hadoop 的场景——MapReduce 处理海量记录,每条记录的序列化开销累加显著,Writable 把这个开销降到最低。

Writable 的简单示例(自定义类型):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
public class PairWritable implements Writable {
private Text first;
private IntWritable second;

public PairWritable() {
first = new Text();
second = new IntWritable();
}

@Override
public void write(DataOutput out) throws IOException {
first.write(out);
second.write(out);
}

@Override
public void readFields(DataInput in) throws IOException {
first.readFields(in);
second.readFields(in);
}
}

注意 write 和 readFields 的字段顺序必须一致——这是手写序列化的责任,框架不检查。如果顺序写错(write 是 first/second,readFields 是 second/first),反序列化会得到错误数据但不报错,调试困难。

Protobuf:跨版本兼容的选择

Protobuf 是 Google 设计的跨语言序列化格式。Hadoop 2.x 起 RPC 协议改用 Protobuf(上一篇讲过),Hadoop 3.x 起 RPC 强制 Protobuf。

Protobuf 的核心抽象是 schema(.proto 文件)。Schema 定义消息结构,编译器生成各语言的桩代码:

1
2
3
4
5
6
message User {
required int32 id = 1;
required string name = 2;
optional string email = 3;
repeated string tags = 4;
}

编译器生成 Java、Python、C++、Go 等语言的类,每个类知道怎么序列化/反序列化自己的字段。Protobuf 的 wire format 紧凑(varint 编码 + field tag),通常比 Java Serializable 小 3-5 倍。

Protobuf 的版本兼容机制是 optional 字段——schema 添加新字段不会破坏旧版本反序列化(旧版本忽略未知字段)。这让 Hadoop 可以平滑滚动升级。

Protobuf 的代价:

CPU 开销比 Writable 高约 10-20%(解析 field tag、varint 解码)。

Schema 必须 precompile——.proto 文件编译生成桩代码,不支持运行时动态 schema。

不适合"任意位置拆分"——Protobuf 消息长度在前缀,必须从消息边界开始解析(不是任意位置)。这让 Protobuf 不适合 HDFS 文件内容的直接存储,但适合 RPC(每条 RPC 是完整消息)。

Avro:Schema 嵌入的序列化

Avro 是 Hadoop 生态系统的另一种序列化格式,由 Doug Cutting(Hadoop 创始人)2009 年创建。Avro 与 Protobuf 类似(schema-based、跨语言),但有几个关键差异:

1
2
3
4
5
6
7
8
                              Protobuf                          Avro
───────────────── ────────────────── ─────────────────────
Schema 位置 Schema 编译进代码 Schema 嵌入数据文件
(schema + data 同放)
字段 ID 字段编号(1, 2, 3...) 字段名(无编号)
版本兼容 optional 字段保证 schema evolution 规则
MapReduce 友好 否(消息长度前缀,不可拆分) 是(block 结构,可拆分)
典型用途 RPC 数据存储(Hive / Pig / Parquet 底层)

Avro 的核心创新是"schema + data 同放"。Avro 文件包含 schema 头 + 多个 data block,每个 block 是一批按 schema 序列化的记录。这种结构让 Avro 文件天然支持 MapReduce 拆分——每个 Map Task 处理一个 block 边界开始。

Avro 解析时以 writer schema 与 reader schema 配对:writer 中存在、reader 中删除的字段会被忽略;reader 新增而旧 writer 没有的字段才需要 default,否则读取旧数据会失败。

Avro 在 Hive / Pig 生态里广泛使用,作为 Sequence File 的替代品。不要把这一点外推成 “HBase 内部元数据也默认用 Avro”:HBase 的核心存储路径是 WAL、MemStore、HFile 和 ZooKeeper / Master 元数据协调,序列化实现要按 HBase 版本源码分别核对。

三种序列化的选择

Hadoop 生态里同时存在三种序列化,按场景选择:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
RPC 参数和返回值:Protobuf(Hadoop 2.x 起强制)
- 跨版本兼容最重要
- 消息大小适中,性能差异不显著

MapReduce 中间数据(Shuffle):Writable 或 Protobuf
- 默认 Writable(速度优先)
- 自定义 Partitioner 时用 Writable(接口要求)

HDFS 数据文件:Avro / Parquet / ORC
- Avro:行式存储,schema 嵌入,适合 OLTP-like 查询
- Parquet:列式存储,OLAP 友好
- ORC:Hive 原生列式存储

用户自定义对象:Writable
- MapReduce 框架的 Mapper/Reducer 输出必须是 Writable
- 不用 Writable 意味着需要自定义 Serialization 类(额外工作)

实际项目中混合使用很常见——MapReduce 内部用 Writable(速度优先),最终数据存到 HDFS 用 Avro / Parquet(schema 友好),RPC 用 Protobuf(兼容性)。

压缩:节省空间与传输

压缩(compression)解决"字节流怎么变更小"。Hadoop 集群里的数据从 GB 到 PB 级,每字节的存储成本和网络带宽成本都被压缩放大。

Hadoop 支持的压缩格式:

1
2
3
4
5
6
7
          压缩比倾向       速度倾向          可拆分      备注
Gzip 高 慢 否 兼容性好,CPU 开销高
BZip2 高 很慢 是 可拆分,但吞吐通常不是首选
LZO 中 快 是(需索引) 早期 Snappy 替代
Snappy 中 很快 否 常用于中间数据
LZ4 中 很快 否 偏速度
Zstd 中高 可调 否 2.9.0 / 3.0.0-alpha2 起支持

注意"可拆分"这一列。MapReduce 处理大文件时按 InputSplit 切分,每个 Split 由一个 Map Task 处理。如果文件被 Gzip 压缩,从压缩流中间开始读不可能(Gzip 必须从头解压到目标位置)。所以 Gzip 文件无法拆分——整个文件作为一个 InputSplit,Map Task 数 = 1,无法并行。

BZip2 是个例外——它在压缩流里嵌入了 stream marker(同步点),从 marker 位置可以开始独立解压。这让 BZip2 文件可以按 marker 拆分,每个 Split 从 marker 开始。

Snappy / LZ4 / Zstd 都不可拆分——但这通常不构成问题,因为这些格式主要用于 MapReduce 中间数据(中间数据按 record 边界拆分,不需要 stream marker)或者用于 Sequence File / Parquet(这些格式内部自带 block 结构)。

Hadoop 从 2.9.0 和 3.0.0-alpha2 起提供 Zstd codec。它通常在压缩比与速度之间取得较好折中,但是否优于 Snappy、LZ4 或 Gzip 取决于数据、CPU 预算、native library 可用性和下游兼容性,不能写成 Hadoop 3.x 的统一推荐。

压缩在 Hadoop 各个环节的应用

压缩不只是"压缩文件"——Hadoop 在多个环节都应用压缩,每个环节的策略不同:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
1. HDFS 静态数据(at-rest)
- 用 Sequence File / Avro / Parquet / ORC 的内置 codec
- 典型:Parquet + Snappy(列式存储 + 快压缩)
- 影响:存储成本(节省幅度要按真实压缩率、单价和保留周期计算)

2. MapReduce 中间数据(Shuffle)
- 配置 io.compression.codecs(注册 codec 列表)
- 配置 mapreduce.map.output.compress=true(默认 false)
- 配置 mapreduce.map.output.compress.codec=org.apache.hadoop.io.compress.SnappyCodec
- 影响:Shuffle 网络流量(具体降幅取决于数据可压缩性)

3. MapReduce 最终输出
- 配置 mapreduce.output.fileoutputformat.compress=true
- 配置 mapreduce.output.fileoutputformat.compress.codec
- 配置 mapreduce.output.fileoutputformat.compress.type(BLOCK / RECORD)
- 影响:HDFS 存储占用

4. Shuffle HTTP 传输
- 中间数据压缩后通过 HTTP 传给 Reduce
- 自动跟随中间数据压缩配置

5. RPC
- 默认不压缩(CPU 敏感)
- SASL QOP 可提供认证、完整性或加密,但这不是压缩

生产集群的典型配置:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
<!-- mapred-site.xml -->
<property>
<name>mapreduce.map.output.compress</name>
<value>true</value>
</property>
<property>
<name>mapreduce.map.output.compress.codec</name>
<value>org.apache.hadoop.io.compress.SnappyCodec</value>
</property>
<property>
<name>mapreduce.output.fileoutputformat.compress</name>
<value>true</value>
</property>
<property>
<name>mapreduce.output.fileoutputformat.compress.codec</name>
<value>org.apache.hadoop.io.compress.GzipCodec</value>
</property>
<property>
<name>mapreduce.output.fileoutputformat.compress.type</name>
<value>BLOCK</value>
</property>

这种组合让 Map 输出在 Shuffle 时用 Snappy 快压缩(CPU 敏感),最终输出用 Gzip 高压缩比(存储敏感)。BLOCK 类型让 Sequence File 按 block 压缩(不是按 record),更高效。

实验:观察压缩效果

实验状态:UNVERIFIED_RUNTIME。下面是验证步骤,本轮没有连接真实 Hadoop 集群执行。

提交一个 word count 作业,配置不同的 Map 输出压缩,观察 Shuffle 网络流量变化:

1
2
3
4
5
6
7
8
# 不压缩
hadoop jar hadoop-mapreduce-examples-*.jar wordcount /input /output1

# Snappy 压缩
hadoop jar hadoop-mapreduce-examples-*.jar wordcount -D mapreduce.map.output.compress=true -D mapreduce.map.output.compress.codec=org.apache.hadoop.io.compress.SnappyCodec /input /output2

# Gzip 压缩
hadoop jar hadoop-mapreduce-examples-*.jar wordcount -D mapreduce.map.output.compress=true -D mapreduce.map.output.compress.codec=org.apache.hadoop.io.compress.GzipCodec /input /output3

在 ApplicationMaster Web UI 的 Counter 标签页观察三个作业的 Reduce shuffle bytes

1
2
3
不压缩:        以实际 Counter 为准
Snappy 压缩: 通常少于不压缩,具体比例取决于数据
Gzip 压缩: 通常压缩率更高,但 CPU 开销也更高

Shuffle 字节减少的代价是 CPU 占用上升。对比作业总执行时间时,要同时看 Reduce shuffle bytes、CPU 时间和作业 wall time。没有在同一集群、同一输入上跑基准前,不能把 Snappy、Gzip 或 Zstd 写成绝对最优。

文件级压缩效果用 hadoop fs -du 观察:

1
2
3
4
5
6
7
8
9
10
# 写入未压缩 Sequence File
hadoop jar hadoop-mapreduce-examples-*.jar wordcount -D mapreduce.output.fileoutputformat.compress=false /input /output-raw

# 写入 Snappy Sequence File
hadoop jar hadoop-mapreduce-examples-*.jar wordcount -D mapreduce.output.fileoutputformat.compress=true -D mapreduce.output.fileoutputformat.compress.codec=org.apache.hadoop.io.compress.SnappyCodec /input /output-snappy

# 查看大小
hadoop fs -du -h /output-raw /output-snappy
# /output-raw 以实际输出为准
# /output-snappy 以实际输出为准

压缩是否省钱要按真实压缩率、CPU 开销、存储单价和保留周期计算。对 PB 级集群,这个差异足够大,但不能用未实测的百分比直接外推成固定金额。

模式提炼

Hadoop 的序列化与压缩体现的设计模式:

1
2
3
4
5
6
7
模式:多序列化格式 + 多压缩格式 + 正交组合 + 按场景选择

- 序列化层抽象出 Serialization 接口,支持 Writable / Protobuf / Avro 等多种实现
- 压缩层抽象出 CompressionCodec 接口,支持 Gzip / BZip2 / Snappy / LZ4 / Zstd 等
- 序列化格式与压缩格式正交——任意组合
- 不同场景选不同组合:RPC(Protobuf + 无压缩)、Shuffle(Writable + Snappy)、存储(Avro + Gzip)
- 拆分性(splittable)是 MapReduce 场景的关键约束

这个模式不只是 Hadoop。数据库系统同样支持多种存储格式(行存 / 列存)和多种压缩(字典 / RLE / Delta),让用户按查询模式选择。Kafka 的 message 格式(v1 / v2)也涉及类似权衡——压缩放在 record batch 级别,让批量压缩效率高。

数据湖格式(Parquet / ORC / Iceberg / Hudi)是这种模式的现代化版本——内置列式存储 + 多种压缩 + schema evolution,把序列化和压缩组合做成了文件格式的一部分。

工程迁移表

Hadoop 概念 Parquet ORC Kafka Database
Writable(行存) Record Row message v2 row
Avro(schema 嵌入) File footer file footer header schema catalog
Sequence File(Hadoop 自带) Parquet File ORC File log segment tablespace
Snappy codec Snappy Snappy Snappy LZ4 / Snappy
可拆分性 row group stripe partition partition
压缩组合 column-wise + codec column-wise + codec batch + codec page-level

注意 Parquet / ORC 这一列。这两个列式存储格式是 Hadoop 生态系统的现代化产物——它们把"列式存储 + 压缩 + schema 演化"做成了一等公民。Parquet 在 Spark 生态流行,ORC 在 Hive 生态流行,两者都是 Hadoop 文件格式(Sequence File、Avro)的进化版本。

常见误解

误解一:“Writable 就是 Java Serializable 的 Hadoop 版本”。Writable 与 Java Serializable 是完全不同的机制——Java Serializable 用反射自动序列化,Writable 让用户自己实现。Writable 更快更紧凑,但要求用户手写序列化代码。

误解二:“压缩总是好的”。压缩节省存储和网络,但消耗 CPU。如果作业是 CPU 密集型(大量计算),压缩反而拖慢——节省的 I/O 时间 < 压缩时间。生产经验:I/O 密集型作业(ETL)开压缩,CPU 密集型作业(机器学习推理)关闭压缩。

误解三:“Gzip 是最佳压缩格式”。Gzip 历史悠久且兼容性好,但速度通常不如 Snappy、LZ4 或 Zstd。新部署应按数据、CPU、可拆分需求和读取端兼容性做基准测试,不能无条件指定 Zstd。

误解四:“压缩文件无法 MapReduce 拆分”。只有 Gzip / Snappy / LZ4 / Zstd 等流式压缩不可拆分。BZip2 有 stream marker 可以拆分。Sequence File / Avro / Parquet / ORC 即使压缩也可以拆分(block 结构)。所以"压缩 + 可拆分"不矛盾,关键是选对格式。

误解五:“Avro 比 Protobuf 新所以更好”。Avro(2009)和 Protobuf(Google 2001 内部,2008 开源)都历史悠久。两者设计目标不同——Avro 为数据存储设计(schema 嵌入文件),Protobuf 为 RPC 设计(schema 编译进代码)。不能简单比较"哪个更好",要看场景。

练习

  1. 提交一个 MapReduce 作业,对比启用 Snappy 压缩和不压缩两种情况下 Shuffle 字节数和作业总时间。验证 Snappy 在你的集群上是节省时间还是增加时间。

  2. 写一个自定义 Writable 类(例如包含两个字段的 Pair),故意让 write 和 readFields 的字段顺序不一致,运行 MapReduce 观察反序列化错误的奇怪表现。然后修复顺序重新运行。

  3. apache/hadoop 源码里找到 Writable.javaCompressionCodec.javaCodecPool.javaZStandardCodec.java(hadoop-common-project 模块),观察序列化和压缩接口的核心定义。

  4. 思考题:如果让 Hadoop 把所有内部序列化从 Writable 切换到 Protobuf(包括 MapReduce 中间数据),会带来什么好处和坏处?为什么 Hadoop 团队没有做这种切换?

系列导航

序号 主题 状态
00 导读:节点总会失败
01 HDFS 架构与三层切分
02 文件写入路径与流水线
03 文件读取路径与副本选择
04 NameNode 内存模型与启动恢复
05 HDFS HA 与脑裂防御
06 HDFS 3.x 演进与纠删码
07 YARN 架构与三方契约
08 YARN 资源模型、Container 与 NodeLabel
09 YARN 调度器对比:FIFO、Capacity、Fair
10 YARN 应用程序生命周期
11 YARN HA 与 Federation
12 YARN Timeline Service v2
13 MapReduce 编程模型与分而治之
14 MapReduce Shuffle 全流程
15 MRv2 on YARN ApplicationMaster 与 Task Attempt
16 Hadoop RPC 协议栈 上一篇
17 序列化与压缩 本篇
18 Hadoop 安全 Kerberos Token ProxyUser 下一篇
19 监控与运维 Metrics JMX 日志聚合
20 Hadoop 生态 Hive HBase Pig
21 Hadoop 与对象存储 Kubernetes 演进对比
22 Hadoop 设计遗产从 GFS MapReduce 到云原生

参考资料