Apache Hadoop — 在一个 Pod 里搭起并运维 HDFS 与 YARN
用 Python 编写的 MapReduce 汇总访问日志
目标
用 Hadoop Streaming 编写 Python 的 Mapper 和 Reducer,按状态码统计三天的访问日志。用计数器衡量 Combiner 把 shuffle 缩减了多少,并处理按日期选择 Reducer 的分区(KeyFieldBasedPartitioner)、统计损坏行的用户计数器,以及用故意崩溃的 Mapper 体验任务失败与重试。
为什么重要
Streaming 把 Map 和 Reduce 换成了标准输入输出的契约。Mapper 从标准输入接收输入行,向标准输出输出“键值”行,框架按键排序后送入 Reducer 的标准输入。Reducer 收到的不是按键分好的组,而是排好序的行流,所以必须自己找出键发生变化的地方。只要遵守这个契约,就可以用任何语言编写 MapReduce,并且在本地用 cat | mapper | sort | reducer 做完全一样的测试。
Combiner 是在 Map 端预先合并的 Reducer。要保证结果不变,运算必须满足结合律和交换律(求和、最大值可以,平均值不行),并且输入和输出的形态必须相同。挂得正确,shuffle 能缩小到几十分之一。
默认的分区是对整个键取哈希。如果键是“日期路径”,而你想在一个 Reducer 里看到某一天的所有路径,就必须改成只按键的第一列来分区。而在生产环境中,Streaming 作业悄悄出错最常见的原因是损坏的行——丢掉它们没关系,但不统计丢掉了多少行就不行。
步骤
- 编写 /root/hdp/streaming/mapper.py。把一行按空白拆分,如果列数不少于 10,第九列(从 0 数起的第 8 号)是三位数字,并且不以
#开头,就输出상태코드<TAB>1(占位符为状态码),否则丢弃。 - 编写 /root/hdp/streaming/reducer.py。接收按键排序的
키<TAB>수行,把相同键的数相加,输出키<TAB>합(占位符依次为键、数量与合计;数量不一定是 1)。 - 把
/data/logs/access-2026-03-01.log到03.log上传到 HDFS 的 /user/root/streaming/in/,以作业名hdp-status运行 Streaming 作业,结果写入 /user/root/streaming/status。 - 给同一个作业用
-combiner挂上 Reducer,以作业名hdp-status-comb、输出 /user/root/streaming/status_comb 运行,并把该作业的MAP_OUTPUT_RECORDS、COMBINE_OUTPUT_RECORDS、REDUCE_INPUT_RECORDS按{"map_output": …, "combine_output": …, "reduce_input": …}的格式写入 /root/hdp/streaming/shuffle.json。 - 让 /root/hdp/streaming/mapper2.py 输出
날짜(yyyy-MM-dd)<TAB>경로<TAB>1(占位符依次为日期与路径),把第 1 步的 reducer.py 原样作为 Reducer,以作业名hdp-daily、2 个 Reducer、键两列、只按第一列分区的方式运行,结果写入 /user/root/streaming/daily。 - 在第 1 步的 Mapper 中,为每个损坏的行加一行把用户计数器
Lab、BadLines加 1 的代码(向标准错误写入reporter:counter:Lab,BadLines,1),以作业名hdp-bad运行,结果写入 /user/root/streaming/bad。 - 编写遇到
/event/spring路径就以异常结束的 /root/hdp/streaming/crash_mapper.py,以作业名hdp-crash、Map 尝试 1 次(-D mapreduce.map.maxattempts=1)运行,输出到 /user/root/streaming/crash,让它失败,把该作业 ID 写入 /root/hdp/streaming/crash.txt,然后用同一个作业名,使用第 1 步的 Mapper 重新运行,输出到 /user/root/streaming/crash-ok,让它成功。 - 在 /root/hdp/streaming/report.md 中写出
## 표준 입출력 계약、## 컴바이너와 분할、## 실패와 카운터三个小节。第二节写入第 4 步的map_output和reduce_input。
参考
- Streaming 作业:
mapred streaming -D mapreduce.job.name=<이름> -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input <입력> -output <출력>(占位符依次为作业名称、输入目录与输出目录)。-D选项要放在其他选项之前。 - 本地测试:
head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | python3 reducer.py。 - 键两列、按第一列分区:通用选项
-D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1放在-files之前,命令选项-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2放在-mapper之后。顺序错了会因“Unrecognized option”而中止。 - 作业历史保存在
/tmp/hadoop-yarn/staging/history/done_intermediate/root/(作业名中的-在文件名里会写成%2D)。失败的作业也会保留。 - 常见错误:以为 Reducer 每个键只被调用一次(它是逐行被调用的流),把 Combiner 输出的形态做得与输入不同,以及没有删除输出目录就重新运行。
- 官方文档:Hadoop Streaming · MapReduce Tutorial · mapred-default.xml
Mapper——接收一行,输出键和值
编写 /root/hdp/streaming/mapper.py。对标准输入的每一行按空白拆分,如果列数不少于 10,第 9 列(从 0 数起为 8)是三位数字,并且不以 # 开头,就向标准输出输出 상태코드<TAB>1(占位符为状态码),否则什么都不输出。
请先用 head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | uniq -c 测试。评分器会把你的 Mapper 直接在混入了损坏行的新输入上运行。
Reducer——在排好序的流中找出键变化的地方
编写 /root/hdp/streaming/reducer.py。从标准输入接收按键排序的 키<TAB>수 行,把相同键的数相加,向标准输出输出 키<TAB>합(占位符依次为键、数量与合计)。数量不一定是 1。
Reducer 收不到按键分好的列表。只有排好序的行流入,所以在与上一行的键不同的那一刻,就输出上一个键的合计并重新开始。不要忘了输出最后一个键。写成数量不必是 1 的形式,它就可以直接当 Combiner 用。
提交到 YARN
用 lab-hadoop start yarn 启动 YARN,把 /data/logs/access-2026-03-01.log 到 03.log 上传到 HDFS 的 /user/root/streaming/in/,然后用 mapred streaming -D mapreduce.job.name=hdp-status -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input /user/root/streaming/in -output /user/root/streaming/status 运行。
用 -files 传递的脚本会被复制到每个任务的工作目录。所以 -mapper 里只写文件名,不带路径。评分器会检查输出是否与直接从原始日志统计出的值相同,以及历史记录中是否有成功的 hdp-status 作业。
用 Combiner 缩减 shuffle
在第 3 步的作业上加上 -combiner 'python3 reducer.py',以作业名 hdp-status-comb、输出 /user/root/streaming/status_comb 运行。从该作业的历史中取出 TaskCounter 的 MAP_OUTPUT_RECORDS、COMBINE_OUTPUT_RECORDS、REDUCE_INPUT_RECORDS,按 {"map_output": 정수, "combine_output": 정수, "reduce_input": 정수} 的格式写入 /root/hdp/streaming/shuffle.json(占位符均为整数)。
求和满足结合律和交换律,所以把 Reducer 原样当作 Combiner,结果也相同。请看经过 Combiner 之后传入 shuffle 的行数(即 Reduce 输入)比 Map 输出少多少。
按日期选择 Reducer
让 /root/hdp/streaming/mapper2.py 对每个有效行输出 날짜(yyyy-MM-dd)<TAB>경로<TAB>1(占位符依次为日期与路径),把 reducer.py 作为 Reducer,以作业名 hdp-daily 运行,使用 -D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1(通用选项,放在前面)和 -partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2(命令选项,放在后面),结果写入 /user/root/streaming/daily。
键是前两列(日期、路径),所以排序按两列,但分区只按第一列(日期)。这样,某一天的所有路径都去往同一个 Reducer,一个输出文件完整地包含一天的数据。如果 reducer.py 是按最后一个制表符来拆分的,即使键有两列,也能照常工作。
用用户计数器统计损坏的行
让 mapper.py 对每个损坏的行向标准错误写入 reporter:counter:Lab,BadLines,1,以作业名 hdp-bad 运行,结果写入 /user/root/streaming/bad。该作业的计数器 Lab、BadLines 必须与三个日志中损坏的行数相等。
Streaming 会把标准错误中的 reporter:counter:<그룹>,<이름>,<증가량> 行当作计数器更新(占位符依次为计数器组、名称与增量)。这个计数器会留在作业历史中,作业结束之后,也能回答“丢掉了多少行”。
会崩溃的 Mapper——失败与重新运行
编写遇到 /event/spring 路径就以异常结束的 /root/hdp/streaming/crash_mapper.py,以作业名 hdp-crash、-D mapreduce.map.maxattempts=1 运行,输出到 /user/root/streaming/crash,让它失败,并把该作业 ID 写入 /root/hdp/streaming/crash.txt。然后用同一个作业名,使用第 1 步的 mapper.py 重新运行,输出到 /user/root/streaming/crash-ok,让它成功。
Map 任务失败时,框架会换一次尝试重新运行(默认 4 次)。把尝试次数缩减为 1 次,第一次失败就是作业失败。失败的作业也会留下历史,所以到底什么因为什么而崩溃,可以用 yarn logs -applicationId application_<같은 숫자> 读取任务的标准错误来确认(占位符为与作业 ID 相同的数字)。
记录契约、Combiner 和失败
在 /root/hdp/streaming/report.md 中写出 ## 표준 입출력 계약、## 컴바이너와 분할、## 실패와 카운터 三个小节。第二节以数字写入第 4 步的 map_output 和 reduce_input。
请写出:Reducer 收到的是什么;Combiner 把 shuffle 缩减到了几分之一;以及是如何发现损坏的行和失败的任务的。