3.2 Map 阶段与 API 编程实践 本节摘要:本节从零写一个可提交的词频统计作业,覆盖作业驱动的六个配置步骤、输入格式与 RecordReader 如何把块变成记录、Writable 序列化体系,以及三类典型实战模式(清洗过滤、Reduce 端连接、二次排序)的实现要点。 第一个完整作业:六步装配 一个 MapReduce 作业 = 两个函数 + 一份装配说明书。装配用 Job 对象完成,固定六步: 提交与观察: 日志里 "number of splits:3" 值得多看一眼——它就是 3.1 节说的"Map 数由数据量决定"在运行时的显影。 输入格式:块如何变成记录 Map 拿到的"一条记录"从哪来?
本节摘要:本节从零写一个可提交的词频统计作业,覆盖作业驱动的六个配置步骤、输入格式与 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 节小文件治理)。
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 排序与分组顺序,这是下一个模式的基础。
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 节)而不是日志,下游才可核对"输入行数 − 输出行数 = 爬虫+脏行"的数据账。
两个数据集按 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 的归并阶段,是零成本利用框架的经典范例。
代码能跑了,但作业在集群里到底怎么运转?下一节打开引擎盖:作业运行时序与 Shuffle 六阶段。