暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

SCOPE 优化器如何生成并行执行计划

分布式系统之美 2021-04-13
534

Incorporating Partitioning and Parallel Plans into the SCOPE Optimizer周靖人 2010 年发表于 ICDE 的论文。论文描述了 SCOPE 优化器如何考虑并行相关信息,生成优秀的并行执行计划,优化分布式计算的性能。SCOPE 是用来做数据分析的系统,SCOPE 优化器会将用户输入的 SQL-Like 语句,转换成执行计划,并下发到集群的机器上执行。

我们先考虑下,让一个不考虑并行信息的优化器,能生成并行计划,应该怎么做?最简单的方式就是在生成计划后,插入一个事后步骤(post-processing)来考虑并行信息,把这个计划变成并行的。比如简单粗暴的在每两算子间插入数据(exchange)交换用的算子,然后把原本的算子都并行起来执行,利用插入的交换算子分发数据。但这样简单的做法很难得到优秀的计划,下面是一个例子:

计划 a 就是最简单的做法,在所有算子中插入数据交换算子。不同情况 a 和 b 各有优劣,取决于众多因素,如:数据量大小,R.a 和 S.a 的分布,Join 的选择率等。且计划 b 的正确性还需要一个条件来保证:R.c,S.d 分别函数依赖 R.a,S.a,因为这样才能使得 (R.c, S.d) 值相同的数据在一个分片中,以确保 HashAgg 的正确性。

可见得到优秀并行计划需要考虑的因素众多,仅靠一个简单的事后步骤是远不够的;更好的做法是教会优化器准确的“理解”并发概念,将这些因素全盘考虑。接下来就介绍作者怎么教会 SCOPE 优化器“理解”并行,能生成正确的执行计划。其中 2~6 小节是理论部分,对分布式环境下的数据性质做出定义,并推导出一些演化规则,为并行计划的正确性打下理论基础;7~9 利用利用优化器本身的结构,把上面的理论知识“教”给优化器。


SCOPE 优化器


先简单介绍 SCOPE 的优化流程,经典优化器结构,这里不再赘述,请结合注释理解:


数据交换


不同的数据交换方式,可以通过组合不同的分配和聚合得到。我们先对不同的分片和聚合方式做一个简单梳理。先是分片(partition)算子,他把输入一分为多,所有的分片算子都是 FIFO 的:

然后是聚合(merge)算子,他把多个分片合并成一个,总结如下:

这里 Sort-Concat Merge 主要是为了和 Range Partition 进行配合,可以比较低成本的保持数据原有的有序性。


数据性质


接下来对数据性质(Structural Properties)进行定义,例子当中的数据为:(1, 1), (1, 2), (2, 2), (2, 6), (3, 7), (4, 3):

对上面提到的 4 种数据性质做一个思考:

  • grouping 和 ordering 描述一批“挨着”的数据,而 partitioning 描述多份“分散”的数据。

  • grouping 和 ordering 是局部性质(local property),而 partitioning 是全局性质(global property)。

  • 可以把 non-ordered partitioning 当做全局的 grouping,把 ordered partitioning 当做全局的 ordering。

把局部和全局性质写在一起,用来完整的定义数据性质,如下:

分别用 g 和 o 来表示 grouping 和 ordering;左边部分表示全局性质如:
  • {C1}^g:表示按照 C1 做了 non-ordered partitioning。
  • {C1^o, C2^o}:表示按照 C1 和 C2 做了 ordered partitioning。
右边部分表示局部性质如:
  • {C3^o}:表示按照 C3 排序。
  • {{C1, C2}^g}:表示按照 C1, C2 聚合。
  • {{C1, C2}^g, C3^o}:先按照 C1, C2 聚合,然后在每个 group 内再按照 C3 排序。
需要注意的是,如有多个 local 性质,需要在前一个的基础上,满足下一个。如数据为 (1, 1, 10), (1, 1, 2), (2, 2, 1),满足上面 case 3 的情况有:
  • [(1, 1, 2), (1, 1, 10), (2, 2, 1)]
  • [(2, 2, 1), (1, 1, 2), (1, 1, 10)]
上面两种情况中,整体看 C3 不是有序的,但在各自的 group 内,C3 是有序的。下面是一个独立的例子:

这份数据的性质可以表示为:


演化规则


接下来,在我们上面的形式化定义上,推导一些转换规则;原文一共有 8 条规则,这里我们选几个简单的当做例子看一看。

这个规则比较简单,按照上小节说明:需要在上一个的基础上满足下一个;因此局部性质是可以后缀裁剪的;

相当于我们的分片函数不考虑最后一列,只要前缀相同则放入一个分片内,最后一列对分片结果无影响,因此 non-ordered partitioning 有后缀扩展性。

这两条规则来自于 “如果一批数据他们满足 sorting 性质,那在相应列上也有 grouping 性质”。

口述证明下 (5):“如果数据满足左边的性质,则 C1, C2 ... Cn 列相同的行,一定在同一个分片内,则也一定满足右边的性质”。


数据交换后的性质变化


先来看进过 Partition 后的数据性质变化,如下:

  • Hash 和 Range 模式比较简单,不赘述。
  • Non-Deterministic 模式下,全局性是空集,表示得不到任何性质保证。
  • Broadcast 模式下,全局性的符号表示把数据被全量复制,每个 partition 都有全量数据。
Partition 算子处理数据是是 FIFO 模式的,因此数据的局部顺序在处理前后是一致的。接下来是 Merge 的:

Merge 算子对数据性质的影响,取决于 Merge 的类型和数据原本的局部性质,限于篇幅我们只看一个简单的例子,Sort Merge;Y => S^o 表示输入数据的局部性 Y 能推导出 S^o,也就是输入数据已经在 S 相关的列上有序,则最后聚合结果也满足 S^o。

比如输入数据有 2 个分片,满足性质 {C1^g; C2^o}:{[(1, 1), (3, 3), (1, 5)], [(2, 2), (2, 10)]}。现在按照 C2 做 Sort Merge,结果为:[(1, 1), (2, 2), (3, 3), (1, 5), (2, 10)],其满足 {; C2^o}。


算子的性质要求


接下来分析每种算子对输入数据不同的性质要求。每种算子分为两种模式 partitioned 和 non-partitioned,不同模式有不同要求。

我们还是选一个较为简单的算子 Stream Aggregate 来看一下;

  • non-partitioned:需要数据在 G 上满足 grouping 性质,才能保证聚合的正确性,因此 local property 部分是 {G^g, *};

  • partitioned:全局性质的含义是:“需要要求聚合列的值相同的行,在同一个分片中”,这个要求的形式化表示就是:输入数据的全局性质 X 包含的列,被被聚合的列集合 G 所包含。(这个推论可以根据之前介绍的规则 2 和 5 得到)


性质匹配


因为我们定义了新的数据性质,性质匹配的算法也需要更新,对应的就是前面优化器介绍流程图中的 PropertyMatch 函数。因为任意性质 P 都是由全局性质和局部性质组成,所以只要 P1 的局部和全局性质满足 P2 的,则认为 P1 满足 P2。整个过程就是利用之前我们的推导规则,对性质进行转换,看能否得到另一个性质。

一个简单的例子:

  • P1 = {C1^o;*}

  • P2 = {{C1, C2}^g;*}

我们对 P1 做转换:

1. 根据规则 5:P1 => {{C1}^g; *}

2. 根据规则 2:P1 => {{C1, C2}^g; *}

由于我们将 P1 转换成了 P2,则可以认定 P1 性质满足 P2;


优化器规则


目前为止,理论工作已经完成,接下来要真正的教会优化器利用这些知识。这里我们添加一条规则,用来生成并行的计划;此规则会被优化器在 LogicalTransform 阶段考虑,使得优化器无缝生成并行计划:

上面的过程其实就是对 expr 做下面这 3 种形式的转换:

1. 即使外部要求不并行,通过添加 FullMerge,使得算子本身也能并行起来,同时满足外部要求的性质。

2. 即使外部要求并行,通过添加 Partition 算子,使得算子本身可以以非并行的方式满足外部要求的性质。

3. 使得算子并行的时候不用考虑最外部对性质的要求,如最外部要求 partition 数为 2,算子本身的并发可以任意进行,最后外部性质会被 repartition 保证。

优化器能够生成正确的并发计划后,就可以把这些并发计划纳入代价模型中考虑。


实验


最后是优化器效果的一个实验,原文大概是想执行如下一条 SQL:

有 start 和 end 表用来记录用户操作的开始和结束时间,需要统计这些用户的累计操作时长。下面是优化前后的 Plan 对比:

初始计划在每两个算子之间,都插入了 Repartition 算子,使得每个算子都并行了起来,最后执行了 21 分钟。优化后,在算子 4 的地方,根据规则 (2),可得按照 GUID 进行分片后,就能同时满足两个 StreamAgg 和一个 MergeJoin 的性质要求,就不需要再 repartition 一次了。优化后的计划最终跑了 10 分钟,性能提升了一倍。


总结


论文严谨的对分布式环境下的数据性质进行了定义,并在其上推导出多个演化规则,及对其他算子的影响,为并行计划的正确性打下了理论基础。

接着利用优化器本身的结构,通过添加转换规则,优雅无缝的把这些知识灌输进优化器的“大脑”,“教会”优化器生成正确的并行计划。理论和工程部分都做得很漂亮。

更多的细节详见论文原文




Tip:

上文划线部分均有跳转,由于微信外链限制,

大家可以点击【阅读原文】进入知乎专栏

查看原文、与作者留言互动~


关于投稿:

我们通过【知乎专栏“分布式系统之美”】接收投稿请求,专栏编辑组将在后台进行审稿,通过后将第一时间发布在知乎专栏上~

欢迎大家点击【阅读原文】关注我们的知乎专栏,更希望志趣相同的小伙伴们加入我们,一起创作、分享!




文章转载自分布式系统之美,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论