在跑Hadoop任务的时候,你是不是也碰到过这种怪事:Map阶段明明很快,到了Shuffle阶段就瘫了,尤其是Sort那一块,日志就像卡住一样,几个小时都不往前挪。很多人第一反应是集群太小,或者数据倾斜,但折腾半天后才发现,问题可能藏在一个不起眼的小类里——你自定义的那个Writable。
一、一个让人挠头的排序慢问题
有位朋友跟我吐槽,他写了个用户信息类用来当Map的输出key,里面有名字、年龄、工资三个字段。数据量不大,也就几千万条,结果排序阶段比预期慢了几十倍。他把代码翻来覆去看了半天,最后怀疑是自定义Writable的问题。我们先看看他的代码,说不定你也能看到自己的影子。
// 技术栈:Java + Hadoop MapReduce
import org.apache.hadoop.io.Writable;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
/**
* 自定义Writable,代表一条用户记录
*/
public class PersonWritable implements Writable, Comparable<PersonWritable> {
private String name; // 姓名
private int age; // 年龄
private long salary; // 工资
// 空构造器,反序列化时必须要有
public PersonWritable() {
}
public PersonWritable(String name, int age, long salary) {
this.name = name;
this.age = age;
this.salary = salary;
}
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public int getAge() {
return age;
}
public void setAge(int age) {
this.age = age;
}
public long getSalary() {
return salary;
}
public void setSalary(long salary) {
this.salary = salary;
}
/**
* 序列化:把对象写入二进制输出流
*/
@Override
public void write(DataOutput out) throws IOException {
out.writeUTF(name); // 注意这里:writeUTF会先写两个字节的长度,再写内容
out.writeInt(age);
out.writeLong(salary);
}
/**
* 反序列化:从二进制输入流读回字段
*/
@Override
public void readFields(DataInput in) throws IOException {
this.name = in.readUTF(); // 读的时候每次都会新建一个String对象
this.age = in.readInt();
this.salary = in.readLong();
}
/**
* 比较逻辑:按名字、年龄、工资的顺序排序
*/
@Override
public int compareTo(PersonWritable o) {
int result = this.name.compareTo(o.name);
if (result != 0) {
return result;
}
result = Integer.compare(this.age, o.age);
if (result != 0) {
return result;
}
return Long.compare(this.salary, o.salary);
}
}
代码看着没问题,可跑起来就是慢。到底慢在哪儿?咱们得把Hadoop的序列化和排序机制扒一层皮才能看清楚。
二、序列化机制到底坑在哪
2.1 排序的时候,Hadoop在干什么
MapReduce的排序通常发生在Map输出写到磁盘以及Shuffle合并的时候。为了排序,Hadoop必须把序列化好的字节流进行反复比较。默认情况下,比较器是WritableComparator,它接收两个字节数组,然后试图从中还原出Writable对象,再调用compareTo方法。
但这里有个关键点:如果你没有为自定义Writable注册一个专门的RawComparator,那么Hadoop就只能走“反序列化再比较”的笨路子。每一次比较,都要把字节流重新变成PersonWritable对象,然后调用compareTo。几千万条数据,几亿次比较,每次都要重新创建对象、读取字符串、分配内存,这速度能快吗?
2.2 反序列化是个“无底洞”
我们再看一眼刚才的readFields方法。你别小看一个readUTF,它内部要读取两个字节判断长度,然后根据长度创建字节数组,再用DataInputStream的魔改版UTF-8解析成一个String。这一套下来,涉及多次数组复制和字符解码。
更糟糕的是,compareTo需要比较两个对象,所以每次比较至少要创建两个PersonWritable对象。对象创建多了,JVM的GC就遭了殃。你会发现,任务日志里Full GC频繁,CPU飙得老高,但排序就是不走。
2.3 序列化后的字节也不省心
writeUTF并不是最紧凑的格式。它用两字节存长度,然后对每个字符按UTF-8编码。如果名字大部分是英文,还勉强能接受,一旦有中文或者表情符,每个字符可能占3~4个字节,整个键的大小立马膨胀。键大了,I/O传输量变大,比较时读取的数据量也变大,雪上加霜。
有人可能会说:我写一个RawComparator,直接基于字节比较,不反序列化,不就行了?思路对,但实现起来没那么轻松。你得自己解析二进制布局,还得处理字符串和可变长度整数,稍有偏差结果就错了。下面这种硬核玩法,真不是一般人能hold住的。
// 技术栈:Java + Hadoop MapReduce
import org.apache.hadoop.io.WritableComparator;
import org.apache.hadoop.io.WritableUtils;
import java.io.IOException;
/**
* 自定义RawComparator的雏形,想直接比较字节数组
* 但需要手工解析每个字段,代码又长又容易出错
*/
public class PersonRawComparator extends WritableComparator {
protected PersonRawComparator() {
super(PersonWritable.class);
}
@Override
public int compare(byte[] b1, int s1, int l1, byte[] b2, int s2, int l2) {
try {
// 先跳过name的UTF-8长度标记(2字节),再跳过name的字节数
int nameLen1 = WritableUtils.decodeVIntSize(b1[s1]) + readUnsignedShort(b1, s1);
// 算了,这里每个字段都要自己算偏移量,
// 字符串编码还可能是变长的,一不留神就bug
return super.compare(b1, s1, l1, b2, s2, l2);
} catch (IOException e) {
throw new IllegalArgumentException(e);
}
}
private int readUnsignedShort(byte[] bytes, int offset) {
return ((bytes[offset] & 0xFF) << 8) | (bytes[offset + 1] & 0xFF);
}
}
看到没?为了跳过名字字段,你得知道Hadoop内部怎么写UTF-8长度。实际实现要比这个复杂得多,光是字符串的比较就要考虑长度前缀、编码方式,还有可能的负数等等。与其花大力气造这轮子,不如换一个天生就支持高效比较的序列化方案。
三、更优替代:让Avro来救场
在Hadoop生态里,有一个非常友好的序列化框架叫Avro。它把数据模式单独写在schema里面,序列化后的二进制非常紧凑,而且最重要的是——Hadoop为Avro专门提供了一个基于字节的RawComparator,叫AvroKeyComparator。它可以直接比较两个Avro序列化后的字节流,完全不需要反序列化成对象。这意味着排序的时候,你根本不需要创建对象,速度自然快出一个量级。
3.1 定义一个Avro模式
Avro的模式是用JSON写的,比如我们继续用那个用户信息模型:
{
"type": "record",
"name": "Person",
"namespace": "com.example",
"fields": [
{"name": "name", "type": "string"},
{"name": "age", "type": "int"},
{"name": "salary", "type": "long"}
]
}
这段schema描述了一个记录类型,里面有三个字段。你可以把这个JSON单独保存成person.avsc文件,也可以用的时候直接用字符串解析。Avro的二进制编码比writeUTF紧凑得多,字符串带了长度信息,但用的是变长编码,对于常见字符更省空间。
3.2 在MapReduce里使用AvroKey
Avro官方提供了AvroKey和AvroValue两种包装类,可以直接当Map输出Key和Value。下面这个Mapper的例子,会把文本行转成Avro记录,然后输出一个AvroKey。
// 技术栈:Java + Hadoop MapReduce + Avro
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.mapred.AvroKey;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
import java.io.IOException;
public class PersonAvroMapper extends Mapper<LongWritable, Text, AvroKey<GenericRecord>, NullWritable> {
private Schema schema;
private AvroKey<GenericRecord> outKey = new AvroKey<>();
@Override
protected void setup(Context context) {
// 在实际项目中,schema一般是从.avsc文件里读的,这里为了演示直接写字符串
String schemaJson = "{\n"
+ " \"type\": \"record\",\n"
+ " \"name\": \"Person\",\n"
+ " \"namespace\": \"com.example\",\n"
+ " \"fields\": [\n"
+ " {\"name\": \"name\", \"type\": \"string\"},\n"
+ " {\"name\": \"age\", \"type\": \"int\"},\n"
+ " {\"name\": \"salary\", \"type\": \"long\"}\n"
+ " ]\n"
+ "}";
schema = new Schema.Parser().parse(schemaJson);
}
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
// 假设输入格式是:姓名,年龄,工资
String[] parts = value.toString().split(",");
// 创建一个泛型记录
GenericRecord record = new GenericData.Record(schema);
record.put("name", parts[0]);
record.put("age", Integer.parseInt(parts[1]));
record.put("salary", Long.parseLong(parts[2]));
// 包成AvroKey并输出
outKey.datum(record);
context.write(outKey, NullWritable.get());
}
}
然后,在Driver里需要告诉Hadoop使用Avro的序列化器,并且给AvroKey设置好对应的schema。这里的配置方式根据Avro版本略有不同,但大方向是明确声明Map输出Key的schema。
// 技术栈:Java + Hadoop MapReduce + Avro
import org.apache.avro.Schema;
import org.apache.avro.mapred.AvroKey;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.Configured;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;
public class PersonSortJob extends Configured implements Tool {
@Override
public int run(String[] args) throws Exception {
Configuration conf = getConf();
Job job = Job.getInstance(conf, "person sort with avro");
job.setJarByClass(getClass());
// 指定Mapper、Reducer(这里省略Reducer,实际肯定有)
job.setMapperClass(PersonAvroMapper.class);
// 注意:如果要使用AvroKey作为输出Key,必须把map输出key类型设为AvroKey
job.setMapOutputKeyClass(AvroKey.class);
job.setMapOutputValueClass(NullWritable.class);
// 设置Avro的Map输出Key schema
String schemaJson = "{\n"
+ " \"type\": \"record\",\n"
+ " \"name\": \"Person\",\n"
+ " \"namespace\": \"com.example\",\n"
+ " \"fields\": [\n"
+ " {\"name\": \"name\", \"type\": \"string\"},\n"
+ " {\"name\": \"age\", \"type\": \"int\"},\n"
+ " {\"name\": \"salary\", \"type\": \"long\"}\n"
+ " ]\n"
+ "}";
Schema schema = new Schema.Parser().parse(schemaJson);
// 这一行的作用是告诉Hadoop,AvroKey里的schema长什么样子
org.apache.avro.mapred.AvroJob.setMapOutputKeySchema(job, schema);
// 输入输出路径
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
return job.waitForCompletion(true) ? 0 : 1;
}
public static void main(String[] args) throws Exception {
int exitCode = ToolRunner.run(new PersonSortJob(), args);
System.exit(exitCode);
}
}
当排序发生时,Hadoop会调用AvroKeyComparator的compare(byte[], ...)方法。这个比较器直接解析Avro的二进制格式,按照schema里的字段顺序一个一个比较字节。整个过程不会new任何Java对象,也不会触发GC。你可以想象成,之前是拆两个箱子看内容,现在直接隔着一层玻璃看箱子里贴的标签,高下立判。
3.3 其他替代方案值不值得考虑
除了Avro,还有Protobuf、Thrift等序列化框架,它们也都有自己的比较器。不过Protobuf和Thrift需要额外生成Java类,代码体积更大。Avro生成的是schema,配合Hadoop更自然。还有Hadoop自带的SequenceFile,但它主要影响存储格式,不直接解决Key排序性能问题。所以如果你的排序性能卡在Writable上,Avro算是目前性价比最高的选择了。
四、应用场景与优缺点
4.1 什么时候可以继续用自定义Writable
如果你的数据结构特别简单,比如只有一个IntWritable或Text,那完全没必要自定义。如果字段就一两个,而且数据量很小,自定义Writable也够用。另外,你已经在一套稳定的作业里用惯了,不想引入新依赖,那优化优化也能跑。
4.2 自定义Writable的典型毛病
- 序列化格式不紧凑:
writeUTF字符串存储效率低,数字用定长编码,白白浪费空间。 - 缺乏高效的RawComparator:每次比较都走反序列化,CPU和内存开销极大。
- 容易写错equals和hashCode:一旦没写,Hadoop里各种Map、Reduce缓存就会出奇奇怪怪的bug。
- 手动编码维护成本高:字段一变,write、readFields、compareTo全要改。
4.3 Avro的闪光点
- 二进制紧凑:变长整数+zigzag编码,数字小的时候只占一个字节,能省不少空间。
- 天生有RawComparator:
AvroKeyComparator直接基于字节比较,排序效率几乎翻倍。 - Schema演化方便:多个作业可以共享同一个schema,字段可增可删,兼容性好。
- 与Hadoop集成度极高:官方支持MapReduce,也支持Spark和Flink。
4.4 Avro的坑也得说
- 多了个schema文件要维护,对开发流程有要求。
- 泛型记录看起来不直观,调试时得靠工具转成文本。
- 如果Avro版本和Hadoop版本不匹配,会有序列化异常,需要踩坑。
五、注意事项与总结
当你再次面对“排序慢成狗”的问题时,先不要急着扩容。回头看看你的自定义Writable,有没有实现RawComparator,序列化是否紧凑,readFields里有没有高频的字符串操作。这些都是潜在的性能杀手。
如果你的作业对排序时效要求很高,我强烈建议切换到Avro。你只需要定义一个schema,把Map输出Key包成AvroKey,然后挪个配置,剩下的交给Avro的底层比较器。省下来的不止是时间,还有你抓头发的那份心酸。
序列化是分布式计算里的隐形地雷,平时看不见,一碰就炸。希望这篇文章能帮你拆掉这颗雷。
评论
围绕“常见自定义Writable导致MapReduce排序效率骤降,序列化机制深挖与更优替代方案”参与讨论