Hive on Spark 企业级调优(二)SQL 优化(下):数据倾斜与进阶调优
前言
上一篇我们走完了 SQL 的"常规体检"——聚合提前做、数据少读、Join 选对算法、文件治好、并行度调对、CBO 兜底。这些手段的共同点是:计划在执行前就定死了。
但生产环境不讲武德。
你设了 200 个 Reduce,199 个 3 秒跑完,第 200 个跑了 40 分钟——因为 user_id = -1 的脏数据有 2 亿条全挤在一个 Task 里。CBO 统计信息说小表有 200 万行,实际过滤完只剩 3MB,白白多了一次全量 Shuffle。同样的 SQL,你的 Shuffle 比同事慢 5 倍,不是计划的问题,是 Java 序列化带着类名、字段名、继承链一起传,体积膨胀了 10 倍。
这一篇,我们解决"计划定死之后"的所有问题:
第 12 章 · 数据倾斜:5 种典型场景 × SQL 改写 + 5 套参数方案,从"加盐打散"到"AQE 自动拆分",逐个给代码。
第 13 章 · AQE 自适应执行:让引擎跑着跑着自己改计划——动态合并分区、动态切换 Join、自动拆分倾斜。
第 14 章 · 序列化与内存:Kryo 切换、内存模型、GC 调优——"最后一公里"决定了任务是从 15 分钟变 8 分钟,还是从 OOM 变跑通。
最后,我会把三篇的所有配置串成一份完整的 ETL 调优模板,一个文件,拿来就能贴进生产。
12. 数据倾斜
12.1 什么是数据倾斜
数据倾斜(Data Skew) 是指在分布式计算中,由于数据分布不均匀,导致 **少数 Task 处理的数据量远远大于其他的 Task ,**这些 “长尾” Task 成为整个作业的瓶颈。
外在表现形式:
- Spark UI 中,某个 Stage 的 Task 进度长时间停留在 99%
- 大部分 Task 在几秒内完成,但 1~2 个 Task 以运行数十分钟甚至数小时
- 最终可能因 内存溢出(OOM)导致任务失败
根本原因:
Shuffle 时,相同的 key 的数据被发送到同一个 Task,
如果某一个 key 对应的数据量是 其他key 的1000倍,
那么处理该 key 的 Task 就要处理 1000倍 的数据,导致数据倾斜。
如何排查:
30 秒速查:数据倾斜五步法
- 看 Spark UI → 哪个 Stage 有长尾 Task
- 查 Key 分布 → 确认倾斜 Key 是谁
- 判断场景 → Group By / Join / NULL / DISTINCT / 类型
- 选方案 → SQL 改写(根治) or 参数(兜底)
- 验证 → 对比 Task 耗时分布是否均匀
12.2 数据倾斜的典型场景
| 场景 | 触发操作 | 示例 |
|---|---|---|
| Group By 倾斜 | GROUP BY | 某个 province = ‘广东省’ 占了 80% 的数据 |
| Join 倾斜 | JOIN ON | 某个 user_id = 0(默认值)关联了百万条 |
| Count Distinct 倾斜 | COUNT(DISTINCT) | 全局去重,所有的数据到一个 Task |
| 空值/默认值 | JOIN ON | NULL 值全部哈希到同一个分区 |
| 数据类型不一致 | JOIN ON | INT vs STRING 隐式转化导致分布变化 |
12.3 如何定位数据倾斜
-- 方法1:查看Join Key的数据分布
SELECT
user_id,
COUNT(*) AS cnt
FROM orders
GROUP BY user_id
ORDER BY cnt DESC
LIMIT 20;
-- 方法2:查看Group By Key的分布
SELECT
city,
COUNT(*) AS cnt
FROM orders
GROUP BY city
ORDER BY cnt DESC
LIMIT 20;
-- 方法3:查看执行计划中的Task数量和数据量
EXPLAIN EXTENDED
SELECT city, COUNT(*) FROM orders GROUP BY city;
12.4 优化手段一 : 改写 SQL
12.4.1 场景一: Group By 数据倾斜
问题SQL:
-- 假设 city='北京' 的数据占了总量的60%
SELECT
city,
COUNT(*) AS order_cnt,
SUM(amount) AS total_amount
FROM orders
GROUP BY city;
优化改写: 两阶段聚合(加盐打散 + 去盐合并)
-- 第一阶段:给 key 加随机前缀,打散到多个 Task
SELECT
city,
SUM(partial_cnt) AS order_cnt,
SUM(partial_amount) AS total_amount
FROM(
SELECT
city,
COUNT(*) AS partial_cnt,
SUM(amount) AS partial_amount
FROM(
-- 加盐:将 '北京' 变成 ‘0_北京’,‘1_北京’,‘2_北京’,... ‘9_北京’,
SELECT
CONCAT( CAST(FLOOR(RAND() * 10) AS STRING), '_', city) AS city,
amount
FROM orders
) salted
GROUP BY city
) partial_result
-- 第二阶段: 去掉前缀,合并结果
GROUP BY SUBSTR(city, INSTR(city,'_') + 1);
更简洁的写法(利用 Hive 的内置函数)
SELECT
city,
SUM(cnt) AS order_cnt,
SUM(amt) AS total_amount
FROM (
SELECT
city,
COUNT(*) AS cnt,
SUM(amount) AS amt
FROM (
-- 加随机前缀 0~9
SELECT
CONCAT( CAST(FLOOR(RAND() * 10) AS STRING) , '_', city) AS city,
amount
FROM orders
) t1
GROUP BY city
) t2
-- 去掉前缀,合并结果
GROUP BY SPLIT(city, '\\_')[1];
12.4.2 场景二:Join 数据倾斜(大表 JOIN 大表,某些 Key 数据量极大)
问题SQL:
-- user_id=0 是默认值/脏数据,有500万条
-- 其他user_id正常分布,每个几千条
SELECT
o.order_id,
o.amount,
u.user_name
FROM orders o
JOIN users u ON o.user_id = u.user_id;
优化改写 方案A:过滤倾斜 Key 后单独处理
-- 思路: 将倾斜 Key 的数据分离出来单独处理,在 UNION ALL
-- part 1 :非倾斜数据正常join
SELECT
o.order_id,
o.amount,
u.user_name
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.user_id != 0
UNION ALL
-- part 2 :倾斜 key 单独处理(如果 user_id = 0 在 users 表中只有 1条,可直接 MapJoin)
SELECT /*+ MAPJOIN(u) */
o.order_id,
o.amount,
u.user_name
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.user_id = 0;
优化改写 方案B:打散加盐(适用于倾斜Key在两边都多的情况)
-- 大表加随机后缀,小表膨胀N倍
SELECT
o.order_id,
o.amount,
u.user_name
FROM (
-- 大表:给倾斜Key加随机后缀 0~9
SELECT
order_id,
amount,
CASE
WHEN user_id = 0
THEN CONCAT(CAST(user_id AS STRING), '_', CAST(FLOOR(RAND()*10) AS STRING))
ELSE CAST(user_id AS STRING)
END AS user_id_salted
FROM orders
) o
JOIN (
-- 小表:将倾斜Key膨胀10倍(炸裂函数)
SELECT
user_name,
CONCAT(CAST(user_id AS STRING), '_', CAST(t.pos AS STRING)) AS user_id_salted
FROM users
LATERAL VIEW EXPLODE(ARRAY(0,1,2,3,4,5,6,7,8,9)) t AS pos
WHERE user_id = 0
UNION ALL
SELECT
user_name,
CAST(user_id AS STRING) AS user_id_salted
FROM users
WHERE user_id != 0
) u ON o.user_id_salted = u.user_id_salted;
12.4.3 场景三:大量空值(NULL) 导致的数据倾斜
问题SQL:
-- user_id 有大量 NULL 值,NULL 的哈希值相同,全部进入同一个Task
SELECT
o.order_id,
o.amount,
u.user_name
FROM orders o
LEFT JOIN users u ON o.user_id = u.user_id;
优化改写:给 NULL值赋值随机 key
SELECT
o.order_id,
o.amount,
u.user_name
FROM (
SELECT
order_id,
amount,
-- NULL值赋予随机数,使其分散到不同Task
COALESCE(user_id, CONCAT('null_', CAST(FLOOR(RAND() * 100) AS STRING))) AS user_id
FROM orders
) o
LEFT JOIN users u ON o.user_id = u.user_id;
注意: 由于 users 表中不存在
null_0、null_1这样的 Key,所以 LEFT JOIN 结果中这些行的 user_name 仍为 NULL,语义不变。
12.4.4 场景 四:COUNT DISTINCT 导致的全局倾斜
问题SQL:
-- 所有数据都要到一个Task做去重,严重倾斜
SELECT
COUNT(DISTINCT user_id) AS uv
FROM orders
WHERE dt = '2026-07-25';
优化改写:GROUP BY 代替 DISTINCT
-- 方案1:子查询去重
SELECT COUNT(*) AS uv
FROM (
SELECT user_id
FROM orders
WHERE dt = '2026-07-25'
GROUP BY user_id
) t;
-- 方案2:多个COUNT DISTINCT同时存在时
-- ❌ 低效
SELECT
COUNT(DISTINCT user_id) AS uv,
COUNT(DISTINCT order_id) AS order_cnt
FROM orders;
-- ✅ 高效改写
SELECT
SUM(CASE WHEN tp = 'user' THEN 1 ELSE 0 END) AS uv,
SUM(CASE WHEN tp = 'order' THEN 1 ELSE 0 END) AS order_cnt
FROM (
SELECT
user_id AS id,
'user' AS tp
FROM orders
GROUP BY user_id
UNION ALL
SELECT
order_id AS id,
'order' AS tp
FROM orders
GROUP BY order_id
) t;
12.4.5 场景 五:不同数据类型隐式转换导致倾斜
问题SQL:
-- orders.user_id 是 BIGINT,users.user_id 是 STRING
-- 隐式转换可能导致哈希分布异常
SELECT o.order_id, u.user_name
FROM orders o
JOIN users u ON o.user_id = u.user_id; -- 隐式转换!
优化改写:显示统一数据类型
SELECT o.order_id, u.user_name
FROM orders o
JOIN users u ON CAST(o.user_id AS STRING) = u.user_id;
-- 或者
SELECT o.order_id, u.user_name
FROM orders o
JOIN users u ON o.user_id = CAST(u.user_id AS BIGINT);
12.5 优化手段二:修改配置参数
12.5.1 参数方案 1 : 开启 Group By 数据倾斜自动优化
参数设置
-- 核心参数 : 两阶段聚合
SET hive.groupby.skewindata = true;
-- 原理:
-- 第一阶段:随机将数据分发到不同的Reducer,进行局部聚合
-- 第二阶段:按正常 key 分发,进行全局聚合
-- 相当于自动完成了 "加盐 + 聚合 + 去盐 + 聚合" 的过程
SELECT city, COUNT(*), SUM(amount)
FROM orders
GROUP BY city;
流程示意图

12.5.2 参数方案 2 : 开启 Skew join 优化
参数设置及原理
-- 核心参数
SET hive.optimize.skewjoin = true;
-- 倾斜阈值:某个Key的数据量超过此值则认为是倾斜Key(默认100000=10万)
SET hive.skewjoin.key = 100000;
-- 原理:
-- Hive在运行时检测到倾斜Key后,自动将其拆分:
-- 1. 倾斜Key的数据走Map Join(如果另一张表够小)
-- 2. 非倾斜Key的数据走正常Common Join
-- 3. 最终UNION ALL合并结果
SELECT o.order_id, o.amount, u.user_name
FROM orders o
JOIN users u ON o.user_id = u.user_id;
流程示意图

12.5.3 参数方案 3:调整并行度缓解倾斜
-- Spark引擎下
SET spark.sql.shuffle.partitions = 200; -- 默认200,可增大
12.5.4 参数方案 4:Map Join 参数(小表倾斜场景)
-- 如果Join的一侧是小表,直接Map Join避免Shuffle
SET hive.auto.convert.join = true;
SET hive.mapjoin.smalltable.filesize = 50000000; -- 调大到50MB
-- 对于Left/Right Join中的倾斜
SET hive.auto.convert.join.noconditionaltask = true;
SET hive.auto.convert.join.noconditionaltask.size = 50000000;
12.5.5 参数方案 5:Spark AQE(Adaptive Query Execution)自动处理倾斜
-- Spark 3.0+ 的 AQE 倾斜处理(Hive on Spark 3.x 支持)
SET spark.sql.adaptive.enabled = true;
SET spark.sql.adaptive.skewJoin.enabled = true;
-- 倾斜因子:某分区数据量 > 中位数 × 此因子 则认为是倾斜
SET spark.sql.adaptive.skewJoin.skewedPartitionFactor = 5.0;
-- 倾斜阈值:数据量超过此值才处理(默认256MB)
SET spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 268435456;
-- 强制拆分的最小分区大小
SET spark.sql.adaptive.forceOptimizeSkewedJoin = true;
12.6 数据倾斜优化总结表
| 倾斜场景 | SQL改写方案 | 参数方案 |
|---|---|---|
| Group By 倾斜 | 加盐两阶段聚合 | hive.groupby.skewindata=true |
| Join 倾斜(已知Key) | 分离倾斜Key + UNION ALL | hive.optimize.skewjoin=true |
| Join 倾斜(未知Key) | 加盐+膨胀 | hive.skewjoin.key=100000 |
| NULL 值倾斜 | COALESCE + RAND() | 过滤NULL后单独处理 |
| COUNT DISTINCT | GROUP BY 改写 | hive.map.aggr=true |
| 类型不一致 | 显式CAST | — |
| 通用 | — | AQE: spark.sql.adaptive.skewJoin.enabled=true |
12.7 实战综合案例
业务场景: 电商平台,订单表(10亿条)JOIN 用户表(1亿条),其中 user_id=-1(未登录用户)有 2 亿条订单
-- ==================== 优化前(运行3小时+,最终OOM) ====================
SELECT
o.order_id,
o.amount,
o.order_time,
u.user_name,
u.city,
u.vip_level
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.dt = '2026-07-25';
-- ==================== 优化后(运行15分钟) ====================
-- 参数设置
SET hive.auto.convert.join = true; -- 开启 Map Join 自动转化
SET hive.mapjoin.smalltable.filesize = 50000000; -- 小表数据阈值
SET hive.optimize.skewjoin = true; -- 开启 skew Join
SET hive.skewjoin.key = 1000000; -- 倾斜 Key 阈值
SET spark.sql.shuffle.partitions = 500; -- 并行度大小
SET spark.executor.memory = 14g; -- 显示指定 Executor 内存大小
-- SQL改写:分离倾斜Key
-- Part 1: 正常用户(非倾斜)
SELECT
o.order_id,
o.amount,
o.order_time,
u.user_name,
u.city,
u.vip_level
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.dt = '2026-07-25'
AND o.user_id != -1
UNION ALL
-- Part 2: 未登录用户(倾斜Key),users表中user_id=-1只有1条,用MapJoin
SELECT /*+ MAPJOIN(u) */
o.order_id,
o.amount,
o.order_time,
u.user_name,
u.city,
u.vip_level
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.dt = '2026-07-25'
AND o.user_id = -1;
13. AQE(Adaptive Query Execution)专题
13.1 AQE的作用
CBO 的根本缺陷
CBO 在编译时做决策 → 依赖 Metastore 中的统计信息 → 统计信息可能不准
典型翻车场景:
CBO 估算:表 A 过滤后有 500万行 → 选择 Shuffle Join
实际运行:过滤后只剩 3MB → 本应 Broadcast Join,白白多了一次 Shuffle
CBO 估算:GROUP BY 后产生 200 个分区刚好
实际运行:过滤后数据极少,200 个分区每个只有 1KB → 大量空 Task 浪费调度
AQE 核心思想:
**不看"估计",看"实际"。**每次跑完一个 Shuffle Map Stage,拿到真实数据量,在决定下一步怎么跑
CBO:编译时决策(一次性)
SQL -> 统计信息 -> 生成计划 -> 执行(计划不变)
AQE :运行时决策(根据阶段调整)
SQL -> 生成初始计划 -> 执行 Stage 1 -> 收集真实统计,动态调整 Stage 2 的计划 -> 执行 Stage 2 -> ...
13.2 AQE 的三大核心能力
13.2.1 动态合并 Shuffle 分区
问题: spark.sql.shuffle.partitions 是固定值(默认200),数据少时产生大量空 Task。
-- 过滤后数据极少,但 Shuffle 仍产生 200 个分区
SELECT city, COUNT(*) FROM orders WHERE dt = '2026-07-27' GROUP BY city;
-- 实际数据只有 50MB → 200 个分区 → 每个分区 256KB → 199 个 Task 几乎空跑
AQE 做法:
执行前:200 个 Shuffle 分区(固定)
Map Stage 执行完毕 → 发现实际数据只有 50MB
↓
AQE 自动合并:200 个分区 → 合并为 1 个分区(目标 64MB/分区)
↓
Reduce Stage:只启动 1 个 Task → 节省 199 次调度开销
13.2.2 动态切换 JOIN 策略
问题:CBO 基于统计信息决定 JOIN 策略,但统计可能不准。
SELECT a.*, b.product_name
FROM fact_orders a
JOIN dim_product b ON a.product_id = b.product_id
WHERE b.category = 'electronics'; -- 过滤后 b 极小,但 CBO 不知道
AQE 做法:
初始计划:Sort-Merge Join(两侧都 Shuffle)
Map Stage 执行完毕 → 发现 b 过滤后实际只有 3MB
↓
AQE 动态切换:Sort-Merge Join → Broadcast Hash Join
↓
b 直接广播,a 不需要 Shuffle → 省掉一次全量 Shuffle
13.2.3 自动处理数据倾斜
问题:某个 Key 的数据量远超平均值,单个 Task 被撑爆。
SELECT user_id, SUM(amount)
FROM fact_orders
GROUP BY user_id;
-- user_id = 'VIP_001' 有 5亿条,其他 user 平均 100 条
-- 倾斜分区 50GB vs 正常分区 10MB → 单 Task OOM 或极慢
AQE 做法:
Map Stage 执行完毕 → 检测到各分区大小:
Partition 0: 12MB
Partition 1: 8MB
Partition 42: 50GB ← 倾斜!
Partition 43: 15MB
↓
AQE 自动将 Partition 42 拆分为多个子分区:
Partition 42-a: 10GB
Partition 42-b: 10GB
Partition 42-c: 10GB
Partition 42-d: 10GB
Partition 42-e: 10GB
↓
5 个 Task 并行处理 → 耗时降为原来的 1/5
13.3 核心配置参数
13.3.1 总开关
SET spark.sql.adaptive.enabled = true;
-- Spark 3.0~3.1 默认 false,Spark 3.2+ 默认 true
-- Hive on Spark 中建议显式开启
13.3.2 分区合并
-- 开启分区合并
SET spark.sql.adaptive.coalescePartitions.enabled = true;
-- 合并后的目标分区大小(最核心参数)
SET spark.sql.adaptive.advisoryPartitionSizeInBytes = 128m;
-- 默认 64MB,建议 64MB~256MB
-- AQE 按此大小计算需要多少个分区
-- 合并后最小分区数量
SET spark.sql.adaptive.coalescePartitions.minPartitionSize = 16MB;
-- 初始 Shuffle 分区数(AQE 会在此基础上合并)
SET spark.sql.shuffle.partitions = 2000;
-- 建议设大一些(如 2000),让 AQE 自动合并
-- 设小了 AQE 无法"拆分",只能合并
分区合并生效场景
-- 数据量波动大的日常查询
-- 工作日数据 50GB,节假日数据 500MB
SELECT city, SUM(amount) AS total
FROM dwd_orders
WHERE dt = '2026-07-27' -- 周日,数据极少
GROUP BY city;
-- 无 AQE:shuffle.partitions=200 → 200 个 Task,199 个空跑
-- 有 AQE:检测到实际 500MB → 合并为 4 个分区(500MB/128MB)→ 4 个 Task
13.3.3 JOIN 策略切换
-- 开启动态 JOIN 切换(默认跟随总开关)
SET spark.sql.autoBroadcastJoinThreshold = 26214400; -- 25MB
-- 运行时发现某侧 < 此值 → 自动转为 Broadcast Join
-- 设为 -1 则禁用此功能
动态 JOIN 切换场景
-- CBO 统计信息过期,误判表大小
SELECT o.order_id, o.amount, u.user_name
FROM fact_orders o
JOIN dim_users u ON o.user_id = u.user_id
WHERE u.vip_level = 5; -- 实际过滤后 u 只剩 2MB
-- CBO(统计过期):认为 u 有 200W 行 → Sort-Merge Join → 两侧 Shuffle
-- AQE(运行时):Map 跑完发现 u 只有 2MB → 自动切 Broadcast Join
-- o 侧无需 Shuffle → 性能提升数倍
13.3.4 倾斜处理
-- 开启倾斜 JOIN 优化
SET spark.sql.adaptive.skewJoin.enabled = true;
-- 倾斜判定:分区大小 > 中位数 × 此倍数 → 认为倾斜
SET spark.sql.adaptive.skewJoin.skewedPartitionFactor = 5;
-- 倾斜判定:分区大小 > 此绝对值 → 认为倾斜
SET spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 268435456;
-- 两个条件同时满足才判定为倾斜
-- 即:size > median × 5 AND size > 256MB
倾斜处理场景
-- 大卖家订单倾斜
SELECT
s.shop_name,
COUNT(o.order_id) AS order_cnt,
SUM(o.amount) AS gmv
FROM dim_shop s
JOIN fact_orders o ON s.shop_id = o.shop_id
WHERE o.dt = '2026-07-27'
GROUP BY s.shop_name;
-- shop_id = 'TOP_SELLER_001' 占全天订单的 30%
-- 无 AQE:该分区 15GB,单 Task 处理 20 分钟,其余 Task 1 分钟完成
-- 有 AQE:检测到倾斜 → 拆分为 60 个子分区(每个 256MB)→ 并行处理 → 2 分钟
13.4 AQE 和 CBO 的关系

14. 序列化与内存调优
14.1 序列化

14.2 存储层序列化(SerDe / 文件格式)
| 格式 | 类型 | 核心优势 | 适用 |
|---|---|---|---|
| TextFile | 行式 | 无 | ❌ 避免使用 |
| ORC | 列式 | 列裁剪 + 谓词下推 + Stripe 跳过 | ⭐ Hive 首选 |
| Parquet | 列式 | 嵌套类型好、Spark 原生优化 | ⭐ Spark 首选 |
| Avro | 行式 | Schema 演进 | 接口数据 |
SQL 示例:
-- 建表时指定 ORC + SNAPPY 压缩
CREATE TABLE dwd_orders (
order_id BIGINT, user_id BIGINT, amount DOUBLE, city STRING
)
PARTITIONED BY (dt STRING)
STORED AS ORC
TBLPROPERTIES ("orc.compress" = "SNAPPY");
同样查询
SELECT city WHERE amount > 100,ORC 比 TextFile 快 10~15 倍(只读相关列 + 跳过无关 Stripe)。
14.3 网络层序列化(Spark Serializer)
决定 Shuffle / 广播时数据怎么传输
| 序列化器 | 体积 | 速度 | 建议 |
|---|---|---|---|
| Java(默认) | 大(含类名字段名) | 慢 | ❌ |
| Kryo | 小 5~10 倍 | 快 5~10 倍 | ✅ 必须切换 |
-- 切换为 Kryo(最重要的一个配置)
SET spark.serializer = org.apache.spark.serializer.KryoSerializer;
-- Shuffle / 广播压缩
SET spark.shuffle.compress = true;
SET spark.broadcast.compress = true;
SET spark.io.compression.codec = lz4;
14.4 内存模型

Hive on Spark 很少用 cache → 把 Storage 比例调小,Execution 调大。
14.5 内存核心配置参数
14.5.1 内存
SET spark.executor.memory = 14g; -- 堆内存(4~16G)
SET spark.executor.memoryOverhead = 2048; -- 堆外(防 YARN 杀容器)
SET spark.executor.cores = 4; -- 并行 Task 数
SET spark.executor.instances = 30; -- Executor 数量
SET spark.driver.memory = 10g; -- Driver(广播收集、结果汇总)
-- 内存比例(Hive 场景推荐)
SET spark.memory.fraction = 0.7; -- Spark 内存占比(默认 0.6)
SET spark.memory.storage.fraction = 0.3; -- Storage 占比(默认 0.5,调小)
14.5.2 Shuffle 分区数(最常用调优手段)
SET spark.sql.shuffle.partitions = 1000; -- 默认 200,大查询必须调大
核心公式:每个 Task 数据量 = 总 Shuffle 数据 / partitions。
单个 Task 数据量控制在 128MB~512MB 为宜。
14.5.3 序列化 + 压缩
SET spark.serializer = org.apache.spark.serializer.KryoSerializer;
SET spark.shuffle.compress = true;
SET spark.io.compression.codec = lz4;
SET hive.exec.compress.intermediate = true; -- Hive 中间结果压缩
14.5.4 GC
-- JVM GC 参数(通过 Spark 传递)
SET spark.executor.extraJavaOptions = -XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:G1HeapRegionSize=16m -XX:InitiatingHeapOccupancyPercent=35;
-- 各参数含义:
-- -XX:+UseG1GC 使用 G1 垃圾回收器(推荐)
-- -XX:MaxGCPauseMillis=200 目标最大 GC 停顿 200ms
-- -XX:G1HeapRegionSize=16m G1 Region 大小
-- -XX:InitiatingHeapOccupancyPercent=35 堆占用 35% 时触发并发标记
-- Driver 的 JVM 参数
SET spark.driver.extraJavaOptions = -XX:+UseG1GC -XX:MaxGCPauseMillis=200;
14.5.5 AQE(Spark 3.0+,自动处理倾斜)
SET spark.sql.adaptive.enabled = true;
SET spark.sql.adaptive.skewJoin.enabled = true;
14.6 内存常见问题
| 报错 / 现象 | 原因 | 解法 |
|---|---|---|
OOM: Java heap space(Executor) | 单 Task 数据太多 | 调大 shuffle.partitions 或 executor.memory |
Container killed by YARN | 堆外内存不够 | 调大 memoryOverhead |
OOM(Driver) | 广播表太大 / 结果集太大 | 检查 MapJoin 表大小;加 LIMIT |
| GC Time 占比 > 30% | 内存紧张频繁 GC | 加内存 / 减数据量 / G1GC |
| 个别 Task 极慢(数据倾斜) | 某 Key 数据量远超平均 | AQE / 加盐打散 / hive.groupby.skewindata=true |
到这里,Hive on Spark 调优的完整链路已经闭合了:
第一篇(地基) 第二篇(常规优化) 第三篇(进阶优化)
━━━━━━━━━━━━━━━ ━━━━━━━━━━━━━━━━━━ ━━━━━━━━━━━━━━━━━━
YARN 资源规划 Explain 执行计划 数据倾斜 5 种场景
Spark 引擎部署 Map-side 预聚合 SQL 改写(加盐/分离/改写)
Executor 内存/Cores 分区裁剪 / 列裁剪 参数方案(skewjoin/AQE)
Hive on Spark 配置 谓词下推 AQE 三大能力
四种 Join 算法 序列化(Kryo/LZ4)
小文件治理 内存模型 & GC
并行度调参
CBO 成本优化
如果把这个系列浓缩成一句话:
调优的本质,就是在数据流动的每一个环节,让"该少的少、该快的快、该并的并、该拆的拆"。
- 该少的少:列裁剪、谓词下推、Map-side 聚合、分区裁剪——减少数据量
- 该快的快:Kryo 序列化、LZ4 压缩、ORC 列存——减少传输和 I/O 开销
- 该并的并:并行度、AQE 动态合并、小文件合并——减少调度浪费
- 该拆的拆:数据倾斜加盐、AQE 拆分倾斜分区——消除长尾
三篇 14 章,核心参数不超过 50 个。真正决定效果的,不是记住每一个参数,而是看到执行计划 / Spark UI 的那一刻,知道瓶颈在哪一段,然后去对应的章节找药方。
完整的配置模板,以下模板覆盖三篇所有调优参数,按 集群资源 → 内存模型 → 序列化压缩 → Hive SQL 优化 → 运行时并行 → AQE 自适应 → 数据倾斜 → GC 调优 → 统计信息维护 九层组织。每个参数均标注默认值、推荐值、作用说明和调整时机。
-- ================================================================
-- Hive on Spark ETL 调优模板
-- 适用:Hive 3.x + Spark 3.x + YARN
-- ================================================================
-- ======================== 一、资源 ========================
-- Executor 堆内存,建议 8g~16g,默认 1g
SET spark.executor.memory = 14g;
-- 堆外内存,防 YARN 杀容器,建议堆内10%~20%,默认 384m
SET spark.executor.memoryOverhead = 2048;
-- 每个 Executor 并行 Task 数,建议 4~5,默认 1
SET spark.executor.cores = 4;
-- Executor 总数 = (集群总核-1)/cores,建议按集群算,默认 2
SET spark.executor.instances = 30;
-- Driver 内存,广播收集/结果汇总,建议 4g~10g,默认 1g
SET spark.driver.memory = 10g;
-- Driver 最大结果集,建议 2g~4g,默认 1g
SET spark.driver.maxResultSize = 2g;
-- ======================== 二、内存模型 ========================
-- Spark 可用内存占堆比例,建议 0.7,默认 0.6
SET spark.memory.fraction = 0.7;
-- Storage 占比,Hive 少用 cache 调小,建议 0.3,默认 0.5
SET spark.memory.storage.fraction = 0.3;
-- ======================== 三、序列化 & 压缩 ========================
-- Kryo 序列化,Shuffle 体积减 5~10 倍,建议必须切,默认 JavaSerializer
SET spark.serializer = org.apache.spark.serializer.KryoSerializer;
-- Shuffle 输出压缩,建议 true,默认 true
SET spark.shuffle.compress = true;
-- 广播变量压缩,建议 true,默认 true
SET spark.broadcast.compress = true;
-- 压缩算法,lz4 最快 / zstd 最小,建议 lz4,默认 lz4
SET spark.io.compression.codec = lz4;
-- Hive 中间结果压缩,建议 true,默认 false
SET hive.exec.compress.intermediate = true;
-- ======================== 四、Map-side 预聚合 ========================
-- Map 端局部聚合,减少 Shuffle 量,建议 true,默认 true
SET hive.map.aggr = true;
-- 聚合表最多占 Map 内存比例,建议 0.5,默认 0.5
SET hive.map.aggr.hash.percentmemory = 0.5;
-- 聚合收益低于此比例自动放弃,建议 0.5,默认 0.5
SET hive.map.aggr.hash.min.reduction = 0.5;
-- 内存达此比例强制刷盘防 OOM,建议 0.9,默认 0.9
SET hive.map.aggr.hash.force.flush.memory.threshold = 0.9;
-- ======================== 五、谓词下推 ========================
-- WHERE 下推到存储层(含 ORC 行级过滤),建议 true,默认 true
SET hive.optimize.ppd = true;
-- ======================== 六、Join 优化 ========================
-- 小表自动转 Map Join,建议 true,默认 true
SET hive.auto.convert.join = true;
-- 小表判定阈值,建议 25MB~100MB,默认 25MB
SET hive.mapjoin.smalltable.filesize = 25000000;
-- 编译期直接决定 Map Join,建议 true,默认 true
SET hive.auto.convert.join.noconditionaltask = true;
-- 多表 Join 小表总大小上限,建议与上面一致,默认 10MB
SET hive.auto.convert.join.noconditionaltask.size = 25000000;
-- Bucket Join,需分桶表,建议 true,默认 false
SET hive.optimize.bucketmapjoin = true;
-- SMB Join,需排序分桶表,建议 true,默认 false
SET hive.optimize.bucketmapjoin.sortedmerge = true;
-- ======================== 七、CBO 成本优化 ========================
-- CBO 总开关,建议 true,默认 true
SET hive.cbo.enable = true;
-- 使用表级统计,建议 true,默认 false
SET hive.compute.query.using.stats = true;
-- 使用列级统计,建议 true,默认 false
SET hive.stats.fetch.column.stats = true;
-- 使用分区级统计,建议 true,默认 true
SET hive.stats.fetch.partition.stats = true;
-- ======================== 八、小文件治理 ========================
-- 输出时合并小文件,建议 true,默认 true
SET hive.merge.sparkfiles = true;
-- 合并后目标文件大小,建议 256MB,默认 256MB
SET hive.merge.size.per.task = 268435456;
-- 平均文件小于此值才触发合并,建议 64MB,默认 16MB
SET hive.merge.smallfiles.avgsize = 67108864;
-- 读表时合并小文件到同一 Task,建议保持,默认 CombineHiveInputFormat
SET hive.input.format = org.apache.hadoop.hive.ql.io.CombineHiveInputFormat;
-- ======================== 九、并行度 ========================
-- Shuffle 分区数(最核心),开 AQE 设大让其自动收,建议 1000,默认 200
SET spark.sql.shuffle.partitions = 1000;
-- 读表单 partition 上限,建议 256MB,默认 128MB
SET spark.sql.files.maxPartitionBytes = 268435456;
-- 小文件打开虚拟代价,建议 4MB,默认 4MB
SET spark.sql.files.openCostInBytes = 4194304;
-- ======================== 十、AQE 自适应执行 ========================
-- AQE 总开关,建议 true,默认 false(3.0~3.1)/true(3.2+)
SET spark.sql.adaptive.enabled = true;
-- 自动合并过小 Shuffle 分区,建议 true,默认 true
SET spark.sql.adaptive.coalescePartitions.enabled = true;
-- 合并目标分区大小,建议 128MB,默认 64MB
SET spark.sql.adaptive.advisoryPartitionSizeInBytes = 134217728;
-- 合并后最小分区,建议 16MB,默认 1MB
SET spark.sql.adaptive.coalescePartitions.minPartitionSize = 16777216;
-- 运行时某侧小于此值自动转 Broadcast,建议 25MB,默认 10MB
SET spark.sql.autoBroadcastJoinThreshold = 26214400;
-- 自动拆分倾斜分区,建议 true,默认 true
SET spark.sql.adaptive.skewJoin.enabled = true;
-- 分区 > 中位数 × 此倍数判定倾斜,建议 5.0,默认 5.0
SET spark.sql.adaptive.skewJoin.skewedPartitionFactor = 5.0;
-- 且 > 此绝对值才处理,建议 256MB,默认 256MB
SET spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 268435456;
-- 强制拆分,不因代价高跳过,建议 true,默认 false
SET spark.sql.adaptive.forceOptimizeSkewedJoin = true;
-- ======================== 十一、数据倾斜兜底 ========================
-- Group By 两阶段聚合,开 AQE 时关闭,建议 false,默认 false
SET hive.groupby.skewindata = false;
-- Hive 层 Skew Join,倾斜 Key 走 Map Join,建议 true,默认 false
SET hive.optimize.skewjoin = true;
-- 某 Key 行数超此值判定倾斜,建议 100000,默认 100000
SET hive.skewjoin.key = 100000;
-- ======================== 十二、GC ========================
-- G1GC/停顿200ms/Region16m/堆占35%触发并发标记,建议按此配,默认无
SET spark.executor.extraJavaOptions = -XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:G1HeapRegionSize=16m -XX:InitiatingHeapOccupancyPercent=35;
-- Driver 端 G1GC,建议按此配,默认无
SET spark.driver.extraJavaOptions = -XX:+UseG1GC -XX:MaxGCPauseMillis=200;
-- ======================== 十三、统计信息维护(独立调度,每日执行) ========================
-- 表级统计:行数/文件数/总大小
-- ANALYZE TABLE dwd_orders PARTITION(dt='${bizdate}') COMPUTE STATISTICS;
-- ANALYZE TABLE dim_users COMPUTE STATISTICS;
-- 列级统计:NDV/最大最小值/NULL比例
-- ANALYZE TABLE dwd_orders PARTITION(dt='${bizdate}') COMPUTE STATISTICS FOR COLUMNS;
-- ANALYZE TABLE dim_users COMPUTE STATISTICS FOR COLUMNS;
更多推荐
所有评论(0)