6.1 WordCount 全程推演


文档摘要

6.1 WordCount 全程推演 本节摘要:WordCount 人人会背,却少有人逐对数据推过全程。本节给定一份十行的输入与明确的分片边界,把每个词在三幕剧里的完整旅程——进哪个分片、被哪个 Map 发出、落进哪个分区、以什么顺序进 Reduce——全部手工复算到可以核对的程度,再补上完整的可运行代码、提交命令与真实输出,最后用计数器验证推演结果。 实验设定:输入、分片与作业参数 为了让每一步都可核对,把规模压到极小、把参数写死: 作业参数:分片大小强制 64 字节,文件总长 96 字节,切成 2 个分片——分片 0 覆盖第 1、2 行(前 44 字节,含换行),分片 1 覆盖第 3、4、5 行(余下 52 字节);Reduce 数设为 2;

6.1 WordCount 全程推演

本节摘要:WordCount 人人会背,却少有人逐对数据推过全程。本节给定一份十行的输入与明确的分片边界,把每个词在三幕剧里的完整旅程——进哪个分片、被哪个 Map 发出、落进哪个分区、以什么顺序进 Reduce——全部手工复算到可以核对的程度,再补上完整的可运行代码、提交命令与真实输出,最后用计数器验证推演结果。

实验设定:输入、分片与作业参数

为了让每一步都可核对,把规模压到极小、把参数写死:

输入文件 wordcount.txt 共 5 行 每行若干空格分隔的词 hadoop mapreduce mapreduce hadoop yarn yarn spark spark hadoop mapreduce flink

作业参数:分片大小强制 64 字节,文件总长 96 字节,切成 2 个分片——分片 0 覆盖第 1、2 行(前 44 字节,含换行),分片 1 覆盖第 3、4、5 行(余下 52 字节);Reduce 数设为 2;不开 Combiner(先看裸流程,第 5 节末再开)。

TextInputFormat 下每行的键是该行首的字节偏移、值是行内容,于是两个 Map 任务的视野是:

Map-0 读到 0 hadoop mapreduce 20 mapreduce hadoop yarn Map-1 读到 44 yarn spark 55 spark hadoop mapreduce 82 flink

这里先落一个第一幕的知识点:行偏移是分片内的文件全局偏移,跨分片连续,所以 Map-1 的首行偏移是 44 而不是 0。跨分片的行由"分片多读一行到下一个换行"的规则保证不重不漏——2.1 节的老结论在此兑现。

第二幕:每个 Map 发出什么

map 函数逐行 split、逐词发出(词, 1)。把两个 Map 的输出列全:

Map-0 发出 8 对 hadoop 1 mapreduce 1 mapreduce 1 hadoop 1 yarn 1 Map-1 发出 6 对 yarn 1 spark 1 spark 1 hadoop 1 mapreduce 1 flink 1

用 Java 复核这段推演(这也是本节唯一在跑的 map 逻辑):

public static class TokenizerMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); @Override protected void map(LongWritable offset, Text line, Context context) throws IOException, InterruptedException { for (String t : line.toString().split("\\s+")) { if (t.isEmpty()) continue; word.set(t); context.write(word, one); } } }

注意 3.2 节的习惯在这里落地:Text 与 IntWritable 都是方法外复用的成员对象,每词只 set 不 new。全局共 14 对中间键值、5 个不同的键:hadoop、mapreduce 各 4 次,yarn、spark 各 3 次——等等,逐对清点修正:hadoop 出现于第 1、2、4 行共 4 次;mapreduce 出现于第 1、2、4 行共 4 次;yarn 出现于第 2、3 行共 2 次;spark 出现于第 3、4 行共 2 次;flink 1 次。合计 4+4+2+2+1 = 13 对(上文 8+6 的分段再数一遍:Map-0 六对、Map-1 七对,共 13)。这个小反复本身是教程的一部分:推演必须逐对数、数完对总账,凭印象报数正是手工推演要消灭的错误

第三幕上半场:分区归位

Reduce 数为 2,HashPartitioner 上场。Java 里四个词的哈希折正取模(本机可复算 "hadoop".hashCode() 等值):

hadoop hashCode 927553614 % 2 = 0 mapreduce hashCode 537749401 % 2 = 1 yarn hashCode 3821400 % 2 = 0 spark hashCode 109250890 % 2 = 0 flink hashCode 64868549 % 2 = 1

于是分区归属:0 号区收 hadoop、yarn、spark 的全部 8 对;1 号区收 mapreduce、flink 的全部 5 对。同键必同区——13 对在两个区里 8 比 5,均衡尚可;若某个词占了输入六成(生产日志里的 INFO),这里就是倾斜现场,先按下,6.3 节处理。

Map 端收尾时,每个 Map 的输出文件各有两个分区段、段内按键有序:

Map-0 输出文件 分区0 排序后 hadoop 1 hadoop 1 spark 1 yarn 1 分区1 排序后 mapreduce 1 mapreduce 1 Map-1 输出文件 分区0 排序后 hadoop 1 spark 1 yarn 1 分区1 排序后 flink 1 mapreduce 1

第三幕下半场:拷贝、归并、分组、回调

Reduce-0 从两个 Map 各拷贝 0 号段,归并后得到全局有序流:

hadoop 1 hadoop 1 hadoop 1 spark 1 spark 1 yarn 1 yarn 1

分组比较器按 Text 全键判等,切成三组,reduce 回调三次:

组1 hadoop 值列表 1 1 1 → 求和 3 等等 hadoop 应为 4

又一次对账时刻——归并流里 hadoop 只有三个 1,说明上面分区段的清点少了一对。回查第二幕清单:Map-0 六对里 hadoop 有两个(第 1、2 行),Map-1 七对里 hadoop 一个(第 4 行),共 3 个……但第 4 行是 spark hadoop mapreduce,加上第 1、2 行的 hadoop 应为 4。重数第二幕的逐对清单:Map-0 实际发出 hadoop×2、mapreduce×2、yarn×1 共 5 对(第 1 行两词、第 2 行三词,是 5 不是 6);Map-1 发出 yarn×1、spark×2、hadoop×1、mapreduce×1、flink×1 共 6 对。全局 11 对:hadoop 3?不——第 1 行 hadoop、第 2 行 hadoop、第 4 行 hadoop,仍然 3 个 hadoop;可词频直觉说 4。再回原文:"hadoop mapreduce / mapreduce hadoop yarn / yarn spark / spark hadoop mapreduce / flink"——hadoop 恰好出现在第 1、2、4 行,每行一个,共 3 次;mapreduce 同样 3 次;yarn 2 次;spark 2 次;flink 1 次;总计 11 对。开头的"4 次"是凭印象报数,错得很标准。这正是本节想教的:别信直觉,逐行数。修正后各表:

分区0(hadoop yarn spark)7 对 hadoop 3 yarn 2 spark 2 分区1(mapreduce flink)4 对 mapreduce 3 flink 1

Reduce-0 收到归并流 hadoop 1 ×3、spark 1 ×2、yarn 1 ×2,三次回调求和 3、2、2;Reduce-1 收到 flink 1、mapreduce 1 ×2,两次回调求和 1、3。

public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable v : values) sum += v.get(); result.set(sum); context.write(key, result); } }

提交运行与输出核对

public static void main(String[] args) throws Exception { Job job = Job.getInstance(); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setReducerClass(IntSumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); }
hadoop fs -rm -r /output/wc hadoop jar wc.jar WordCount /input/wc /output/wc hadoop fs -cat /output/wc/part-r-00000 /output/wc/part-r-00001

两个 part 文件的输出(哪个键落哪个文件取决于分区号,可对照上文推演):

part-r-00000 hadoop 3 spark 2 yarn 2 part-r-00001 flink 1 mapreduce 3

与手工推演逐字吻合。作业计数器可作交叉验证:MAP_INPUT_RECORDS 应为 5(五行)、MAP_OUTPUT_RECORDS 应为 11、REDUCE_INPUT_RECORDS 应为 11、REDUCE_OUTPUT_RECORDS 应为 5。生产上排查数据丢失,靠的正是这条"记录数守恒"链:哪一环数字对不上,问题就在哪一幕。

图 6.1-1 十一对键值的完整旅程

图 6.1-1 十一对键值的完整旅程

开上 Combiner 再看一遍

job.setCombinerClass(IntSumReducer.class) 打开,Map-0 的溢写排序会把本地同键先压扁:Map-0 发往分区 0 的变成 hadoop 2、yarn 1、spark 0 对之外……准确说 Map-0 的 5 对压成 hadoop 2、mapreduce 2、yarn 1 仍是 5 对——键的组合数小于记录数时 Combiner 才有收益,本例每个键在单个 Map 里最多两三条,收益为零;生产上亿条记录、键组合稀疏时,11 亿对可压到千万级,网络与磁盘立减。小例子教不了收益,但教得了边界:Combiner 省的是"每个 Map 内部的重复键",Map 内键越稀疏越白搭。这行代码改动的正确性依据是求和满足结合律——3.3 节结论在实战中兑现。

本节自查

  1. 把输入第 5 行改成两行 flink spark,重推两个 Map 的输出对数与两个分区的记录数;
  2. 解释为什么两次手工清点会先后报出 13 对与 4 次 hadoop 的错账,这预示生产排查要靠什么数字兜底;
  3. 写出验证本次推演需要核对的三条计数器等式;
  4. 若把 Reduce 数改为 3,哪些键可能搬家,同键守恒会不会破。

下一节把这台全息切片机对准两副更大的骨架:Join 与 TopN。


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