0%

序列化

Hadoop 序列化:为什么不用 Java 自带 Serializable?

在分布式系统里,数据要在节点之间网络传输,也要落盘存储。序列化就是干这个的——把内存里的对象变成字节流,到了目的地再还原回来。

Java 自带 Serializable,那 Hadoop 为什么还要搞一套自己的序列化?

因为 Java 原生的序列化太重了——它会把类的元数据(类名、父类信息、字段签名等)也序列化进去,字节流体积大、速度慢。 在海量数据场景下,这种开销是灾难性的。

Hadoop 的 Writable 接口是轻量级序列化的解决方案:只序列化数据本身,不附带任何元数据。 收方和发方约定好”数据的顺序和类型是什么”,然后直接读、直接写。

Writable 长什么样?

Writable 接口就两个方法:

1
2
3
4
public interface Writable {
void write(DataOutput out) throws IOException; // 对象 → 字节流
void readFields(DataInput in) throws IOException; // 字节流 → 对象
}

对比 Java Serializable:

Java Serializable Hadoop Writable
序列化内容 数据 + 类元数据 仅数据
字节流大小
速度
实现复杂度 自动(implements 就行) 手写 write/readFields
跨语言支持 不支持 不支持

所以 Writable 是”用开发者的少量手写代码,换取生产环境的大量性能提升”。

常用 Writable 类型

Hadoop 封装了常用 Java 类型的 Writable 版本:

Java 类型 Writable 类型 备注
int IntWritable 4 字节
long LongWritable 8 字节
float FloatWritable 4 字节
double DoubleWritable 8 字节
boolean BooleanWritable 1 字节
String Text UTF-8 变长编码
byte[] BytesWritable 字节数组
null NullWritable 占位,不占空间

MapReduce 的 Key/Value 必须实现 Writable(或 WritableComparable)。

自定义 Writable:手写一个 UserInfo

假设你要在 MapReduce 里传递用户信息:

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
public class UserInfo implements Writable {
private long userId;
private String userName;
private int age;

// 无参构造函数(必须!反射创建对象用)
public UserInfo() {}

public UserInfo(long userId, String userName, int age) {
this.userId = userId;
this.userName = userName;
this.age = age;
}

@Override
public void write(DataOutput out) throws IOException {
out.writeLong(userId);
Text.writeString(out, userName); // Text 工具类写字符串
out.writeInt(age);
}

@Override
public void readFields(DataInput in) throws IOException {
userId = in.readLong();
userName = Text.readString(in);
age = in.readInt();
}
}

三个关键点:

  1. 读写顺序必须严格一致——write 先写 userId,readFields 就先读 userId。错一个顺序,数据全乱。
  2. 无参构造函数必须有——MapReduce 用反射创建对象,没有无参构造会报错。
  3. 字符串用 Text 工具类读写——Text.writeString()Text.readString() 处理 UTF-8 变长编码,比自己手写 writeUTF 更高效。

需要排序?用 WritableComparable

MapReduce 的 Shuffle 阶段会对 Key 做排序。如果自定义类型要做 Key,需要实现 WritableComparable(Writable + Comparable):

1
2
3
4
5
6
7
8
9
public class UserInfo implements WritableComparable<UserInfo> {
// 同上...

@Override
public int compareTo(UserInfo other) {
// 先按 userId 升序
return Long.compare(this.userId, other.userId);
}
}

排序顺序决定了 Shuffle 后的分组顺序,也影响 GroupingComparator 的行为。

Writable 的局限:什么时候该换 Avro?

Writable 好用,但它有三个硬伤:

1. 不能跨语言

Writable 是 Java 专属的。如果你的下游是 Python 写的程序,没法反序列化 Writable。

2. 版本兼容性差

UserInfo 里加一个 email 字段,老的序列化数据就读不了了。因为没有字段标识,收方不知道”多了一个字段”。

3. 手写代码多

每个自定义类型都要手写 write/readFields,字段多了很烦。

Avro 是更现代的序列化方案:

  • Schema 驱动:数据格式单独定义(JSON 格式的 Schema),数据不含字段描述
  • Schema 演化:新旧 Schema 可以兼容——加字段、删字段、改默认值都可以
  • 跨语言:Schema 可以生成 Java、Python、C++ 等语言的类

Avro 示例

定义 Schema(user.avsc):

1
2
3
4
5
6
7
8
9
{
"type": "record",
"name": "UserInfo",
"fields": [
{"name": "userId", "type": "long"},
{"name": "userName", "type": "string"},
{"name": "age", "type": "int"}
]
}

在 MapReduce 里用 Avro:

1
2
3
4
5
6
7
8
// 读 Avro 文件
job.setInputFormatClass(AvroKeyInputFormat.class);

// Mapper 里直接拿 GenericRecord
public void map(AvroKey<GenericRecord> key, ...) {
GenericRecord user = key.datum();
long userId = (Long) user.get("userId");
}

加了新字段怎么办?

1
2
// 新 Schema 加上 email,且有默认值
{"name": "email", "type": "string", "default": ""}

老的 Avro 数据用新 Schema 读,email 字段自动填充 "",不会报错。

Writable vs Avro:怎么选?

场景 推荐方案 原因
MapReduce 内部 Key/Value Writable 轻量、性能最好
数据要跨语言消费 Avro 多语言支持
Schema 可能频繁变化 Avro 支持向后兼容
追求极致性能 Writable 无元数据开销
长期存储(数据湖) Avro / Parquet 支持演化、可压缩
数据量小、简单场景 Writable 无需额外工具

实际项目中的组合用法:

  • MapReduce 内部 Key 用 Writable(性能优先)
  • HDFS 上长期存储的数据用 Avro 或 Parquet(跨语言 + Schema 演化)

Writable 性能优化几点经验

  1. 复用对象,别 new

TextIntWritable 这些对象在 Map/Reduce 的循环里不要反复 new。定义成成员变量,每次 set() 新值:

1
2
3
4
5
6
private Text outKey = new Text();
private IntWritable outValue = new IntWritable(1);

// map() 里复用
outKey.set(word);
context.write(outKey, outValue);
  1. 字符串用 Text,不用 String

Text 是 Hadoop 自己的变长字符串编码,序列化效率比 String 高。

  1. 可变长度字段放前面

write 顺序里,变长字段(Text、BytesWritable)放前面,定长字段放后面,有利于流式读取。

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