一、先搞懂啥是Flink实时ETL里的“数据倾斜”
很多刚接触实时计算的朋友,跑Flink任务时都会遇到一个怪事儿:任务跑着跑着突然慢得像蜗牛,甚至直接卡死。打开Flink的WebUI一看,好几个任务槽(就是跑计算的“工位”)的CPU、内存都拉满了,另外几个却闲得摸鱼,这就是数据倾斜在搞鬼。
举个最常见的例子:你要统计实时订单的商品销量,得按商品ID分组(也就是KeyBy)来算每个商品卖了多少。结果突然有个爆款商品,比如“某网红奶茶”的订单量是其他商品的几百倍,所有带这个商品ID的订单,都会被分到同一个工位上处理,这个工位瞬间被压垮,其他工位又帮不上忙,整个任务的速度就被这个慢工位拖死了。
说白了,数据倾斜就是“工作量分配不均”,本来大家轮着搬砖,结果有人搬了一卡车,其他人只搬了几块砖,效率自然上不去。
二、第一步:定位数据倾斜的根因
要解决问题,得先找到问题在哪。很多人一看到任务慢就瞎调参数,其实大部分时候都是KeyBy阶段的分配出了问题,所以第一步得从KeyBy的分布开始查。
2.1 用Flink自带工具查Key分布
Flink的WebUI有个很实用的功能,能直接看每个Key的处理量。操作很简单:打开WebUI,找到你的任务,点“SubTasks”(子任务)标签,再选你怀疑倾斜的算子(比如KeyBy之后的聚合算子),就能看到每个子任务的“Records Received”(收到的记录数)。如果有一个子任务的记录数是其他的几十上百倍,那肯定是Key分布不均。
举个具体的场景:假设你有8个子任务(8个工位),其中7个每个收到1万条订单,第8个收到100万条,那第8个对应的Key就是倾斜的根源。
2.2 排查倾斜Key的具体特征
找到倾斜的子任务后,接下来要确定是哪个Key导致的。你可以在KeyBy之后加一个打印日志的步骤,把Key和对应的数量打出来。这里用Java写个简单的测试代码,技术栈统一用Java 1.8 + Flink 1.17(因为这是目前最稳定的版本之一)。
技术栈:Java 1.8, Flink 1.17
// 假设你的订单数据是Order类,里面有商品ID(itemId)和订单金额(amount)
public class Order {
private String itemId;
private double amount;
// 构造函数、getter、setter省略
}
// 在Flink任务的KeyBy之后加这个算子,打印每个Key的数量
orders.keyBy(Order::getItemId)
.process(new ProcessFunction<Order, Order>() {
private transient Map<String, Long> keyCountMap; // 临时存储每个Key的数量
@Override
public void open(Configuration parameters) throws Exception {
keyCountMap = new HashMap<>(); // 任务启动时初始化
}
@Override
public void processElement(Order value, Context ctx, Collector<Order> out) throws Exception {
String itemId = value.getItemId();
// 统计当前Key的数量,每次加1
keyCountMap.put(itemId, keyCountMap.getOrDefault(itemId, 0L) + 1);
// 每10秒打印一次,避免日志太多
if (ctx.timerService().currentProcessingTime() % 10000 == 0) {
// 只打印数量超过1万的Key,方便找到倾斜的那个
keyCountMap.entrySet().stream()
.filter(entry -> entry.getValue() > 10000)
.forEach(entry -> System.out.println("倾斜Key:" + entry.getKey() + ",数量:" + entry.getValue()));
}
out.collect(value); // 把数据传给下一个算子
}
});
跑这个代码后,你就能从日志里看到,比如“itemId=奶茶001”的数量是120万,其他Key最多才2万,那这个奶茶001就是导致倾斜的罪魁祸首。
三、第二步:解决数据倾斜的核心方案
找到问题后,我们有三个常用的解决办法:调整KeyBy的逻辑、自定义分区器、两阶段聚合,下面一个个说。
3.1 调整KeyBy逻辑:给倾斜Key加“随机前缀”
这是最简单的办法,适合倾斜Key只有少数几个的场景。原理很简单:本来所有带“奶茶001”的订单都分到一个工位,现在我们给它加个随机前缀,比如“1_奶茶001”“2_奶茶001”“3_奶茶001”,这样不同前缀的同一个商品订单,就会被分到不同的工位,把压力分摊开。
举个具体的代码例子,还是用刚才的订单场景: 技术栈:Java 1.8, Flink 1.17
orders.keyBy(order -> {
String itemId = order.getItemId();
// 如果是倾斜的Key(比如奶茶001),加随机前缀;其他Key不变
if ("奶茶001".equals(itemId)) {
// 生成1-10的随机数,作为前缀,把压力分摊到10个工位
int randomPrefix = new Random().nextInt(10) + 1;
return randomPrefix + "_" + itemId;
} else {
return itemId;
}
})
// 这里做第一次聚合,比如统计每个前缀分组的销量
.sum("amount")
// 接下来要把不同前缀的同一个商品销量加起来,所以要再按原Key分组
.keyBy(sumOrder -> sumOrder.getItemId().split("_")[1]) // 去掉前缀,得到原商品ID
.sum("amount"); // 第二次聚合,得到最终的商品总销量
这个办法的优点是简单,不用改太多代码,缺点是如果倾斜Key太多,比如有几百个倾斜Key,加前缀的逻辑会很麻烦,而且随机数的范围不好控制,范围太大可能导致聚合效率低,范围太小还是会有倾斜。
3.2 自定义分区器:给倾斜Key单独分配工位
如果倾斜Key的数量不多,你也可以自己写分区器,专门给倾斜Key分配特定的工位,其他Key按默认逻辑分配。Flink的分区器是用来决定数据分到哪个子任务的,默认的分区器是按Key的哈希值取模,我们可以重写这个逻辑。
代码例子: 技术栈:Java 1.8, Flink 1.17
// 自定义分区器,继承Flink的Partitioner类
public class CustomPartitioner implements Partitioner<String> {
@Override
public int partition(String key, int numPartitions) {
// numPartitions是子任务的总数,比如8
if ("奶茶001".equals(key)) {
// 给奶茶001分配第0号工位(专门留一个工位给它)
return 0;
} else if ("蛋糕001".equals(key)) {
// 给蛋糕001分配第1号工位
return 1;
} else {
// 其他Key按默认的哈希取模分配
return Math.abs(key.hashCode()) % numPartitions;
}
}
}
// 在Flink任务中使用自定义分区器
orders.keyBy(Order::getItemId)
// 调用partitionCustom方法,传入自定义分区器和Key的类型
.partitionCustom(new CustomPartitioner(), Order::getItemId)
.sum("amount");
这个办法的优点是精准,能专门给倾斜Key分配资源,缺点是如果倾斜Key经常变(比如今天是奶茶001,明天是面包002),你得频繁修改代码,不太灵活。
3.3 两阶段聚合:最通用的解决方案
如果倾斜Key很多,或者你不知道哪个Key会倾斜,两阶段聚合是最靠谱的办法。原理是把聚合分成两步:第一步先做“局部聚合”,把同一个Key的多个数据先在本地合并成一个中间结果,再把中间结果做“全局聚合”,这样能大大减少需要传输的数据量,避免单个工位压力过大。
举个具体的例子,比如统计订单的商品销量,两阶段聚合的流程是:
- 第一步:每个子任务先统计自己收到的订单,比如A子任务里奶茶001卖了1万,B子任务里奶茶001卖了1.2万,都先记下来,生成中间结果(比如“奶茶001:1万”“奶茶001:1.2万”)。
- 第二步:把所有子任务的中间结果,再按商品ID分组,加起来得到最终的销量(1万+1.2万=2.2万)。
代码例子: 技术栈:Java 1.8, Flink 1.17
// 第一步:局部聚合,用sum算子做本地统计
orders.keyBy(Order::getItemId)
.sum("amount") // 每个子任务先统计自己的商品销量
// 把中间结果转成一个新的类,方便后续处理
.map(sumOrder -> new IntermediateResult(sumOrder.getItemId(), sumOrder.getAmount()))
// 第二步:全局聚合,再按商品ID分组,加总所有中间结果
.keyBy(IntermediateResult::getItemId)
.sum("amount");
// 中间结果类,用来存储局部聚合的结果
public class IntermediateResult {
private String itemId;
private double amount;
// 构造函数、getter、setter省略
}
这个办法的优点是通用,不管哪个Key倾斜都能解决,因为局部聚合已经把数据量压缩了,缺点是需要多一步聚合,代码稍微复杂一点,但现在Flink的一些高级聚合算子(比如Window)已经内置了两阶段聚合的逻辑,用起来也很方便。
四、应用场景、优缺点和注意事项
4.1 应用场景
数据倾斜几乎在所有实时ETL场景中都可能出现,最常见的有:
- 统计类场景:比如商品销量、用户访问量、订单金额等按Key分组的统计。
- 关联类场景:比如两个流关联(比如订单流和商品信息流关联),如果其中一个流的Key分布不均,也会导致关联算子倾斜。
- 窗口计算场景:比如按时间窗口统计,窗口内某个Key的数据量特别大,也会导致窗口计算倾斜。
4.2 各方案的优缺点对比
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 加随机前缀 | 简单易实现,代码改动小 | 倾斜Key多的时候麻烦,随机范围难控制 | 倾斜Key少且固定的场景 |
| 自定义分区器 | 精准分配资源,效率高 | 倾斜Key变化时要改代码,不灵活 | 倾斜Key少且固定的场景 |
| 两阶段聚合 | 通用,不管哪个Key倾斜都能解决 | 代码稍复杂,多一步聚合 | 倾斜Key多、不确定的场景,通用场景 |
4.3 注意事项
- 随机前缀的范围要合适:比如你有8个子任务,随机前缀的范围设为10就差不多,范围太小还是会有倾斜,范围太大可能导致局部聚合的效率低。
- 自定义分区器要注意子任务总数:比如你留了第0号工位给奶茶001,那子任务总数不能小于2,否则第0号工位还是会处理其他Key的数据。
- 两阶段聚合要注意中间结果的大小:如果中间结果太大,还是会有问题,所以局部聚合要尽量把数据压缩到最小。
- 避免空Key:很多时候倾斜是因为Key为空,比如订单的商品ID为空,所有空Key都会分到同一个工位,所以要先过滤掉空Key,或者给空Key加一个默认值。
五、文章总结
数据倾斜是Flink实时ETL中最常见的问题之一,解决它的核心思路是“先定位,再选方案”:先通过Flink的WebUI和日志找到倾斜的根源,再根据倾斜Key的数量和变化情况,选择合适的解决方案。
如果倾斜Key少且固定,可以选加随机前缀或者自定义分区器;如果倾斜Key多或者不确定,两阶段聚合是最靠谱的选择。另外,解决数据倾斜不是一劳永逸的,你需要定期监控任务的运行情况,及时调整方案。
最后,给大家一个小技巧:跑Flink任务时,一定要先做压测,模拟真实的流量情况,提前发现数据倾斜的问题,避免上线后出问题。
评论
围绕“Flink实时ETL中数据倾斜的定位与治理:从KeyBy分布到自定义分区器及两阶段聚合方案的详细解析与实战完整指南及优化技巧”参与讨论