ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Hadoop 数据倾斜:加盐(Salting)

Hadoop 数据倾斜:加盐(Salting) 加盐本质对热点 key 增加随机前缀打散 shuffle两阶段聚合。适用group by、count/sum类聚合倾斜求平均值不能直接用加盐。分角色开发代码 / SQL 层面改造、运维参数、平台配置、监控不改动业务逻辑。原理回顾热点 keyhot_key100w 条全部落到同一个 task。加盐rand_N_hot_keyN 取 0~N-1把 1 个热点拆成 N 个不同 key分散到 N 个 task 做局部聚合第二阶段去掉随机数全局聚合。一、开发侧工作改业务 SQL / 代码核心解决手段场景 1Hive SQL group by 聚合倾斜加盐完整案例业务 SQL倾斜版本-- 原始SQLuser_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_idhot_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_idhot_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);场景 2Spark 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()场景 3Spark 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()场景 4Join 类型数据倾斜加盐大表 join 热点 keyjoin 加盐规则大表热点 key 加随机盐小表必须复制膨胀同等份数才能关联上示例大表 A小表 Bjoin keyuser_iduser_idhot01 是热点。--大表A热点key加0‑9随机盐 select case when user_idhot01 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_idhot01 union all select null as salt,user_id as salt_uid,* from tableB where user_id!hot01之后两张表用salt_uid做 join。总结开发侧做识别热点 key通过任务监控、日志、采样找到倾斜的 keySQL/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 历史服务方便定位热点 keyHive/Spark 历史服务器保存已完成任务 UI。2. 开启框架自带自动倾斜优化内置自动加盐打散不用开发写盐值Spark AQE自适应执行运维开启参数Spark3AQE 内部会自动检测倾斜分区底层自动加盐拆分倾斜分区不需要开发手写 rand ()。 运维在spark‑defaults.conf配置或者任务 session 参数#开启AQE总开关 spark.sql.adaptive.enabledtrue #开启倾斜处理自动打散倾斜分区内部自动加盐机制 spark.sql.adaptive.skewJoin.enabledtrue #判定为倾斜的阈值分区大小大于中位数3倍以上判定倾斜 spark.sql.adaptive.skewJoin.skewedPartitionFactor3 #倾斜分区最小大小 spark.sql.adaptive.skewedPartitionThresholdInBytes256MB #动态合并小分区 spark.sql.adaptive.coalescePartitions.enabledtrue⚠️注意AQE 自动倾斜优化只对 SparkSQL 生效原生 RDD 代码无效RDD 必须开发手动加盐。Hive 侧运维参数Hive 没有自动加盐只能做缓解不能根治Hive 不支持自动加盐运维只能做缓解--倾斜join优化对join倾斜做拆分不是加盐 set hive.optimize.skewjointrue; set hive.skewjoin.key100000; --group by倾斜优化开启map端局部聚合combiner set hive.map.aggrtrue; --倾斜任务允许更多重试次数 set mapreduce.map.maxattempts4; set mapreduce.reduce.maxattempts4;hivehive.optimize.skewjoin是把热点 key 单独走 MapJoin不是加盐。3. 资源层面兜底运维当倾斜无法快速修改代码时临时兜底调大 reduce/executor 内存、堆外内存防止 OOMspark.executor.memory8g spark.executor.memoryOverhead2gMRmapreduce.reduce.memory.mb8192 mapreduce.reduce.java.opts-Xmx6g只能扛住不能解决倾斜只是不容易 OOM任务依然跑很慢。调整并行度适当增大 task 数量。4. 平台规范与支持给开发输出规范文档数据倾斜处理方案文档加盐使用场景、示例 SQL平台拦截识别慢任务给开发告警推送临时应急如果任务反复失败可以调整队列资源优先级协助定位热点 key运维拉取任务日志、SparkUI提取倾斜 key 交给开发开发再做加盐改造。5. 运维的边界❌运维不能修改业务 SQL 代码不会手动写加盐逻辑。✅运维可以开启 AQE 自动加盐SparkSQL、监控告警、定位热点 key、调资源参数、输出规范文档。重点区分Spark AQE 的自动倾斜拆分框架底层自动实现加盐打散开发不用写 rand ()运维开启参数即可生效。三、开发 vs 运维职责对比总结角色做什么案例操作局限开发业务代码改造手动加盐识别热点 key验证结果正确性写两阶段 SQL/RDDjoin 倾斜大小表同步膨胀只给热点 key 加盐控制盐值数量需要修改业务代码多一轮 shuffle资源上涨RDD 代码只能手动加盐运维监控告警定位倾斜 key开启 Spark AQE 自动倾斜底层自动加盐Hive 倾斜缓解参数调资源兜底输出开发规范文档配置 spark.sql.adaptive.skewJoin.enabledtrue监控 task 运行时长分布拉取日志输出热点 key调 executor 内存Hive、Spark RDD 不支持自动加盐只能缓解无法根治业务逻辑导致倾斜不能改业务代码精简版开发识别热点 key业务 SQL/RDD 手动加盐两阶段聚合group by 给热点 key 拼接随机前缀打散局部聚合后去除盐做全局聚合join 倾斜时大表加盐小表同步膨胀相同份数。控制盐值数量优先只给热点 key 加盐评估 shuffle 资源开销。运维不能改业务代码做监控识别慢任务、定位倾斜 keySparkSQL 开启 AQE框架底层自动加盐打散倾斜分区Hive 只能开启 skewjoin 做缓解无自动加盐调大内存并行度做临时兜底输出开发规范。RDD 代码 AQE 不生效只能开发手动加盐。补充坑点加盐不能直接用于 avg 求平均值sum/count 可以分开聚合最后再相除。Spark RDD 代码 AQE 无效必须开发手动写加盐逻辑。Hive 没有自动加盐只能开发手写 SQL 加盐。盐值不是越大越好盐值过多产生大量小 task消耗资源。AQE 自动倾斜优化只解决 SparkSQL 的 group by/join 倾斜底层就是框架封装的加盐。
返回列表