3.2 Map阶段与API编程实践


文档摘要

3.2 Map 阶段与 API 编程实践 本节摘要:本节从零写一个可提交的词频统计作业,覆盖作业驱动的六个配置步骤、输入格式与 RecordReader 如何把块变成记录、Writable 序列化体系,以及三类典型实战模式(清洗过滤、Reduce 端连接、二次排序)的实现要点。 第一个完整作业:六步装配 一个 MapReduce 作业 = 两个函数 + 一份装配说明书。装配用 Job 对象完成,固定六步: 提交与观察: 日志里 "number of splits:3" 值得多看一眼——它就是 3.1 节说的"Map 数由数据量决定"在运行时的显影。 输入格式:块如何变成记录 Map 拿到的"一条记录"从哪来?

3.2 Map 阶段与 API 编程实践

本节摘要:本节从零写一个可提交的词频统计作业,覆盖作业驱动的六个配置步骤、输入格式与 RecordReader 如何把块变成记录、Writable 序列化体系,以及三类典型实战模式(清洗过滤、Reduce 端连接、二次排序)的实现要点。

第一个完整作业:六步装配

一个 MapReduce 作业 = 两个函数 + 一份装配说明书。装配用 Job 对象完成,固定六步:

public class WordCount { 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 key, Text value, Context ctx) throws IOException, InterruptedException { // key 是行偏移 本例不用 value 是一行文本 for (String w : value.toString().split("\\s+")) { word.set(w); ctx.write(word, one); // 中间键值对交给框架 } } } public static class SumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override protected void reduce(Text word, Iterable<IntWritable> vals, Context ctx) throws IOException, InterruptedException { int sum = 0; for (IntWritable v : vals) sum += v.get(); result.set(sum); ctx.write(word, result); } } public static void main(String[] args) throws Exception { Job job = Job.getInstance(new Configuration(), "wordcount"); job.setJarByClass(WordCount.class); // 1 定位jar job.setMapperClass(TokenizerMapper.class); // 2 装配Map job.setReducerClass(SumReducer.class); // 3 装配Reduce job.setMapOutputKeyClass(Text.class); // 4 中间类型 job.setMapOutputValueClass(IntWritable.class); job.setOutputKeyClass(Text.class); // 5 输出类型 job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); // 6 路径 FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }

提交与观察:

export HADOOP_CLASSPATH=$HADOOP_CLASSPATH:. hadoop WordCount /data/docs /out/wordcount # 典型日志节选: # INFO mapreduce.JobResourceUploader: Uploading job.jar ... # INFO mapreduce.JobSubmitter: number of splits:3 ← 分片数即Map数 # INFO mapreduce.Job: Running job job_1724000000000_0002 # INFO mapreduce.Job: map 100% reduce 100% ← 进度两段式 hdfs dfs -cat /out/wordcount/part-r-00000 | head # hadoop 1523 # the 2044

日志里 "number of splits:3" 值得多看一眼——它就是 3.1 节说的"Map 数由数据量决定"在运行时的显影。

输入格式:块如何变成记录

Map 拿到的"一条记录"从哪来?答案是 InputFormat,它干两件事:把输入切成逻辑分片(InputSplit)、为每个分片造一个 RecordReader 把字节流转成键值对。

InputFormat 切分方式 记录形态 典型用途
TextInputFormat 按块切,行记录跨块自动接续 行偏移 + 行文本 通用文本日志
KeyValueTextInputFormat 同上 首个分隔符前为 key 已是 k-v 的文本
SequenceFileInputFormat 按块切 二进制键值对 中间结果、打包小文件
NLineInputFormat 每 N 行一分片 行偏移 + 行文本 控制Map数均匀

TextInputFormat 里有个容易被忽视的精巧设计:行跨块不截断。一条 500 字节的记录恰好横跨两个块,第一个分片的 RecordReader 会向下一个块"借"到行尾。这保证了"一条记录只被一个 Map 处理一次",也再次印证 1.2 节"块是物理切分、分片是逻辑切分"的分层。

另一个工程高频问题是控制分片大小。默认一个分片等于一个块(128MB),当输入是 1000 个 1MB 小文件时会产生 1000 个 Map 任务——调度开销淹没有效计算。两种对策:用 CombineFileInputFormat 把多个小文件打包进一个分片;或摄入侧就把小文件写进 SequenceFile(呼应 2.4 节小文件治理)。

Writable:可序列化的类型体系

MapReduce 进程间传的一切都必须可序列化。框架自带一套紧凑的 Writable 类型:IntWritable、LongWritable、Text(UTF-8 字符串)、DoubleWritable、NullWritable,以及容器 MapWritable、ArrayWritable。对比 Java 自带序列化,Writable 格式小数倍、速度快数倍,代价是失去语言中立——这也是为什么生态后来的重量级组件改用 Thrift/Protobuf 做跨语言协议。

自定义复杂数据只需实现 Writable 接口:

public class AccessLog implements Writable { private Text user; private LongWritable bytes; @Override public void write(DataOutput out) throws IOException { user.write(out); // 逐字段委托 bytes.write(out); } @Override public void readFields(DataInput in) throws IOException { user = new Text(); // 反序列化前先重置 user.readFields(in); bytes = new LongWritable(); bytes.readFields(in); } }

要作为 key 参与排序,还需实现 WritableComparable 并给出 compareTo——compareTo 的字典序将直接决定 Shuffle 排序与分组顺序,这是下一个模式的基础。

实战模式一:Map-only 清洗过滤

ETL 里最常见的作业其实没有 Reduce。需求:从原始日志中过滤掉爬虫流量、规范化时间戳:

public static class CleanMapper extends Mapper<LongWritable, Text, Text, NullWritable> { @Override protected void map(LongWritable k, Text line, Context ctx) throws IOException { String[] f = line.toString().split("\t"); if (f.length < 5) return; // 脏行丢弃 if (f[3].contains("spider")) return; // 爬虫过滤 String normTs = f[0].replace('/', '-'); // 时间规范化 ctx.write(new Text(String.join("\t", normTs, f[1], f[2], f[4])), NullWritable.get()); } } // 提交时不设置Reducer: job.setNumReduceTasks(0); // 无Shuffle直落盘 吞吐高一个档位

注意计数器:被丢弃的行数应该写进自定义计数器(3.4 节)而不是日志,下游才可核对"输入行数 − 输出行数 = 爬虫+脏行"的数据账。

实战模式二:Reduce 端连接

两个数据集按 key 对齐是数仓永恒需求。订单表 join 用户维表:

// 两个输入路径 一个MultipleInputs分别挂Mapper MultipleInputs.addInputPath(job, orderPath, TextInputFormat.class, OrderMapper.class); MultipleInputs.addInputPath(job, userPath, TextInputFormat.class, UserMapper.class); // OrderMapper 输出: userId → tag:order,amount ctx.write(new Text(userId), new Text("O" + amount)); // UserMapper 输出: userId → tag:userName,city ctx.write(new Text(userId), new Text("U" + name + "," + city)); // Reduce 端:先扫一遍分离维表与事实 protected void reduce(Text userId, Iterable<Text> vals, Context ctx) { String userInfo = null; List<String> orders = new ArrayList<>(); for (Text v : vals) { String s = v.toString(); if (s.charAt(0) == 'U') userInfo = s.substring(1); else orders.add(s.substring(1)); } if (userInfo == null) return; // 无维表匹配 事实行丢弃 for (String o : orders) ctx.write(userId, new Text(userInfo + "\t" + o)); }

要点两个。其一,维表必须能完整放进 Reduce 内存——超限就得换 Map 端连接(维表经 DistributedCache 广播到每个 Map,3.4 节展开)。其二,Iterable 不能缓存为 List 后继续迭代(框架复用对象),因此事实行先入 List、维表行存标量是标准写法。

实战模式三:二次排序

需求:取每个用户最近 3 次访问。朴素做法在 Reduce 端对全量值排序,内存可能爆。正确解法是把"时间"编码进 key,让框架替你排:

// 复合key:userId + ts 实现WritableComparable // compareTo:先比userId 再比ts倒序 // 分区器:只按userId分区 —— 保证同用户同Reduce job.setPartitionerClass(UserIdPartitioner.class); // 分组器:只按userId相等分组 —— 保证一次reduce调用拿到该用户全部记录且按时间有序 job.setGroupingComparatorClass(UserIdGroupingComparator.class);

三个组件各司其职:分区决定"去哪个Reduce",分组决定"几次reduce调用",key排序决定"调用内的顺序"。二次排序把"组内排序"外包给 Shuffle 的归并阶段,是零成本利用框架的经典范例。

本节要点回顾

  • 六步装配一个作业:jar、Mapper、Reducer、中间类型、输出类型、路径;日志里 splits 数即 Map 数;
  • InputFormat 决定记录从哪来,行跨块自动接续保证记录恰好处理一次;小文件用 CombineFileInputFormat 打包分片;
  • Writable 体系紧凑高效,key 需 WritableComparable,compareTo 字典序即 Shuffle 序;
  • Map-only 省整个 Shuffle,清洗过滤类作业首选;
  • Reduce 端连接靠 tag 区分来源,维表须小于 Reduce 内存,否则换广播连接;
  • 二次排序三件套:分区器管去向、分组器管聚合边界、key 排序管组内顺序。

代码能跑了,但作业在集群里到底怎么运转?下一节打开引擎盖:作业运行时序与 Shuffle 六阶段。


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