在跑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官方提供了AvroKeyAvroValue两种包装类,可以直接当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会调用AvroKeyComparatorcompare(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

如果你的数据结构特别简单,比如只有一个IntWritableText,那完全没必要自定义。如果字段就一两个,而且数据量很小,自定义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的底层比较器。省下来的不止是时间,还有你抓头发的那份心酸。

序列化是分布式计算里的隐形地雷,平时看不见,一碰就炸。希望这篇文章能帮你拆掉这颗雷。