3.1 Map 阶段的执行机制


文档摘要

3.1 Map 阶段的执行机制 本节摘要:Map 任务拿到分片后走一条固定流水线:初始化 InputFormat 与 RecordReader,逐条读出键值对交给用户 map 函数,输出写进环形缓冲区,缓冲到阈值就排序、(可选压缩)溢写成磁盘分段,最后把所有分段按分区归并成待拉取的输出文件。理解这条流水线,就理解了 Shuffle 前半段的全部机制。 从任务开机说起 调度器把 Map 任务分配到某个节点后,第一步不是跑你的代码,而是初始化运行环境:启动任务 JVM、加载作业的配置与 jar、构造用户指定的 InputFormat,由它打开分片对应的文件(通常借助 HDFS 副本选择拿到本地或近端的流)。

3.1 Map 阶段的执行机制

本节摘要:Map 任务拿到分片后走一条固定流水线:初始化 InputFormat 与 RecordReader,逐条读出键值对交给用户 map 函数,输出写进环形缓冲区,缓冲到阈值就排序、(可选压缩)溢写成磁盘分段,最后把所有分段按分区归并成待拉取的输出文件。理解这条流水线,就理解了 Shuffle 前半段的全部机制。

从任务开机说起

调度器把 Map 任务分配到某个节点后,第一步不是跑你的代码,而是初始化运行环境:启动任务 JVM、加载作业的配置与 jar、构造用户指定的 InputFormat,由它打开分片对应的文件(通常借助 HDFS 副本选择拿到本地或近端的流)。接着 InputFormat 创建 RecordReader,从此你的 map 函数眼里就只剩下"一条条键值对",字节流、块边界、压缩解压,全被下面两层消化掉了。

这条分工链的责任边界值得记牢:

组件 职责 用户可否替换
InputFormat 校验输入、生成分片、创建读取器 可以,作业设置里指定
RecordReader 把分片字节流切成一个个键值对 可以,随自定义 InputFormat 提供
Mapper 与 map 对单条键值对做业务加工 核心,必写
输出收集器 收集 map 输出写入缓冲区 框架内部,不可替换

以最常用的 TextInputFormat 为例,它生成的键是行首在分片中的字节偏移(LongWritable),值是行内容(Text)。换二进制场景可用 SequenceFileInputFormat,键值类型随文件头记录;再特殊的结构(比如定长二进制记录、多行聚合记录)就自己实现一对 InputFormat 与 RecordReader。

// 自定义 RecordReader 的骨架 只列关键方法 public class FixedLenRecordReader extends RecordReader<LongWritable, BytesWritable> { private long pos; // 当前读取位置 private long end; // 分片结束边界 public boolean nextKeyValue() { if (pos >= end) return false; // 读固定长度的一条记录填入 value 更新 pos key.set(pos); return true; } }

环形缓冲区:第二幕的主舞台

map 函数每调一次 context.write,数据并没有直接写磁盘,而是进入一个环形缓冲区(默认约 100MB)。它是 Map 端性能设计的核心:

为什么是环形?因为要同时容纳两股写入——键值数据的序列化字节,与每条记录的元信息(键起点、值起点、分区号)——两股数据各自从数组两端向中间生长,空间利用率高于"头部分配"。缓冲区写到八成(默认阈值 0.8)触发一次溢写

  1. 锁定已写的八成区域,新写入继续往剩余空间走,互不阻塞;
  2. 在锁定区内,按"分区号在前、键排序在后"的顺序对这段数据排序——注意排序是在内存里对元信息排序再顺序写出,不是搬数据;
  3. 若配置了 Combiner,排序后相邻的同键记录先做一次本地归并;
  4. 若开了中间压缩(上一章的 snappy),此时压缩;
  5. 写成一个磁盘上的溢写分段文件,然后释放锁定区。

一个处理大量输出的 Map 任务会溢写多次,产生多个分段。任务收尾时,把所有分段按分区归并成一个文件(每个分区内部有序),这就是 Map 端的最终输出,静静等待 Reduce 端的 HTTP 请求来拉取。

图 3.1-1 环形缓冲区与溢写循环

图 3.1-1 环形缓冲区与溢写循环

三个关键参数

mapreduce.task.io.sort.mb 缓冲区大小 默认 100MB mapreduce.map.sort.spill.percent 溢写阈值 默认 0.8 mapreduce.task.io.sort.factor 归并时一次合并的流数 默认 10

调优直觉:任务输出量大、日志里溢写次数高(作业计数器有专门指标),优先加大 sort.mb(注意任务内存上限同步放宽);分段特别多时提高 sort.factor 让收尾归并一次并更多路。

一个容易踩的坑:setup 与 cleanup

Mapper 除了逐条调用的 map,还有任务开始前一次的 setup 和结束后一次的 cleanup。把外部资源的连接(数据库连接、词典加载)放进 setup、在 cleanup 里关闭,而不是在 map 里反复创建——一个处理千万条记录的任务,每条记录建一次连接的开销足以把作业拖死。同理,跨记录的累计状态(如缓存、布隆过滤器)也应挂在实例字段上,map 是"每条调用一次"的纯函数位,不是"从头执行到尾"的脚本。

public class DictMapper extends Mapper<LongWritable, Text, Text, Text> { private Map<String, String> dict; // 任务级缓存 protected void setup(Context context) { dict = loadDictFromHDFS("dict"); // 每个任务只加载一次 } protected void map(LongWritable k, Text line, Context ctx) { String enriched = enrich(line.toString(), dict); // 逐条使用 // ... 写出 } }

Map 端输出的形态

任务结束时,Map 端留下的输出文件内部按分区排列,每个分区内按键有序。这个"已排序"的性质是第三幕 Reduce 端能做高效归并的前提——Map 端多付一次排序,Reduce 端就可以做多路归并而不是全量重排,整条链路的总排序成本被摊薄。设计上的连续性,从这里已经埋下伏笔。

本节要点回顾

  • 流水线五环:InputFormat 初始化、RecordReader 切键值对、map 加工、环形缓冲、溢写归并;
  • 组件边界:InputFormat 管分片与读取器创建,RecordReader 管切记录,map 管业务,职责可替换点在前两者;
  • 环形缓冲区:默认百MB、八成阈值,锁定区内完成分区排序、可选 Combiner 与压缩后落盘;
  • 溢写与归并:多次溢写成多段,任务收尾按分区归并成单一输出,段内分区内部有序;
  • 参数三件套:sort.mb、spill.percent、sort.factor,对应缓冲大小、触发时机、归并路数;
  • setup 与 cleanup:重资源初始化的正确位置,避免每条记录重建开销。

输出已在 Map 端排好序、分好区,下一节先补齐支撑这条流水线的类型系统,再进入第三幕。


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