posts

Aug 07, 2026

💥从17.37GB到能跑完:一条 Impala 查询的三次手术

Estimated Reading Time: 5 minutes (1041 words)

同事扔过来一条 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,引用几次就算几次。optfeemoney 各用一次,整棵子树就老老实实算两遍。

各家不一样,顺便一提。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 正好是这三列。所以:

反过来,如果图省事让 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=unavailablecardinality=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(没统计,它连数字都编不出来);predicateskudu predicates 差一个词,差一个数量级。这些都得对着 SQL 反推。计划本身不会提醒你。

日期维表是标准做法,但下游本来就要跟事实表 inner join 时,用事实表里真实存在的组合当驱动,往往能同时换来更好的 join 算法和一次免费剪枝。前提是等价性推清楚——剪枝键不能超出分区键范围。这条边界含糊不得。

还有一点我之前没怎么意识到:读慢查询顺便就是一次 code review。这次捞出来的两个日期问题,比性能问题更值钱。跑得慢大家都知道要修;算错的查询可能已经悄悄错了几个月。