MapReduce 进阶:倒排索引与用户销售额统计
深入 MapReduce 编程,实现倒排索引(Combiner 优化)和用户销售额统计两个实战案例,理解分布式数据处理模式。
MapReduce 进阶:倒排索引与销售额统计
倒排索引是搜索引擎的基石:从词到文档的映射。通过 MapReduce 构建倒排索引是理解"分而治之"大数据处理范式的绝佳例子。
倒排索引原理
正向索引:文档 → 词(doc1 contains "hello")
倒排索引:词 → 文档列表("hello" → [doc1, doc3, doc5])
搜索引擎的工作方式就是对用户输入的关键词,在倒排索引中快速找到包含该关键词的所有文档。
MapReduce 实现
public class InvertedIndex {
// Map:输出 (word, docId)
public static class TokenizerMapper
extends Mapper<Object, Text, Text, Text> {
private Text word = new Text();
private Text docId = new Text();
public void map(Object key, Text value, Context context)
throws IOException, InterruptedException {
// 假设输入格式:docId content
String[] parts = value.toString().split(" ", 2);
docId.set(parts[0]);
StringTokenizer itr = new StringTokenizer(parts[1]);
while (itr.hasMoreTokens()) {
word.set(itr.nextToken().toLowerCase());
context.write(word, docId);
}
}
}
// Combine:同文档内去重(优化网络传输)
public static class Combiner
extends Reducer<Text, Text, Text, Text> {
public void reduce(Text key, Iterable<Text> values, Context context)
throws IOException, InterruptedException {
Set<String> uniqueDocs = new HashSet<>();
for (Text val : values) {
uniqueDocs.add(val.toString());
}
context.write(key, new Text(String.join(",", uniqueDocs)));
}
}
// Reduce:聚合输出 (word, "doc1,doc2,doc3")
public static class IndexReducer
extends Reducer<Text, Text, Text, Text> {
public void reduce(Text key, Iterable<Text> values, Context context)
throws IOException, InterruptedException {
Set<String> docs = new HashSet<>();
for (Text val : values) {
Collections.addAll(docs, val.toString().split(","));
}
context.write(key, new Text(String.join(",", docs)));
}
}
}
核心模式:按 Key 分组聚合
倒排索引、销售额统计、用户行为分析的本质都是同一个模式:
Map:提取 Key-Value 对
Shuffle:按 Key 分组
Reduce:对每个 Key 的 Value 列表做聚合
// 销售额统计:同样的模式
public static class SalesMapper extends Mapper<...> {
public void map(Object key, Text value, Context context) {
// 输入:orderId,userId,amount
// Map:(userId, amount)
String[] fields = value.toString().split(",");
context.write(new Text(fields[1]), new DoubleWritable(Double.parseDouble(fields[2])));
}
}
public static class SalesReducer extends Reducer<...> {
public void reduce(Text key, Iterable<DoubleWritable> values, Context context) {
// Reduce:对每个用户的金额求和
double total = 0;
for (DoubleWritable val : values) total += val.get();
context.write(key, new DoubleWritable(total));
}
}
理解了这个模式
"按 Key 分组聚合"是大数据处理的核心范式。MapReduce、Spark、Flink、Hive 都是这个思想的变体。理解了它就能理解整个大数据生态的逻辑内核。