Hadoop 数据倾斜:加盐(Salting)
加盐本质:对热点 key 增加随机前缀,打散 shuffle,两阶段聚合。
适用:
group by、count/sum类聚合倾斜;求平均值不能直接用加盐。
分角色:开发(代码 / SQL 层面改造)、运维(参数、平台配置、监控,不改动业务逻辑)。
原理回顾
热点 key hot_key,100w 条,全部落到同一个 task。加盐:rand_N_hot_key,N 取 0~N-1,把 1 个热点拆成 N 个不同 key,分散到 N 个 task 做局部聚合;第二阶段去掉随机数,全局聚合。
一、开发侧工作(改业务 SQL / 代码,核心解决手段)
场景 1:Hive SQL group by 聚合倾斜,加盐完整案例
业务 SQL(倾斜版本)
-- 原始SQL,user_id存在热点,某个user_id数据量巨大,发生倾斜
select user_id,count(*) as pv
from user_log
group by user_id;
加盐改造(两阶段聚合)
--第一阶段:加盐局部聚合,给key拼接随机前缀
with partial_agg as (
select
-- 加盐:随机数0‑9,一共10个分片;非热点key也加盐,简单粗暴
concat(floor(rand()*10),'_',user_id) as salt_userid,
count(*) as cnt
from user_log
group by floor(rand()*10),user_id
)
--第二阶段:去除盐值,全局聚合
select
split(salt_userid,'_')[1] as user_id,
sum(cnt) as pv
from partial_agg
group by split(salt_userid,'_')[1];
说明:
floor(rand()*10)盐值范围 0-9,拆成 10 份;热点 key 分散到 10 个 reduce。
优化:只给热点 key 加盐,普通 key 不加盐,减少数据膨胀。
with partial_agg as (
select
case
when user_id='hot_user001' then concat(floor(rand()*10),'_',user_id)
else user_id
end as salt_userid,
count(*) as cnt
from user_log
group by
case
when user_id='hot_user001' then concat(floor(rand()*10),'_',user_id)
else user_id
end
)
select
if(instr(salt_userid,'_')>0,split(salt_userid,'_')[1],salt_userid) as user_id,
sum(cnt) as pv
from partial_agg
group by if(instr(salt_userid,'_')>0,split(salt_userid,'_')[1],salt_userid);
场景 2:Spark SQL 加盐,和 Hive 几乎完全一样
val sql = """
with partial as (
select concat(floor(rand()*10),'_',user_id) salt_uid,count(*) cnt
from user_log
group by floor(rand()*10),user_id
)
select split(salt_uid,'_')[1] user_id,sum(cnt) pv
from partial
group by split(salt_uid,'_')[1]
"""
spark.sql(sql).show()
场景 3:Spark RDD 代码加盐案例
val originRDD:RDD[(String,Int)] = ... // (userid,1)
// 第一阶段:加盐
val saltedRDD = originRDD.map{
case (key,value)=>
val salt = scala.util.Random.nextInt(10) //盐0‑9
(s"${salt}${key}",value)
}.reduceByKey(+_) //局部聚合
//第二阶段:去掉盐,全局聚合
val resultRDD = saltedRDD.map{
case(saltedKey,cnt)=>
val key = saltedKey.split("")(1)
(key,cnt)
}.reduceByKey(+_)
resultRDD.collect()
场景 4:Join 类型数据倾斜加盐(大表 join 热点 key)
join 加盐规则:大表热点 key 加随机盐;小表必须复制膨胀同等份数,才能关联上 示例:大表 A,小表 B,join key=user_id,user_id=hot01 是热点。
--大表A:热点key加0‑9随机盐
select
case when user_id='hot01' then concat(floor(rand()*10),'_',user_id)
else user_id end as salt_uid,
*
from tableA
--小表B:把热点key复制10份,拼接0‑9的盐值,普通key不变
select explode(array(0,1,2,3,4,5,6,7,8,9)) as salt,concat(salt,'_',user_id) salt_uid,*
from tableB where user_id='hot01'
union all
select null as salt,user_id as salt_uid,* from tableB where user_id!='hot01'
之后两张表用
salt_uid做 join。
总结:开发侧做
- 识别热点 key:通过任务监控、日志、采样,找到倾斜的 key;
- SQL/RDD 改造,实现两阶段加盐聚合;join 倾斜需要大小表配套改造;
- 控制盐值数量:不是越大越好,10、20、50;盐值过大生成大量小 task;
- 优先只对热点 key 加盐,减少全量数据膨胀;
- 测试验证结果正确性:sum/count 正确;avg 不能直接加盐;
- 评估资源:加盐多一轮 shuffle,资源会上涨,预估资源调整参数。
缺点:加盐会多一轮 shuffle,增加 CPU、内存、磁盘 IO 开销。
二、运维可以做什么(不改业务代码,平台 / 参数 / 监控层面)
运维不能写 SQL 改业务逻辑,主要做:事前监控告警、平台参数调优、任务诊断、框架自带倾斜优化开关、资源调优、兜底策略。
运维不会直接写加盐 SQL,但可以开启框架内置的自动加盐能力(Spark AQE)。
1. 事前:监控,识别数据倾斜
运维配置 YARN/Spark/Hive 监控:
- 监控 Stage/Reduce 任务时长分布:同一个 stage task 时长差异巨大(有的几秒,有的几十分钟),判定倾斜;
- YARN 任务指标:查看每个 reduce 输入记录数,单个 task 输入数据远超平均值;
- 采集慢任务告警:任务卡在 99%、task 长时间 running;
- 保存 YARN 日志、SparkUI 历史服务,方便定位热点 key;
- Hive/Spark 历史服务器,保存已完成任务 UI。
2. 开启框架自带自动倾斜优化(内置自动加盐打散,不用开发写盐值)
Spark AQE(自适应执行,运维开启参数,Spark3+)
AQE 内部会自动检测倾斜分区,底层自动加盐拆分倾斜分区,不需要开发手写 rand ()。 运维在spark‑defaults.conf配置,或者任务 session 参数:
#开启AQE总开关
spark.sql.adaptive.enabled=true
#开启倾斜处理,自动打散倾斜分区(内部自动加盐机制)
spark.sql.adaptive.skewJoin.enabled=true
#判定为倾斜的阈值:分区大小大于中位数3倍以上,判定倾斜
spark.sql.adaptive.skewJoin.skewedPartitionFactor=3
#倾斜分区最小大小
spark.sql.adaptive.skewedPartitionThresholdInBytes=256MB
#动态合并小分区
spark.sql.adaptive.coalescePartitions.enabled=true
⚠️注意:AQE 自动倾斜优化,只对 SparkSQL 生效;原生 RDD 代码无效,RDD 必须开发手动加盐。
Hive 侧运维参数(Hive 没有自动加盐,只能做缓解,不能根治)
Hive 不支持自动加盐,运维只能做缓解:
--倾斜join优化,对join倾斜做拆分(不是加盐)
set hive.optimize.skewjoin=true;
set hive.skewjoin.key=100000;
--group by倾斜优化,开启map端局部聚合(combiner)
set hive.map.aggr=true;
--倾斜任务允许更多重试次数
set mapreduce.map.maxattempts=4;
set mapreduce.reduce.maxattempts=4;
hive
hive.optimize.skewjoin是把热点 key 单独走 MapJoin,不是加盐。
3. 资源层面兜底(运维)
当倾斜无法快速修改代码时,临时兜底:
- 调大 reduce/executor 内存、堆外内存,防止 OOM;
spark.executor.memory=8g
spark.executor.memoryOverhead=2g
MR:
mapreduce.reduce.memory.mb=8192
mapreduce.reduce.java.opts=-Xmx6g
只能扛住,不能解决倾斜,只是不容易 OOM,任务依然跑很慢。
- 调整并行度,适当增大 task 数量。
4. 平台规范与支持
- 给开发输出规范文档:数据倾斜处理方案文档,加盐使用场景、示例 SQL;
- 平台拦截:识别慢任务,给开发告警推送;
- 临时应急:如果任务反复失败,可以调整队列资源优先级;
- 协助定位热点 key:运维拉取任务日志、SparkUI,提取倾斜 key 交给开发,开发再做加盐改造。
5. 运维的边界
❌运维不能修改业务 SQL 代码,不会手动写加盐逻辑。
✅运维可以:开启 AQE 自动加盐(SparkSQL)、监控告警、定位热点 key、调资源参数、输出规范文档。
重点区分:Spark AQE 的自动倾斜拆分,框架底层自动实现加盐打散,开发不用写 rand (),运维开启参数即可生效。
三、开发 vs 运维职责对比总结
| 角色 | 做什么 | 案例操作 | 局限 |
|---|---|---|---|
| 开发 |
业务代码改造,手动加盐; 识别热点 key; 验证结果正确性 | 写两阶段 SQL/RDD;join 倾斜大小表同步膨胀;只给热点 key 加盐;控制盐值数量 | 需要修改业务代码;多一轮 shuffle,资源上涨;RDD 代码只能手动加盐 |
| 运维 |
监控告警; 定位倾斜 key; 开启 Spark AQE 自动倾斜(底层自动加盐); Hive 倾斜缓解参数; 调资源兜底; 输出开发规范文档 | 配置 spark.sql.adaptive.skewJoin.enabled=true;监控 task 运行时长分布;拉取日志输出热点 key;调 executor 内存 | Hive、Spark RDD 不支持自动加盐;只能缓解,无法根治业务逻辑导致倾斜;不能改业务代码 |
精简版
开发:识别热点 key,业务 SQL/RDD 手动加盐,两阶段聚合;group by 给热点 key 拼接随机前缀打散,局部聚合后去除盐做全局聚合;join 倾斜时大表加盐,小表同步膨胀相同份数。控制盐值数量,优先只给热点 key 加盐,评估 shuffle 资源开销。
运维:不能改业务代码;做监控,识别慢任务、定位倾斜 key;SparkSQL 开启 AQE,框架底层自动加盐打散倾斜分区;Hive 只能开启 skewjoin 做缓解,无自动加盐;调大内存并行度做临时兜底;输出开发规范。RDD 代码 AQE 不生效,只能开发手动加盐。
补充坑点
- 加盐不能直接用于 avg 求平均值,sum/count 可以分开聚合,最后再相除。
- Spark RDD 代码 AQE 无效,必须开发手动写加盐逻辑。
- Hive 没有自动加盐,只能开发手写 SQL 加盐。
- 盐值不是越大越好,盐值过多产生大量小 task,消耗资源。
- AQE 自动倾斜优化,只解决 SparkSQL 的 group by/join 倾斜,底层就是框架封装的加盐。
更多推荐
所有评论(0)