4.2 Partition 分区 本节摘要:分区是 Shuffle 的分拣口——每条中间键值对在 Map 端就被盖上"归属哪个 Reduce 任务"的章。默认的 HashPartitioner 用键的哈希对 Reduce 数取模,简单但可能倾斜。本节推演哈希分区的数学、演示数据倾斜的成因与症状、给出按业务键自定义分区的完整代码(含抽样 TotalOrderPartitioner),并讨论 Reduce 数量这个旋钮该怎么拧。 分区在流水线上的位置 回看 4.1 节的六步:分区发生在环形缓冲区溢写之前——每条 map 输出在进入缓冲区时就带上了分区号,溢写排序的规则是"分区号在前、键在后",最终 Map 端输出文件按分区切成段。Reduce 任务 i 只从所有 Map 输出里拷贝第 i 段。
本节摘要:分区是 Shuffle 的分拣口——每条中间键值对在 Map 端就被盖上"归属哪个 Reduce 任务"的章。默认的 HashPartitioner 用键的哈希对 Reduce 数取模,简单但可能倾斜。本节推演哈希分区的数学、演示数据倾斜的成因与症状、给出按业务键自定义分区的完整代码(含抽样 TotalOrderPartitioner),并讨论 Reduce 数量这个旋钮该怎么拧。
回看 4.1 节的六步:分区发生在环形缓冲区溢写之前——每条 map 输出在进入缓冲区时就带上了分区号,溢写排序的规则是"分区号在前、键在后",最终 Map 端输出文件按分区切成段。Reduce 任务 i 只从所有 Map 输出里拷贝第 i 段。所以分区决定了数据的最终归属地,也决定了每个 Reduce 任务的负载。
// 框架内置的默认实现 全部逻辑只有一行 public class HashPartitioner<K, V> extends Partitioner<K, V> { public int getPartition(K key, V value, int numPartitions) { return (key.hashCode() & Integer.MAX_VALUE) % numPartitions; } }
这一行值得逐个运算符拆开:
key.hashCode():键的哈希,Java 的 int,可能是负数;& Integer.MAX_VALUE:按位与上 0x7FFFFFFF,把符号位清零——负哈希折成正数,等价于取绝对值但不受 -2147483648 溢出 bug 影响;% numPartitions:对 Reduce 任务数取模,结果落在 0 到 numPartitions-1。于是"同键必同分区"由哈希函数的性质保证:相等对象哈希相等(这就是 3.2 节强调 hashCode 与 equals 语义一致的原因);而"同分区尽量等量"则取决于键的分布。前者是正确性约束,铁律;后者只是性能期望,常常落空。
拿 3.1 节的部门工资例子手工推演:设 Reduce 数为 2,"研发"的哈希折正后是 35,35 % 2 = 1,进 1 号分区;"市场"折正后是 64,64 % 2 = 0,进 0 号分区。五条记录分拣如下:
研发 → hashCode 35 → 分区 1 研发 → hashCode 35 → 分区 1 市场 → hashCode 64 → 分区 0 市场 → hashCode 64 → 分区 0 研发 → hashCode 35 → 分区 1 0 号分区:市场 9000 市场 11000 → Reduce 0 1 号分区:研发 12000 研发 15000 研发 13000 → Reduce 1
注意:哈希值是我为演示编定的,真实值取决于字符串内容,但同键必同号这一点与真实一致——你可以在本机用 ("研发".hashCode() & Integer.MAX_VALUE) % 2 复算任何键的真实去向。
哈希分区的前提是"键的分布大体均匀"。现实数据常有长尾,倾斜就来了:
| 倾斜来源 | 典型场景 | 症状 |
|---|---|---|
| 键频长尾 | 99% 日志的 level 是 INFO | 一个 Reduce 处理亿条,其余闲置 |
| 键空间塌缩 | 全表按城市分区,但 NULL 城市占六成 | NULL 全挤进同一个分区 |
| 哈希聚集 | 复合键 hashCode 写得差,同尾号键扎堆 | 若干分区热、若干分区冷 |
倾斜的 Observable 证据在作业计数器里:Map output records 各任务基本相等,而 Reduce shuffle bytes 某一个任务吃掉总量的大头;最终所有 Reduce 同时领到任务,却有一个拖到最后几分钟才结束——作业总时长等于最慢 Reduce 的时长,其他任务的白白等待全是浪费。
一个可复算的推演:一亿条记录、10 个 Reduce,理想情况每区一千万条、每个 Reduce 跑 10 分钟;若某键占 60%,它所在分区六千万条、该 Reduce 跑 60 分钟,其余九个各 4 分钟。总时长从 10 分钟劣化到 60 分钟,集群利用率不到两成。这个算式是后面所有倾斜对策的收益标尺。
当业务知道比哈希更聪明的分法时,就自己写 Partitioner。经典题:把日志按级别分仓,ERROR 单独一个 Reduce 保它尽快出结果(数据量小),其余级别按哈希散开。
public class LogLevelPartitioner extends Partitioner<Text, Text> { @Override public int getPartition(Text key, Text value, int numPartitions) { String level = key.toString(); if (level.equals("ERROR")) return 0; // 0 号仓专属 ERROR // 其余级别挤进剩下的仓 且同键仍必同仓 return 1 + (level.hashCode() & Integer.MAX_VALUE) % (numPartitions - 1); } }
// 注册 只需一行 job.setPartitionerClass(LogLevelPartitioner.class); job.setNumReduceTasks(5); // 0 号收 ERROR 1 到 4 号散装其余
自定义分区器只有一条红线:同键必须永远返回同分区号(否则同键数据分裂,reduce 结果错误),而不同键是否同区随便。像"按 key 长度分区"这种写法满足红线但倾斜更凶——分区数、键分布、业务语义三者要一起想。
再进一步:若目标是"输出文件全局有序"(比如按时间排序的分片结果),哈希分区帮不上忙,用框架自带的 TotalOrderPartitioner——它先对输入做抽样,根据样本为每个 Reduce 划定键区间,使分区 0 的键全部小于分区 1,依此类推:
job.setPartitionerClass(TotalOrderPartitioner.class); job.setNumReduceTasks(4); // 先用 InputSampler 抽样生成分区边界文件 TotalOrderPartitioner.setPartitionFile(conf, new Path("/boundary/bins")); InputSampler.Sampler<Text, Text> sampler = new InputSampler.RandomSampler<>(0.01, 10000); // 1% 最多一万条 InputSampler.writePartitionFile(job, sampler);
代价是多一趟抽样作业与边界文件的分发;收益是 4 个输出文件"段间有序",拼接即可全局有序,而不必把全部数据压进单个 Reduce。
分区数恒等于 Reduce 任务数,所以 setNumReduceTasks 同时决定了并行度与每区数据量。拧小拧大各有利弊:
| Reduce 数 | 好处 | 代价 |
|---|---|---|
| 0 | 无 Shuffle,Map 输出直接写结果(纯 ETL 场景) | 没有任何归并 |
| 小(1~2) | 输出文件少,下游好读 | 并行度低,长尾风险大 |
| 适中 | 各任务负载均衡,磁盘网络均摊 | 输出文件数增多 |
| 过大 | 单任务更快 | 大量小文件、调度开销、拷贝连接数暴涨 |
经验起点:让每个 Reduce 处理 1 到 5 GB 压缩后的数据,或按"Reduce 数 ≈ Map 任务数 × 0.95 或 1.75"的旧经验估算。设为 0 是个合法且常用的特例——作业没有 reduce 阶段,Map 输出即最终输出,适合过滤、清洗、格式转换类 ETL;此时 map 的输出类型即最终输出类型,且通常配合多输出 MultipleOutputs 把不同类数据写不同目录(6.3 节会用到)。
倾斜的另一半答案(加盐与两阶段聚合)属于实战手法,留到 6.3 节清单里收束。下一节进入 Reduce 函数本体与最终输出的落盘。