加盐本质:对热点 key 增加随机前缀,打散 shuffle,两阶段聚合。

适用:group bycount/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。

总结:开发侧做

  1. 识别热点 key:通过任务监控、日志、采样,找到倾斜的 key;
  2. SQL/RDD 改造,实现两阶段加盐聚合;join 倾斜需要大小表配套改造;
  3. 控制盐值数量:不是越大越好,10、20、50;盐值过大生成大量小 task;
  4. 优先只对热点 key 加盐,减少全量数据膨胀;
  5. 测试验证结果正确性:sum/count 正确;avg 不能直接加盐
  6. 评估资源:加盐多一轮 shuffle,资源会上涨,预估资源调整参数。

缺点:加盐会多一轮 shuffle,增加 CPU、内存、磁盘 IO 开销。


二、运维可以做什么(不改业务代码,平台 / 参数 / 监控层面)

运维不能写 SQL 改业务逻辑,主要做:事前监控告警、平台参数调优、任务诊断、框架自带倾斜优化开关、资源调优、兜底策略

运维不会直接写加盐 SQL,但可以开启框架内置的自动加盐能力(Spark AQE)。

1. 事前:监控,识别数据倾斜

运维配置 YARN/Spark/Hive 监控:

  1. 监控 Stage/Reduce 任务时长分布:同一个 stage task 时长差异巨大(有的几秒,有的几十分钟),判定倾斜;
  2. YARN 任务指标:查看每个 reduce 输入记录数,单个 task 输入数据远超平均值;
  3. 采集慢任务告警:任务卡在 99%、task 长时间 running;
  4. 保存 YARN 日志、SparkUI 历史服务,方便定位热点 key;
  5. 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. 资源层面兜底(运维)

当倾斜无法快速修改代码时,临时兜底:

  1. 调大 reduce/executor 内存、堆外内存,防止 OOM;
spark.executor.memory=8g
spark.executor.memoryOverhead=2g

MR:

mapreduce.reduce.memory.mb=8192
mapreduce.reduce.java.opts=-Xmx6g

只能扛住,不能解决倾斜,只是不容易 OOM,任务依然跑很慢。

  1. 调整并行度,适当增大 task 数量。

4. 平台规范与支持

  1. 给开发输出规范文档:数据倾斜处理方案文档,加盐使用场景、示例 SQL;
  2. 平台拦截:识别慢任务,给开发告警推送;
  3. 临时应急:如果任务反复失败,可以调整队列资源优先级;
  4. 协助定位热点 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 不生效,只能开发手动加盐。

补充坑点

  1. 加盐不能直接用于 avg 求平均值,sum/count 可以分开聚合,最后再相除。
  2. Spark RDD 代码 AQE 无效,必须开发手动写加盐逻辑
  3. Hive 没有自动加盐,只能开发手写 SQL 加盐。
  4. 盐值不是越大越好,盐值过多产生大量小 task,消耗资源。
  5. AQE 自动倾斜优化,只解决 SparkSQL 的 group by/join 倾斜,底层就是框架封装的加盐。

更多推荐