Oracle 2013 年的一篇论文《Adaptive and Big Data Scale Parallel Execution in Oracle 》里介绍了并行执行中使用到的大量优化技术。每个点都值得展开。本文介绍它的第4节,Window Function 并行优化。
Oracle 对以下三类 window function 做执行优化:
先看看常规解决方案,然后分析这些方案存在的不足,最后介绍 Oracle 解决方案。
下层数据读入窗口后进行聚合运算,将运算结果写到对应列上。
可以发现,以上操作都需要读入整个分区后才能开始计算。如果分区分布不均,例如某个 c2 值对应的分区有上亿行,那么会导致这个分区的计算非常慢。这就是通常说的 “分区倾斜”(Skew)。
进一步的,为了支持下面的 query,应该怎样读入分区呢?
SELECT /*y:year q:quarter m:month d:day*/
y, q, m, d, sales,
SUM(sales) OVER (PARTITION BY y,q,m) msales,
SUM(sales) OVER (PARTITION BY y,q) qsales,
SUM(sales) OVER (PARTITION BY y) ysales
FROM fact f;
一种方式是使用三个 Window 算子串联(如下图),但这种方式需要多次传递数据,效率较低。
Oracle 只使用一个 Window 算子,按照 y 做数据分发,如下图:
思考题:为什么不能用 (y,q) 或者 (y,q,m) 做 hash / range 数据分发?
答案:
Oracle RDBMS uses sort-based
execution of window functions and evaluates the three reporting
aggregates using a single “window sort” operation. Because
of this clumping, the parallel plan requires data redistribution
(hash or range) on the common PBY key. the data distribution is
done by “hash” on the common PBY key “year”.
并行计算步骤(算法1):
在这个算法里,我们只能按照 y 做分区,一旦 y 有倾斜,那么 (y,q),(y,q,m) 都会收到拖累。
所以,处理好 y 的倾斜问题非常重要。Oracle 的处理思路是:使用更多的 key 来做分区,比如 (y,q) ,然后把窗口函数的处理拆分成两步:Sort & Consolidator

为什么 能够 拆成 2 步呢?回到窗口函数的本源:
分类讨论:
先看 Reporting:上面过程要求每个分区对应的 Consolidator 是单线程,所有 Sort 的结果都需要发到 Consolidator 。假设分区有倾斜,那这里又会是单点。
为了解决这个问题,有两种做法:
基于上述算法 Reporting 还需考虑 NDV 错估的场景。当真实 NDV 非常大时,会分出非常多的子块,导致需要广播的数据量暴增,最终导致该算法的收益还不如最原始的 Window 并行算法(算法1)。
Report 相对简单一些,因为每个分区里 Consolidate 出来的数据是一个标量,直接 patch 到分区内的每一行上即可。对于 Cumulative 和 Ranking,Consolidate 出来的数据是向量,需要按照一定的规则 patch 到子块上,不同行对应向量里不同的结果值。Patch 算法大致思路是,Consolidate 算法的输出包括两列:value,index,value 是计算的结果, index 指示将结果写到哪个原始行里。
Cumulative 和 Ranking 算法本文不详述,可以看到 Oracle 论文 4.2 节。