同事扔过来一条 SQL,说跑不出来。打开一看:六个 CTE 套着,中间塞了 900 个促销编号的 IN,光 EXPLAIN 就 42 个节点,头两行写着:
Max Per-Host Resource Reservation: Memory=367.88MB Threads=25
Per-Host Resource Estimates: Memory=17.37GB
单节点估 17.37GB。这不是慢,是根本跑不完。
下面是我改它的过程,动了三刀。表名脱敏了;逻辑和计划里的数字是真的。
这条 SQL 在干什么
业务大概是这样:一批促销(下面叫 plan),每个覆盖一段日期、一个渠道、一批商品,带一个返点率。同一天、同渠道、同商品可能被好几个 plan 同时罩住——这时候取返点率最高的那个。再用这个率去乘销售明细算钱,扣掉别处已经分摊过的费用,按月汇总。
脱敏后的骨架:
WITH rbt AS ( -- 促销清单:plan × 商品 × 渠道,带生效区间
SELECT p.plan_code, pc.channel_id, s.item_id,
greatest(s.begin_date, p.start_time) AS begin_date,
date_sub(least(s.end_date, p.create_time), 3) AS end_date,
s.rebate_rate, p.plan_id, p.vendor_id
FROM tmp.item_plan_span s
JOIN t_plan p ON p.plan_id = s.plan_id AND p.method = 'X'
JOIN t_plan_channel pc ON pc.plan_code = p.plan_code AND pc.seq = 0
WHERE p.plan_code IN ('P0001', /* ...900 个... */ 'P0900')
),
expanded AS ( -- 按天展开
SELECT d.period_date, rbt.*,
row_number() OVER (
PARTITION BY d.period_date, rbt.channel_id, rbt.item_id
ORDER BY rbt.rebate_rate DESC, rbt.plan_code DESC) AS rn
FROM rbt
JOIN dim.dim_date d
ON to_date(d.period_date) BETWEEN rbt.begin_date AND rbt.end_date
),
opt AS (SELECT * FROM expanded WHERE rn = 1), -- 取返点最高的
fee AS ( -- 别处已分摊的费用
SELECT f.vendor_id, f.channel_id, f.item_id, f.stat_date, sum(f.fee_amount) fee
FROM t_extra_fee f
JOIN opt ON opt.vendor_id = f.vendor_id
AND opt.channel_id = f.channel_id
AND opt.item_id = f.item_id
AND opt.period_date = f.stat_date
GROUP BY 1,2,3,4
),
money AS (
SELECT from_timestamp(opt.period_date,'yyyy-MM'), opt.channel_id,
sum(greatest(d.price * d.qty * opt.rebate_rate - ifnull(fee.fee,0), 0))
FROM kudu.t_sales_detail d
JOIN opt ON ...四个等值条件...
LEFT JOIN fee ON ...四个等值条件...
GROUP BY 1,2
)
SELECT * FROM money;
看着挺规矩。坑全在执行计划里。
第一刀:那棵子树被算了两遍
42 个节点我从头翻到尾。翻到一半愣了一下——这段好像刚看过。往上对,01–10 和 12–21 几乎是镜像:
01:SCAN HDFS [tmp.item_plan_span] size=1.51GB 12:SCAN HDFS [tmp.item_plan_span] size=1.51GB
02:SCAN KUDU [t_plan] 13:SCAN KUDU [t_plan]
03:SCAN KUDU [t_plan_channel] 14:SCAN KUDU [t_plan_channel]
07:NESTED LOOP JOIN 18:NESTED LOOP JOIN
08:SORT 19:SORT
09:ANALYTIC row_number() 20:ANALYTIC row_number()
1.51GB 扫两遍,三张 Kudu 各扫两遍,嵌套循环两遍,全量排序两遍。
原因很土:Impala 的 CTE 不物化。WITH 就是语法糖,展开成 inline view,引用几次就算几次。opt 被 fee 和 money 各用一次,整棵子树就老老实实算两遍。
各家不一样,顺便一提。PostgreSQL 12 之前 CTE 强制物化,之后默认内联但可以 MATERIALIZED;Spark 会做 CTE 复用;Impala 是老实展开派。跨引擎搬 SQL 很容易踩。
能不能少引用一次?回去读逻辑,发现 fee 里那个 JOIN opt 根本多余:
FROM t_extra_fee f
JOIN opt ON opt.vendor_id = f.vendor_id AND opt.channel_id = f.channel_id
AND opt.item_id = f.item_id AND opt.period_date = f.stat_date
GROUP BY f.vendor_id, f.channel_id, f.item_id, f.stat_date
join 的四个 key 和自己 GROUP BY 的四列一模一样,粒度不变,只是过滤。而 fee 最后只被 money 用同样四个 key 做 LEFT JOIN——这边滤掉的行,那边本来也匹配不上。
删掉之后:
fee AS (
SELECT vendor_id, channel_id, item_id, stat_date, sum(fee_amount) fee
FROM t_extra_fee
WHERE stat_date >= '2026-03-01' AND stat_date < '2026-08-08'
GROUP BY 1,2,3,4
)
删了七个字。节点 42 → 26,内存估算 17.37GB → 9.94GB,线程 25 → 16。
性价比高得有点荒唐。但前提是你得先看见子树在重复——计划里唯一的信号就是「节点编号多了一倍」,EXPLAIN 不会写「注意,这里算了两遍」。所以我现在看长计划,习惯先数一眼节点数跟 SQL 复杂度对不对得上。对不上,就去找对称结构。
第二刀:NESTED LOOP JOIN 骗了我一次
9.94GB 还是跑不动。最扎眼的是这个:
07:NESTED LOOP JOIN [INNER JOIN, BROADCAST]
| predicates: greatest(s.begin_date, p.start_time) <= to_date(d.period_date),
| date_sub(least(s.end_date, p.create_time), 3) >= to_date(d.period_date)
|
|--18:EXCHANGE [BROADCAST]
| 06:SCAN HDFS [dim.dim_date d] size=3.97MB
嵌套循环。条件是 period_date BETWEEN begin_date AND end_date,纯范围,没有等值列做 hash,Impala 只能逐行配对。
第一反应是把日期维表切小。原来没限制,我加了个范围:
JOIN dim.dim_date d
ON d.date_type = 'D'
AND d.period_date BETWEEN '2026-03-01' AND '2026-08-08'
AND to_date(d.period_date) BETWEEN rbt.begin_date AND rbt.end_date
右表从整张维表压到大约 160 行。看着舒服了。一算,砍在棉花上。
嵌套循环贵的不是右表大小,是输出行数。
rbt 有多少行?临时表 1.51GB,row-size=46B,粗估千万级。每行跟 160 天配对,几十亿次比较——光比较其实还好,摊到十几个节点也就几秒。
要命的是配对成功之后:每行 rbt 会展开成「生效区间有多少天」那么多行。区间平均 30 天,千万行进去、几亿行出来。这几亿行还得接着走:
19:EXCHANGE [HASH(period_date, channel_id, item_id)] ← 几亿行过网络重分区
08:SORT ← 几亿行全量排序
09:ANALYTIC row_number()
10:SELECT predicates: row_number() = 1 ← 排完了扔掉 99%
先按三个键 shuffle,再全量排序,排完为了 rn = 1 把绝大部分扔掉。9.94GB 主要吃在这儿。
还有个隐藏成本:predicates 里的 greatest / least / date_sub / to_date。写在 join 条件里,等于每个配对重算一次,不是每行一次。几十亿次函数调用。我当时扫计划时差点漏过这点。
换个思路:日期不该从维表来
这几亿行里到底有多少有用?展开出来的每一行是「某个 plan 在某一天覆盖某商品某渠道」,但最后算钱要跟销售明细 inner join——那天那个商品那个渠道如果根本没销售,这行一定被丢掉。
「哪些天哪个商品哪个渠道有销售」,明细表自己知道。
所以把日期维表整段换掉,用明细里真实存在的组合当日期源:
cal AS (
SELECT DISTINCT to_date(d.sale_time) AS period_date, d.item_id, d.channel_id
FROM kudu.t_sales_detail d
WHERE d.sale_time >= '2026-03-01' AND d.sale_time < '2026-08-08'
),
expanded AS (
SELECT cal.period_date, rbt.*,
row_number() OVER (
PARTITION BY cal.period_date, rbt.channel_id, rbt.item_id
ORDER BY rbt.rebate_rate DESC, rbt.plan_code DESC) AS rn
FROM rbt
JOIN cal ON cal.item_id = rbt.item_id
AND cal.channel_id = rbt.channel_id
AND cal.period_date BETWEEN rbt.begin_date AND rbt.end_date
)
join 从「一个纯范围」变成「两个等值 + 一个范围」,NESTED LOOP JOIN 直接变回 HASH JOIN。
算法换回来只是附带的。真正值钱的是剪枝:每行 rbt 不再展开成「区间有多少天」,而是「这个商品这个渠道在区间内实际有几天出货」。前者可能是 30,后者常常是 3。进 SORT 的行数掉一个数量级。
那些函数也从 join 条件挪出去了——现在在 rbt 里每行算一次。
这个替换是等价的吗
这是全文最需要小心的地方。我当时来回推了三遍,有点强迫症:候选集剪小了,rn = 1 会不会选出不同的 plan?
关键在剪的是哪个维度。cal 剪掉的是整个 (日期, 渠道, 商品) 组合,而 row_number() 的 PARTITION BY 正好是这三列。所以:
- 被剪掉的组合,明细里本来就没有对应记录,最终 inner join 也会丢掉,剪不剪结果一样;
- 留下来的组合,分区内所有 plan 候选一个都没少(
cal的三个 key 不含 plan),rn = 1还是同一个。
反过来,如果图省事让 rbt 直接 join 明细、再在明细上做 row_number(),那就错了。因为 join 会带上 plan_id,候选集会缩成「那天有明细的那些 plan」。假设 plan A 返点最高但那天没出货、plan B 返点低但有出货:原逻辑选 A,A 匹配不上明细 → 不产生金额;错误写法会选 B → 凭空多一笔。
差别就在于:剪枝的键里能不能出现分区键之外的列。能出现,就不安全。这条我后来跟同事复述了两遍,生怕自己当时推漏了。
代价:明细表要扫两遍
cal 和最终算钱都读 t_sales_detail。Impala 不物化 CTE,这张 Kudu 表会被扫两遍。
最干净是落一张中间表,扫一遍后面复用。这次不许建表——临时表在这套环境里另有麻烦——只能接受扫两遍。
算下来还是赚:第一遍只读三列、带 sale_time 范围谓词,Kudu 是列存,读三列比读全表便宜得多;换掉的是几亿行的 shuffle + 全量排序。这笔我签。
第三刀:一字之差的 kudu predicates
改到这儿我盯着 Kudu 扫描看了半天,发现两件事。
00:SCAN KUDU [kudu.t_sales_detail d]
runtime filters: RF004 -> d.plan_id, RF005 -> d.channel_id, RF006 -> d.item_id
整个计划里唯一一个连 predicates 都没有的扫描。销售明细在裸扫全表,只靠 runtime filter 兜底。原因是 join 写的是 to_date(d.sale_time) = opt.period_date——函数套在列上,谓词下不去。补一条裸列静态范围就行:
WHERE d.sale_time >= '2026-03-01' AND d.sale_time < '2026-08-08'
另一个更阴。费用表这个节点:
11:SCAN KUDU [t_extra_fee f]
predicates: stat_date <= '2026-06-30', stat_date >= '2026-03-01'
我以为谓词生效了。对比同一份计划里另外两个节点:
03:SCAN KUDU [t_plan_channel]
kudu predicates: seq = 0, plan_code IN ('P0001', ...)
一个写 predicates,一个写 kudu predicates。差的是「Kudu 在存储层过滤」和「Kudu 全量吐给 Impala,Impala 再过滤」。前者省 IO 和网络,后者只省了后面算子的活。
费用表那个没下推,大概是 stat_date 列类型和字符串字面量对不上,Impala 放弃了。给字面量显式定型就好了:
WHERE f.stat_date >= CAST('2026-03-01' AS timestamp)
AND f.stat_date < CAST('2026-08-08' AS timestamp)
EXPLAIN 里这两个词长得太像,扫一眼很容易当成一回事,IO 量却差着数量级。我差点就当成一回事。
没有统计信息,就得自己开手动挡
全过程里最拖后腿的其实是这个。计划头上一直挂着:
WARNING: The following tables are missing relevant table and/or column statistics.
满屏 cardinality=unavailable 和 cardinality=0。基数未知时 Impala 默认倾向 BROADCAST,于是每一个 join 都 broadcast——包括把千万行的 opt 结果广播到每个节点。
正常该 COMPUTE STATS。这次跑不动:表太大,权限也不够。只能手写 hint:
-- 这两张表被 900 个编号过滤后只剩几百行,几十 KB,广播到 20 个节点也就几百 KB
JOIN /* +BROADCAST */ t_plan p ON ...
JOIN /* +BROADCAST */ t_plan_channel pc ON ...
-- 这三处两边都是千万级,broadcast 等于每个节点都装一份完整 hash 表
JOIN /* +SHUFFLE */ cal ON ...
JOIN /* +SHUFFLE */ kudu.t_sales_detail ON ...
LEFT JOIN /* +SHUFFLE */ fee ON ...
判断标准一句话:右表小到「复制 N 份也无所谓」就 broadcast,否则 shuffle。broadcast 花的是「右表行数 × 节点数」的网络和「右表全量」的单机内存;shuffle 花的是「两边各传一次」的网络和「右表 1/N」的内存。
有个反直觉的点:小表 join 要保留 broadcast,别图省事全改成 shuffle。broadcast join 更容易生成有效的 runtime filter 推给左侧扫描——上面那个 RF010 -> s.plan_id 就是这么来的,让 1.51GB 那张表能提前跳掉大量行。没统计信息的时候,runtime filter 几乎是唯一还在替你干活的优化。别掐掉。
配套调了几个 query option:
SET RUNTIME_FILTER_WAIT_TIME_MS=60000; -- 默认 0,扫描不等 filter 到达就开跑
SET RUNTIME_FILTER_MODE=GLOBAL;
SET MT_DOP=8;
SET MEM_LIMIT=16g;
RUNTIME_FILTER_WAIT_TIME_MS 默认是 0:扫描不等 filter,先跑起来再说。既然现在全靠 filter 剪枝,让它等一会儿值得。
另外发现一个偏方:不用 COMPUTE STATS,只要有 ALTER TABLE 权限就能手写行数,不做任何扫描:
ALTER TABLE tmp.item_plan_span
SET TBLPROPERTIES('numRows'='35000000', 'STATS_GENERATED_VIA_STATS_TASK'='true');
自己 count(*) 一下填进去。哪怕只修最大那张表,也比满屏 cardinality=0 强,上面那堆 hint 能少写几个。
顺手捡到两个 bug
调优时读 SQL 比平时细,捞出两个跟性能无关的问题。说实话,这两个比前面几刀更让我后怕。
日期窗口截错了。第二版加日期范围时写的是 BETWEEN '2026-01-01' AND '2026-06-30',但促销编号覆盖的是 3 月初到 8 月初,临时表名带的日期是 8 月 7 号。这个窗口把 7、8 两月整段砍掉了,1、2 月大概率是空的——白扫两个月,又漏两个月。而且这种错不报错,只是数字小一点。对不上账够查半天。
<= 漏掉一整天。stat_date <= '2026-06-30',如果这列带时分秒(明细那张表的时间字段在 ORM 里是 Date,对应 MySQL datetime),6 月 30 号当天所有非零点记录会被整天漏掉。时间范围一律写 >= 下界 AND < 上界+1天,别用 <=。
还有一处不算 bug,但很可疑:同一个字段在两处 join 里,一处写 opt.period_date = f.stat_date,另一处写 opt.period_date = to_date(f.stat_date)。如果这列带时分秒,匹配结果不同,费用会被静默 null 掉。同一个字段两种写法——现在是我 review SQL 时的固定检查项。
如果排序还是撑不住
剩下唯一的内存尖峰是 row_number() 那次排序。真撑不住还有一招,够脏:把排序键编码成定长字符串,用 hash 聚合的 max() 替代排序,再拆回来。
-- 编码:返点率定长补零 + 分隔符 + 其他要带出来的字段
concat(lpad(CAST(CAST(round(rebate_rate * 1000000) AS bigint) AS string), 19, '0'),
'#', plan_code, '#', CAST(plan_id AS string), '#', CAST(vendor_id AS string)) AS k
-- 聚合取最大,然后拆
SELECT period_date, channel_id, item_id,
CAST(substr(k,1,19) AS decimal(28,6)) / 1000000 AS rebate_rate,
split_part(k,'#',2) AS plan_code,
CAST(split_part(k,'#',3) AS bigint) AS plan_id
FROM (SELECT period_date, channel_id, item_id, max(k) AS k
FROM expanded GROUP BY 1,2,3) t
SORT 换成 AGGREGATE,hash 聚合 spill 起来比全量排序体面得多。
代价是可读性。还有两个前提:返点率不能为负(补零编码对负数会排错),plan_code 必须定长(这里固定 15 位,满足)。我把它留在注释里当备选。前面几刀够用就不上。
复盘
回头看,三刀的顺序其实挺重要——我差点搞反。
先找结构问题,再抠算子。第一刀删七个字省掉一半工作量,这种收益参数调优换不来。如果一上手就去琢磨 MEM_LIMIT 该给多少、MT_DOP 调几,我会在错误基线上浪费很多时间。差点就是那样干的。
EXPLAIN 里最贵的东西往往不写在字面上。「子树算了两遍」只体现为节点编号多一倍;「嵌套循环输出几亿行」体现为一行 cardinality=0(没统计,它连数字都编不出来);predicates 和 kudu predicates 差一个词,差一个数量级。这些都得对着 SQL 反推。计划本身不会提醒你。
日期维表是标准做法,但下游本来就要跟事实表 inner join 时,用事实表里真实存在的组合当驱动,往往能同时换来更好的 join 算法和一次免费剪枝。前提是等价性推清楚——剪枝键不能超出分区键范围。这条边界含糊不得。
还有一点我之前没怎么意识到:读慢查询顺便就是一次 code review。这次捞出来的两个日期问题,比性能问题更值钱。跑得慢大家都知道要修;算错的查询可能已经悄悄错了几个月。