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

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

在 Hadoop Streaming 中,标准输入输出的一行就是全部约定

在 TT Lab 中继续学习

一句话总结

Hadoop Streaming 是把 Mapper 和 Reducer 换成任意读取标准输入、写入标准输出的程序的工具。契约很简单——一行就是一条记录,第一个制表符之前是键,Reducer 收到的是按键排序的行。正因为简单,键发生变化的边界必须由脚本自己去找,失败和计数器也必须通过退出码和标准错误来汇报。

为什么需要它

MapReduce 原本的契约是 Java 接口:要继承 Mapper,对齐 Writable 类型,再打包成 jar。为了统计访问日志中每个状态码的请求数,这样的形式未免太重了。而且,处理数据的人大多早已有了在一台机器上运行的 Python 或 shell 脚本。

Hadoop Streaming 文档填补了这道缝隙。任何可执行文件或脚本,都可以用作 Mapper 和 Reducer,文档中的第一个示例就用 /bin/cat 作 Mapper,用 /usr/bin/wc 作 Reducer。在一台机器上用 cat log | map.py | sort | reduce.py 运行的管道,几乎原样就能扩展到几百台机器上。中间的 sort,相当于由 MapReduce 的 shuffle 来代替。

工作原理

根据文档,Map 任务在启动时,会把指定的可执行文件作为独立进程启动。它把输入分片转成行,送入该进程的标准输入,再把标准输出中产生的行收集起来,转换成键值对。默认规则是第一个制表符字符之前是键,之后是值;如果没有制表符,整行就是键,值为空。Reducer 一侧也一样。框架把排好序、合并好的键值对重新展开成 키\t값 这样的行(占位符依次为键与值),送入 Reducer 进程的标准输入。

这个契约中最重要的事实是:Reducer 收到的不是按键分好的组,而是排好序的行流。Java 的 Reducer 会对每个键调用一次 reduce(키, 값 목록)(占位符依次为键与值列表),而 Streaming 的 Reducer 要一行一行地读取,自己判断“键是否变了”。正如 MapReduce 教程所说,Map 输出经过排序后,再按 Reducer 分开送来,所以相同键的行一定是连在一起到达的。Reducer 就依靠这唯一的保证,记住上一行的键,在键发生变化的那一刻输出合计并重置。

import sys

current, total = None, 0
for line in sys.stdin:
    key, _, value = line.rstrip("\n").partition("\t")
    if key != current:
        if current is not None:
            print(f"{current}\t{total}")
        current, total = key, 0
    total += int(value)
if current is not None:
    print(f"{current}\t{total}")

在循环之后再输出一次最后一个键的那一行,漏写它是最常见的错误。这样一来,结果中按排序位于最末的那个键,就会悄无声息地消失。

Mapper 和 Reducer 脚本必须存在于工作节点上。用 -files 传递后,任务的工作目录中会生成同名的符号链接。文档警告说,通用选项必须放在 Streaming 选项之前,比如 -files 或 -D,否则命令会失败。

调节 shuffle 的旋钮

Combiner。可以用 -combiner 再指定一个可执行文件。如果 Mapper 是按状态码输出 1,就可以把与 Reducer 相同的求和脚本挂成 Combiner,大幅减少传入 shuffle 的行数。Combiner 会运行几次是没有保证的,所以只挂那些应用多次结果也相同的求和类运算。

键字段与分区。当键像 2026-09-19.404 这样由多个字段构成时,有时希望排序按完整的键,而分区只按前面的字段。文档中的 KeyFieldBasedPartitioner 示例就是这样。用 stream.num.map.output.key.fields 指定键有几个字段,用 map.output.key.field.separator 指定字段分隔符,再用 mapreduce.partition.keypartitioner.options=-k1,1 这样的形式指定用于分区的字段。这样,同一天的行都会去往同一个 Reducer,并且在其内部按完整的键顺序排好序送达。想要按日期分文件的结果,或者想在一个 Reducer 里按日期做小计时,就用它。

Reducer 数量与只有 Map 的作业。默认的 Reducer 数量是 mapred-default.xml 中 mapreduce.job.reduces 的默认值 1。根据文档,设为 0 时没有 Reducer,Map 的输出就是结果。如果只是过滤行或转换格式,干脆去掉 shuffle 最快。

如何汇报失败与计数器

Streaming 脚本没有 Java API,只能通过两条通道与框架对话。一条是退出码。根据文档,默认情况下以非零退出码结束的任务会被当作失败。失败的任务会被重试,mapreduce.map.maxattempts 的默认值是 4,所以同一个 Map 失败四次,整个作业就失败。遇到一个损坏的行就抛出异常的 Mapper,会在包含这一行的分片上死四次,并把整个作业拖垮。

另一条是标准错误。把 reporter:counter:<그룹>,<카운터>,<양> 格式的行写入 stderr(占位符依次为计数器组、计数器名称与增量),计数器就会增加;reporter:status:<메시지> 则会改变状态文字(占位符为状态消息)。所以遇到损坏的行,最好不要抛异常崩溃,而是跳过它,同时写出 reporter:counter:logs,malformed,1。作业会一直运行到结束,有多少行损坏,也会以数字留在计数器里。如果混写到标准输出里,这一行就会变成结果记录,所以一定要发送到 stderr。

配置值也可以作为环境变量读取。根据文档,配置名称中的点会变成下划线,例如 mapreduce.job.id 会显示为 mapreduce_job_id。

在现场相遇的样子

第一,本地正确,到了集群上最后一个键却丢了。这是 Reducer 没有在循环之后输出最后一组。在一台机器上用 sort 测试时,用小数据有时碰巧发现不了。

第二,作业重试四次后死掉。日志中有一行编码损坏,或者有列数不足的行。如果 Python Mapper 在那一行抛出异常,同一分片的重试也会死在同一行。重试是为临时性故障设计的,不是为坏数据设计的。

第三,计数器出现在结果文件里。把 reporter:counter 用 print 写到了标准输出。结果里就出现了以 reporter:counter:... 开头的键。

第四,提示脚本在节点上不存在。漏了 -files,或者把它放在了 Streaming 选项之后。同时也要确认执行权限和第一行的解释器声明。

实际工作中真正重要的事

下一项实验要做什么

编写一个 Python Mapper(对访问日志中的每个状态码输出 1),以及一个 Reducer(连续读取排好序的输入并求和),用 Streaming 作业运行。把同一个求和脚本挂成 Combiner,用计数器对比 Map 输出与 Reduce 输入的记录数;再把日期和路径用制表符连成的两列键,通过 KeyFieldBasedPartitioner 只按第一列(日期)分给两个 Reducer。用标准错误增加用户计数器来统计损坏的行,再让一个在特定路径上故意崩溃的 Mapper 以 Map 尝试 1 次来运行,看作业立刻失败,然后用正常的 Mapper 把同名作业重新运行并使其成功。