6.2 Join 与 TopN 经典模式


文档摘要

6.2 Join 与 TopN 经典模式 本节摘要:业务题落到 MapReduce 上,多半是两副骨架的组合:把两张表按键接到一起(Join),或从海量记录里选出前 K 名(TopN)。本节给出 reduce 侧等值连接的完整可运行骨架,推演它的数据搬运账单与倾斜伏笔;再讲树形 TopN 的两级归并、复合键二级排序的配合,以及把两者拼成"分组 TopN"的合成题。 骨架一:reduce 侧等值连接 设定一个可复算的小场景:订单表 orders(oid, uid, amount)要连用户表 users(uid, name): 思路一句话:连接键 uid 作中间键,两表的行都发往同一键,reduce 在同一次回调里凑齐左右两半。

6.2 Join 与 TopN 经典模式

本节摘要:业务题落到 MapReduce 上,多半是两副骨架的组合:把两张表按键接到一起(Join),或从海量记录里选出前 K 名(TopN)。本节给出 reduce 侧等值连接的完整可运行骨架,推演它的数据搬运账单与倾斜伏笔;再讲树形 TopN 的两级归并、复合键二级排序的配合,以及把两者拼成"分组 TopN"的合成题。

骨架一:reduce 侧等值连接

设定一个可复算的小场景:订单表 orders(oid, uid, amount)要连用户表 users(uid, name):

orders users o1,u1,120 u1,赵敏 o2,u2,80 u2,钱磊 o3,u1,60 u3,孙蕙 o4,u3,200

思路一句话:连接键 uid 作中间键,两表的行都发往同一键,reduce 在同一次回调里凑齐左右两半。区分左右的方法是给值打标签——最常用的是拼接前缀:

public static class JoinMapper extends Mapper<LongWritable, Text, Text, Text> { private Text outKey = new Text(); private Text outVal = new Text(); private boolean isOrder; // 由路径判断身份 @Override protected void setup(Context context) { // 拿到本任务处理的输入路径 分片来自哪张表就标哪张 FileSplit split = (FileSplit) context.getInputSplit(); isOrder = split.getPath().getName().startsWith("orders"); } @Override protected void map(LongWritable k, Text line, Context context) throws IOException, InterruptedException { String[] c = line.toString().split(","); if (isOrder) { outKey.set(c[1]); // uid outVal.set("O" + c[0] + "," + c[2]); // O 前缀标订单 } else { outKey.set(c[0]); // uid outVal.set("N" + c[1]); // N 前缀标用户名 } context.write(outKey, outVal); } } public static class JoinReducer extends Reducer<Text, Text, Text, Text> { private Text out = new Text(); @Override protected void reduce(Text uid, Iterable<Text> values, Context context) throws IOException, InterruptedException { String name = null; List<String> orders = new ArrayList<>(); for (Text v : values) { String s = v.toString(); if (s.charAt(0) == 'N') { name = s.substring(1); // 维度侧先记下 } else { orders.add(s.substring(1)); // 事实侧攒列表 } } if (name == null) return; // 无主订单 丢弃或另写 for (String o : orders) { out.set(o + "," + name); context.write(uid, out); // 展开成宽行 } } }

提交与输出(Reduce 设 1 以便核对):

hadoop fs -rm -r /output/join hadoop jar join.jar JoinSide /input/join /output/join hadoop fs -cat /output/join/part-r-00000
u1 o1,120,赵敏 u1 o3,60,赵敏 u2 o2,80,钱磊 u3 o4,200,孙蕙

注意 reduce 里的顺序坑:不能假设 N 标签先到。Shuffle 只保证按键有序,键内的值顺序不定;若值序恰为订单在前,收到第一个订单就急着输出会拿到 null 名。所以必须两阶段:先扫一遍分类,再展开输出。这与 4.3 节"先拷值再处理"同宗——reduce 回调里凡是依赖"看到全部值"的决策,都要攒完再动。

账单与倾斜伏笔:reduce 侧 Join 把两张表的全量行都推进 Shuffle——users 表会被复制到每个有订单的键上重复参与。本例 9 行输入无所谓;生产上 10 亿订单连 1 亿用户,网络搬运以 TB 计。骨架还有个内置的不公平:热门键(头部用户的十万订单)全压进一个回调,倾斜在这副骨架里不是意外而是常态。

骨架一变体:map 侧连接

若一张表小到能塞进任务内存(比如千万行的维表约几 GB),用 DistributedCache 把它广播到每个任务节点,map 里直接查内存拼结果——零 Shuffle

@Override protected void setup(Context context) throws IOException { users = new HashMap<>(); try (BufferedReader r = new BufferedReader( new FileReader("users.txt"))) { // 缓存文件本地可见 String line; while ((line = r.readLine()) != null) { String[] c = line.split(","); users.put(c[0], c[1]); } } } @Override protected void map(LongWritable k, Text line, Context context) { String[] c = line.toString().split(","); String name = users.get(c[1]); if (name == null) return; context.write(new Text(c[1]), new Text(c[0] + "," + c[2] + "," + name)); }

取舍一目了然:reduce 侧 Join 慢而通用(两表都大);map 侧 Join 快而受限(一侧必须内存装下,且维表要在作业启动前静止)。中间地带是半连接:先扫维表抽键、再过滤事实表、最后 reduce 连接——三次作业换一次瘦身 Shuffle,思路本质仍是骨架一的变奏。

骨架二:树形 TopN

"全站销售额前 10 的订单"——最朴素的写法是 Reduce 数设 1,reduce 里维护一个小顶堆,全部数据过一遍取前 K。正确但完全放弃并行。生产写法是两级归并:每个 Reduce 先算本地前 K,再把本地冠军送去总决赛。

public static class TopNReducer extends Reducer<Text, DoubleWritable, Text, DoubleWritable> { private int n; private PriorityQueue<Order> heap; // 小顶堆 堆顶是当前第 N 名 static class Order implements Comparable<Order> { String id; double amt; Order(String id, double amt) { this.id = id; this.amt = amt; } public int compareTo(Order o) { return Double.compare(amt, o.amt); } } @Override protected void setup(Context context) { n = context.getConfiguration().getInt("top.n", 10); heap = new PriorityQueue<>(n); } @Override protected void reduce(Text uid, Iterable<DoubleWritable> values, Context context) { for (DoubleWritable v : values) { if (heap.size() < n) { heap.offer(new Order(uid.toString(), v.get())); } else if (v.get() > heap.peek().amt) { heap.poll(); // 挤掉当前最弱 heap.offer(new Order(uid.toString(), v.get())); } } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { List<Order> all = new ArrayList<>(heap); all.sort((a, b) -> Double.compare(b.amt, a.amt)); // 降序输出 for (Order o : all) { context.write(new Text(o.id), new DoubleWritable(o.amt)); } } }

这个 reducer 同时是两级归并的两级:第一阶段直接当 reducer 用,每个分区出本地 TopN;第二阶段把若干 part 文件作为输入再跑一遍同一个类(N 不变),本地前 K 的前 K 就是全局前 K——因为小顶堆的正确性只依赖"每个局部的前 K 必然包含全局前 K 的候选"这一条数学事实。假设 20 个 Reduce 各自漏掉某条记录,只可能是该记录在本地进不了前 10,而它在本地分区之外还有 19 个分区的对手……推演补全:若某记录属全局前 10 却在某分区落榜,则该分区至少有 10 条记录比它大,与"全局前 10"矛盾,故不可能——两级之后必然收敛

骨架二的进阶:二级排序出场

要"每用户消费最高的前 3 单"(分组 TopN),裸堆就不够了,得请出 4.4 节的二级排序:复合键(uid, amount)排序,分组比较器只按 uid 判等——于是每次 reduce 回调拿到的值列表天然按金额降序,取前三个即答:

// 复合键 3.2 节的 YearTempKey 换字段版 public class UidAmtKey implements WritableComparable<UidAmtKey> { private String uid; private double amt; public void set(String uid, double amt) { this.uid = uid; this.amt = amt; } public void write(DataOutput out) throws IOException { out.writeUTF(uid); out.writeDouble(amt); } public void readFields(DataInput in) throws IOException { uid = in.readUTF(); amt = in.readDouble(); } public int compareTo(UidAmtKey o) { int c = uid.compareTo(o.uid); if (c != 0) return c; return Double.compare(o.amt, amt); // 组内金额降序 } public String getUid() { return uid; } } // 分组比较器 只看 uid public static class UidGrouping extends WritableComparator { public UidGrouping() { super(UidAmtKey.class, true); } @Override public int compare(WritableComparable a, WritableComparable b) { return ((UidAmtKey) a).getUid().compareTo(((UidAmtKey) b).getUid()); } }
public static class GroupTopNReducer extends Reducer<UidAmtKey, Text, Text, Text> { @Override protected void reduce(UidAmtKey key, Iterable<Text> values, Context context) throws IOException, InterruptedException { int rank = 0; String uid = key.getUid(); for (Text v : values) { // 值列表已按金额降序 if (++rank > 3) break; context.write(new Text(uid), v); } } } // 注册 job.setGroupingComparatorClass(UidGrouping.class);

三副零件的配合在此看得清楚:复合键排序决定组内顺序,分组比较器决定组的大小,reduce 只负责消费顺序。不用二级排序也能解(reduce 里攒列表自己排序),但攒列表意味着热组整组进内存——二级排序把这个负担转交给了 Shuffle 的归并排序,4.4 节的投入在这里回本。

图 6.2-1 两副骨架与零件配合

图 6.2-1 两副骨架与零件配合

合成题:先连后选

真实需求常是复合的:"每个用户在 electronics 品类的消费前 3 单"。翻译成骨架:先 reduce 侧 Join 订单与品类维表得到宽表(或维表小时直接 map 侧连),宽表进第二个作业做分组 TopN——复合键(uid, amount)、分组按 uid、回调取前三。两个作业串联,各自是已解过的骨架。这条"把业务题拆成骨架序列"的翻译能力,正是本章比任何单个机制都更想交付的东西。

串联时的工程细节:中间宽表用 SequenceFileOutputFormat 落盘(4.3 节),第二个作业省掉文本解析且类型直传;宽表若巨大,第一个作业的倾斜对策(加盐)提前在 Join 阶段做掉,别把热点带进第二个作业。

本节自查

  1. 把 users 表的 u3 行删掉,推演 Join 输出与"无主订单"的两种处置策略各得到什么;
  2. 解释为什么 Join 的 reduce 回调不能边收边输出,以及这与 4.3 节哪条纪律同源;
  3. 证明树形两级归并的正确性:假设全局第 3 名在第一级某分区落榜,推出矛盾;
  4. 把分组 TopN 改成"每组第一名",二级排序方案里哪一行改动最小;
  5. 设计"全站 TopN 之上再要每组 TopN"的双作业序列,画出数据流草图。

下一节收拢全册最后一件事:跑慢了怎么诊断、怎么拧。


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