TT Lab
开始
学习 学习路径 课程

Apache Hadoop — 在一个 Pod 里搭起并运维 HDFS 与 YARN

在 YARN 上运行词频统计和 TeraSort,并通过计数器解读

在 TT Lab 中继续学习

目标

启动 YARN,用 MapReduce 统计五本书中的单词。作业结束后,打开留在 HDFS 中的作业历史(.jhist)读取计数器,并通过 TeraSort 确认:改变 Reducer 数量时输出如何划分,以及 shuffle 如何保证排序。

为什么重要

MapReduce 如今很少再被新写,但它的模型——Map 输出(键,值),框架按键汇集、排序后传递(shuffle),Reduce 按键合并——已被 Spark、Flink 和 SQL 引擎的 shuffle 全部继承。shuffle 为什么昂贵,Combiner(Map 端预先合并)为什么能减少 shuffle,Reducer 数量为什么会变成输出文件数,在这里看得最直白。 计数器是把作业做了什么用数字记录下来的档案。Map 读了多少行、输出了多少行,Combiner 缩减到了多少行,Reduce 收到了多少个键。结果出现异常时,这是最先要看的地方,而且作业结束后,它依然保存在历史文件中。 这个 Pod 的 YARN 按 Pod 内存(2Gi)把容器设得很小(AM 384MB,Map 和 Reduce 各 256MB)。所以任务是一两个一组依次运行的——慢一些,但结果和计数器是一样的。

步骤

  1. 用 lab-hadoop start yarn 启动 YARN,并把 yarn node -list -all 的输出保存到 /root/hdp/mr/nodes.txt。
  2. 把 /data/books/book-*.txt 五本书上传到 HDFS 的 /user/root/mr/in/。
  3. 用示例 jar 中的 wordcount 统计 /user/root/mr/in,把结果写入 /user/root/mr/wc。
  4. 把第 3 步作业的 ID(job_…)写入 /root/hdp/mr/jobid.txt。
  5. 从该作业的历史文件(.jhist)中取出 TaskCounter 的 MAP_INPUT_RECORDS、MAP_OUTPUT_RECORDS、COMBINE_INPUT_RECORDS、COMBINE_OUTPUT_RECORDS、REDUCE_INPUT_GROUPS、REDUCE_OUTPUT_RECORDS,写入 /root/hdp/mr/counters.json。
  6. 用 3 个 Reducer(-D mapreduce.job.reduces=3)运行同样的 wordcount,把结果写入 /user/root/mr/wc3。
  7. 用 teragen 在 /user/root/mr/tera-in 中生成 10 万行,用 terasort 排序到 /user/root/mr/tera-out,再用 teravalidate 把验证结果写入 /user/root/mr/tera-val。
  8. 在 /root/hdp/mr/report.md 中写出 ## 맵과 리듀스、## 카운터、## 정렬 三个小节。第二节写入第 5 步的 MAP_OUTPUT_RECORDS 和 COMBINE_OUTPUT_RECORDS。

参考

启动 YARN

用 lab-hadoop start yarn 启动 ResourceManager 和 NodeManager,并把 yarn node -list -all 的输出保存到 /root/hdp/mr/nodes.txt。

NodeManager 必须向 ResourceManager 报告自己能提供的内存和核数,才能接收作业。节点必须显示为 RUNNING。请用 lab-hadoop status 确认四个 JVM(NameNode、DataNode、ResourceManager、NodeManager)都已启动。

上传输入

把 /data/books/book-1.txt 到 book-5.txt 上传到 HDFS 的 /user/root/mr/in/。

输入目录中的每个文件都会生成输入分片(split),每个分片对应一个 Map 任务。五个文件就是五个 Map。

运行 wordcount

用 yarn jar /opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.5.0.jar wordcount /user/root/mr/in /user/root/mr/wc 统计单词。

作业结束后,输出目录中会生成 _SUCCESS 和 part-r-00000。Map 为每个单词输出(单词,1),Combiner 先在 Map 端相加,Reducer 再按单词合并。评分器会把你的输出与直接从原文统计出的值对照。

找到作业历史

把第 3 步作业的 ID(job_<숫자>_<숫자>)单独写成一行,保存到 /root/hdp/mr/jobid.txt(占位符为数字)。

运行作业时,控制台会打印“Running job: job_…”。如果错过了,就在 HDFS 的 /tmp/hadoop-yarn/staging/history/done_intermediate/root/ 列表中,找名称里含有 word+count 的 .jhist——即使没有历史服务器,AM 也会把它留在这里。

读取计数器

从第 4 步作业的 .jhist 中,取出 org.apache.hadoop.mapreduce.TaskCounter 组的 MAP_INPUT_RECORDS、MAP_OUTPUT_RECORDS、COMBINE_INPUT_RECORDS、COMBINE_OUTPUT_RECORDS、REDUCE_INPUT_GROUPS、REDUCE_OUTPUT_RECORDS,按 {"이름": 정수} 的格式写入 /root/hdp/mr/counters.json(占位符依次为计数器名称和整数值)。

MAP_OUTPUT_RECORDS 是单词数(含重复),COMBINE_OUTPUT_RECORDS 是各个 Map 的不同单词数之和——也就是 shuffle 实际传递的行数。REDUCE_INPUT_GROUPS 必须等于全部不同单词的个数。这两个数字的比值,就是 Combiner 省下的 shuffle。

三个 Reducer,输出也是三个

用 yarn jar <jar> wordcount -D mapreduce.job.reduces=3 /user/root/mr/in /user/root/mr/wc3 运行带 3 个 Reducer 的作业。

Map 输出按键的哈希值被分配到三个 Reducer 之一(HashPartitioner)。所以一个单词只会出现在一个文件中,每个文件内部是有序的,但文件之间并不有序。

TeraSort——shuffle 造就了排序

用 teragen 100000 /user/root/mr/tera-in 生成 10 万行(每行 100 字节),用 terasort /user/root/mr/tera-in /user/root/mr/tera-out 排序,再用 teravalidate /user/root/mr/tera-out /user/root/mr/tera-val 验证(都使用同一个示例 jar)。

TeraSort 会对输入抽样,把键的范围按 Reducer 数量切开(全局有序分区),每个 Reducer 把自己范围内的数据排好序输出。这样,把输出文件按顺序首尾相接,就是整体排序。TeraValidate 发现错位就会写出“misorder”,没有的话只写校验和。

用数字记录 shuffle

在 /root/hdp/mr/report.md 中写出 ## 맵과 리듀스、## 카운터、## 정렬 三个小节。第二节以数字写入第 5 步的 MAP_OUTPUT_RECORDS 和 COMBINE_OUTPUT_RECORDS 的值。

第一节写 Map 和 Reduce 各有多少个、输出是如何划分的;第二节写 Combiner 把 shuffle 缩减到了几分之一;第三节写 TeraSort 是如何造出整体顺序的。