哎,说起MapReduce,很多刚入行大数据的朋友可能觉得它是个“老古董”。但在我看来,它依然是理解分布式计算最纯粹的教科书。我见过太多人,MapReduce写得飞起,结果作业跑了三天三夜还没完,或者因为数据倾斜直接OOM(内存溢出)炸掉集群。今天我就掰开揉碎,从最基础的WordCount讲起,一路聊到数据去重的坑,最后重点攻克那个让无数工程师头秃的“数据倾斜”问题。咱们不整虚的,直接上干货。
从WordCount说起:别被简单的例子骗了
先回顾一下WordCount,这是MapReduce的“Hello World”。逻辑很简单:Map阶段把每一行文本拆成单词,输出<单词, 1>;Reduce阶段把这些1加起来。看起来人畜无害对吧?
但真正的挑战从这里就开始。假设你要处理一个10TB的日志文件,每个日志行是一个用户点击事件。如果你直接写个标准的WordCount,Map输出<userId, 1>,那么Shuffle阶段就会把所有相同userId的数据发给同一个Reduce Task。
问题来了:热门用户(比如某个网红)可能有几亿次点击,而冷门用户只有几次。这时候,负责处理那个网红userId的Reduce Task会处理海量数据,而其他Reduce Task可能几秒钟就处理完了。结果就是:整个作业时间取决于最慢那个Task,也就是木桶效应。这就是数据倾斜的雏形。
我在某次实际项目中就遇到过这种情况。用户点击日志里,Top 1%的用户贡献了90%的数据量。如果直接用标准WordCount聚合,那个处理Top用户的Reduce Task跑了一整夜,其他Task早就闲得发慌了。
数据去重:你以为你懂,其实你可能掉坑里了
接下来聊聊数据去重。很多人觉得这很简单:Map输出<userId, null>,Reduce里put到一个Set里,最后输出key就行。
但这里有个大坑:内存爆炸。
假设某个userId有1亿条重复记录,Map阶段输出1亿个<userId, null>,经过Shuffle,这些全部到Reducer。如果你用Set去重,需要把这些key都加载到内存。1亿个userId字符串,每个假设32字节(实际上可能更长),那就是3.2GB内存,这还是只处理一个key的情况!如果有多个热门key,Reducer直接OOM。
解决方案一:Map端预聚合
别急,我们有办法。在Map端,我们可以先用一个本地HashMap预聚合。这样,每个Map Task输出的是<userId, count>,而不是<userId, null>。虽然Shuffle的数据量还是很大(因为还要发count),但至少我们知道了重复次数。
但等等,这还不够解决去重问题。如果只是为了去重,Map端预聚合后,Reduce端还是要把所有相同的userId合并。这时候,我们可以用二次排序或者Bloom Filter来优化。
解决方案二:Bloom Filter过滤
Bloom Filter是一种概率型数据结构,可以用来判断一个元素是否在集合中。它的优点是内存占用极小,缺点是有误判率(假阳性)。
在MapReduce中,我们可以这样设计:
- 第一遍MapReduce:用Bloom Filter统计每个userId的出现情况,输出
<userId, count>。 - 第二遍MapReduce:如果count == 1,则保留;如果count > 1,说明有重复,这时候我们再决定怎么处理(可能需要进一步去重逻辑)。
但这种方式需要两遍作业,开销较大。而且,如果我们要的是“去重后的列表”,而不是“每个userId出现次数”,那这个方法就不太合适。
解决方案三:自定义Partitioner + Combiner
让我们回到最根本的问题:数据倾斜。为什么会出现数据倾斜?因为Partitioner默认用的是hash(key) % numReduceTasks。当key分布不均匀时,某些partition会特别大。
解决办法是自定义Partitioner。我们可以根据key的访问热度,把热门key分散到多个partition中。比如,我们可以对userId的hash值进行二次哈希,把热门key映射到不同的Reduce Task。
但更优雅的方案是:利用Combiner + 自定义OutputFormat。
等等,Combiner只能用于可结合的操作(比如求和、求最大值),去重操作是不可结合的。所以Combiner在这里帮不上忙。
那怎么办?让我给你看一个实际的代码案例。
实战:如何处理热点key的数据去重
假设我们有用户点击日志,格式如下:
userId\tclickCount
我们要统计每个userId的去重后的clickCount(其实这里就是求和,但假设场景是去重后的唯一值)。
方案A:Map端预聚合 + 自定义Partitioner
public class DedupMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
private static final IntWritable ONE = new IntWritable(1);
@Override
protected void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
String line = value.toString();
String[] parts = line.split("\t");
String userId = parts[0];
int clickCount = Integer.parseInt(parts[1]);
// Map端预聚合:按userId累加
// 注意:这里其实不需要去重,因为输入已经是聚合后的
// 但假设输入是原始日志,我们需要先去重
context.write(new Text(userId), ONE);
}
}
public class DedupPartitioner extends Partitioner<Text, IntWritable> {
@Override
public int getPartition(Text key, IntWritable value, int numPartitions) {
// 自定义分区策略:避免热点key集中到同一个partition
// 方案:对key的hash值进行二次哈希,分散到不同partition
int hash = key.hashCode();
// 使用多个hash函数,增加分散度
int partition = (hash & Integer.MAX_VALUE) % numPartitions;
return partition;
}
}
等等,我刚才写的Partitioner还是默认逻辑。真正解决热点key的问题是:我们需要在Map端就识别出热点key,并把它们分散。
方案B:两遍MapReduce + Bloom Filter(推荐)
这是我在生产环境中验证过最有效的方案。
第一遍:统计每个userId的频率
public class FrequencyMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
private static final IntWritable ONE = new IntWritable(1);
@Override
protected void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
String line = value.toString();
String userId = line.split("\t")[0];
context.write(new Text(userId), ONE);
}
}
public class FrequencyReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
@Override
protected void reduce(Text key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
int sum = 0;
for (IntWritable val : values) {
sum += val.get();
}
context.write(key, new IntWritable(sum));
}
}
第二遍:根据频率过滤
对于频率<=阈值的key,直接输出;对于频率>阈值的key,我们需要特殊处理。
但这里有个问题:如果某个key频率极高,第二遍的Reducer还是会倾斜。
方案C:使用Secondary Sort + 自定义Comparator
这是最优雅的方案。我们利用MapReduce的排序特性,在Shuffle阶段就把相同key的数据合并,然后在Reducer中只保留第一个value(因为我们要去重)。
public class DedupWritable implements WritableComparable<DedupWritable> {
private Text key;
private IntWritable value;
// 比较逻辑:先按key排序,再按value排序
@Override
public int compareTo(DedupWritable o) {
int cmp = this.key.compareTo(o.key);
if (cmp != 0) {
return cmp;
}
return this.value.compareTo(o.value);
}
}
等等,这个思路有点绕。让我重新梳理。
真正的解决方案:Partitioner + 自定义GroupingComparator
核心思想是:在Shuffle阶段,把相同key的数据分到同一个Reducer,但在Reducer内部,我们只处理一次。
但如何避免内存爆炸?答案是:使用外部排序 + 流式处理。
在Reducer中,我们不把所有数据加载到内存,而是:
- 对输入数据进行排序(使用外部排序,当内存不够时写入磁盘)
- 遍历排序后的数据,跳过重复的key
- 输出唯一key
这需要自定义SortComparator和GroupingComparator。
public class DedupReducer extends Reducer<Text, NullWritable, Text, NullWritable> {
private Text lastKey = new Text();
@Override
protected void reduce(Text key, Iterable<NullWritable> values, Context context)
throws IOException, InterruptedException {
// 由于GroupingComparator保证相同key的数据会分组处理
// 我们只需要输出一次key即可
if (!key.equals(lastKey)) {
context.write(key, NullWritable.get());
lastKey.set(key);
}
}
}
但这里有个关键配置:
// 设置GroupingComparator,使得相同key的数据被分组到同一个Reducer
job.setGroupingComparatorClass(KeyGroupingComparator.class);
// 自定义Partitioner,确保相同key的key被分到同一个Reducer
job.setPartitionerClass(HotKeyPartitioner.class);
数据倾斜的终极解决方案
好了,前面说了一堆,让我直接告诉你生产环境中最有效的方法:
方法一:加盐(Salting)
这是解决数据倾斜最经典的方法。思路是:
- 在Map阶段,给热点key加上随机后缀(盐),把一个大key拆分成多个小key
- 在Reduce阶段,先去掉盐,再聚合
public class SaltedMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
private static final Random RANDOM = new Random();
private static final int SALT_COUNT = 10; // 盐的数量
@Override
protected void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
String line = value.toString();
String userId = line.split("\t")[0];
// 判断是否是热点key(这里可以预先知道热点key列表)
if (isHotKey(userId)) {
// 给热点key加盐,分散到不同partition
int salt = RANDOM.nextInt(SALT_COUNT);
context.write(new Text(userId + "_" + salt), ONE);
} else {
// 非热点key正常输出
context.write(new Text(userId), ONE);
}
}
private boolean isHotKey(String userId) {
// 实际生产中,这里可以是从HDFS加载的热点key列表
return hotKeys.contains(userId);
}
}
然后,我们需要两个Reduce阶段:
- 第一遍Reduce:处理加盐后的key,每个key只输出一次(去重)
- 第二遍Reduce:去掉盐,重新聚合
public class UnsaltReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
@Override
protected void reduce(Text key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
// 去掉盐
String originalKey = key.toString().split("_")[0];
context.write(new Text(originalKey), ONE);
}
}
方法二:使用Hive/Spark替代(如果可能)
说实话,如果业务允许,我强烈建议用Hive或Spark。它们内置了数据倾斜优化:
- Hive:启用
hive.groupby.skewindata=true - Spark:使用
repartition或salting
但如果你必须用原生MapReduce,那前面的方法就是你要用的。
性能调优的实战技巧
除了算法层面,还有几个配置可以大幅提升性能:
1. 调整Reducer数量
默认情况下,MapReduce的Reducer数量是1。对于大数据集,这会导致巨大的瓶颈。
// 根据数据量和平均每个key的数据量来设置
int numReducers = (int) (totalDataSize / (100 * 1024 * 1024)); // 假设每个Reducer处理100MB
job.setNumReduceTasks(Math.max(numReducers, 10)); // 至少10个
2. 启用Combiner
虽然Combiner不能用于去重,但可以用于求和、求最大值等操作。
job.setCombinerClass(SumReducer.class);
3. 调整JVM复用
// 设置每个JVM可以运行多个Map Task,减少JVM启动开销
job.setBoolean("mapreduce.job.jvm.numtasks", 10);
4. 使用压缩
// 启用输出压缩
job.setCompressOutput(true);
job.setOutputCompressorClass(GzipCodec.class);
一个完整的去重案例
让我给你一个完整的、可运行的MapReduce去重代码框架:
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import java.io.IOException;
import java.util.Random;
public class DataDeduplication {
public static class DedupMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
private static final IntWritable ONE = new IntWritable(1);
private static final Random RANDOM = new Random();
private static final int SALT_COUNT = 10;
// 热点key列表(实际生产中可以从HDFS加载)
private static final String[] HOT_KEYS = {"user_123", "user_456"};
@Override
protected void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
String userId = value.toString();
// 判断是否是热点key
boolean isHot = false;
for (String hotKey : HOT_KEYS) {
if (hotKey.equals(userId)) {
isHot = true;
break;
}
}
if (isHot) {
// 给热点key加盐,分散到不同partition
int salt = RANDOM.nextInt(SALT_COUNT);
context.write(new Text(userId + "_" + salt), ONE);
} else {
// 非热点key正常输出
context.write(new Text(userId), ONE);
}
}
}
public static class DedupReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
@Override
protected void reduce(Text key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
// 去掉盐(如果有)
String originalKey = key.toString().split("_")[0];
context.write(new Text(originalKey), ONE);
}
}
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
Job job = Job.getInstance(conf, "Data Deduplication");
job.setJarByClass(DataDeduplication.class);
job.setMapperClass(DedupMapper.class);
job.setReducerClass(DedupReducer.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);
// 设置Reducer数量为热点key数量的10倍
job.setNumReduceTasks(100);
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}
总结:避坑要点
- 数据倾斜是MapReduce的最大敌人:时刻关注你的key分布,提前识别热点key。
- 加盐是解决倾斜的有效手段:把热点key拆分成多个子key,分散到不同Reducer。
- 不要把所有数据加载到内存:使用外部排序、流式处理,避免OOM。
- 合理设置Reducer数量:根据数据量和硬件资源调整。
- 考虑使用更高级的工具:如果可能,用Hive或Spark,它们内置了倾斜优化。
记住,MapReduce虽然“古老”,但它的思想依然有价值。理解这些原理,你在面对更复杂的分布式计算框架时,会更有信心。毕竟,万丈高楼平地起,地基打好了,上面的东西自然稳。
希望这些经验能帮到你。如果在实际项目中遇到具体问题,欢迎继续交流。咱们一起把大数据玩得转!