前言
上一篇我们走完了 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 秒速查:数据倾斜五步法

  1. 看 Spark UI → 哪个 Stage 有长尾 Task
  2. 查 Key 分布 → 确认倾斜 Key 是谁
  3. 判断场景 → Group By / Join / NULL / DISTINCT / 类型
  4. 选方案 → SQL 改写(根治) or 参数(兜底)
  5. 验证 → 对比 Task 耗时分布是否均匀

12.2 数据倾斜的典型场景

场景触发操作示例
Group By 倾斜GROUP BY某个 province = ‘广东省’ 占了 80% 的数据
Join 倾斜JOIN ON某个 user_id = 0(默认值)关联了百万条
Count Distinct 倾斜COUNT(DISTINCT)全局去重,所有的数据到一个 Task
空值/默认值JOIN ONNULL 值全部哈希到同一个分区
数据类型不一致JOIN ONINT 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_0null_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 ALLhive.optimize.skewjoin=true
Join 倾斜(未知Key)加盐+膨胀hive.skewjoin.key=100000
NULL 值倾斜COALESCE + RAND()过滤NULL后单独处理
COUNT DISTINCTGROUP 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.partitionsexecutor.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;

更多推荐