4.2 Partition 分区


文档摘要

4.2 Partition 分区 本节摘要:分区是 Shuffle 的分拣口——每条中间键值对在 Map 端就被盖上"归属哪个 Reduce 任务"的章。默认的 HashPartitioner 用键的哈希对 Reduce 数取模,简单但可能倾斜。本节推演哈希分区的数学、演示数据倾斜的成因与症状、给出按业务键自定义分区的完整代码(含抽样 TotalOrderPartitioner),并讨论 Reduce 数量这个旋钮该怎么拧。 分区在流水线上的位置 回看 4.1 节的六步:分区发生在环形缓冲区溢写之前——每条 map 输出在进入缓冲区时就带上了分区号,溢写排序的规则是"分区号在前、键在后",最终 Map 端输出文件按分区切成段。Reduce 任务 i 只从所有 Map 输出里拷贝第 i 段。

4.2 Partition 分区

本节摘要:分区是 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 数量:一个牵一发动全身的旋钮

分区数恒等于 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 节会用到)。

本节自查

  1. 手算:"orders".hashCode() 在 Java 里等于 106750090,设 4 个 Reduce,它落在哪个分区(提示:106750090 % 4);
  2. 解释为什么 HashPartitioner 要先按位与 Integer.MAX_VALUE 而不是直接 Math.abs;
  3. 上述 LogLevelPartitioner 若把 numPartitions 设为 1 会发生什么,抛异常还是退化;
  4. 设计倾斜对策:某表 60% 行的 city 字段为 NULL,设计一个分区器把 NULL 行加盐打散到多个区,同时保证"非 NULL 同城同行"仍同区——想一想 NULL 行被拆开后 reduce 端要怎么重新合拢。

倾斜的另一半答案(加盐与两阶段聚合)属于实战手法,留到 6.3 节清单里收束。下一节进入 Reduce 函数本体与最终输出的落盘。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U