在分布式数据库中,关联查询不可避免地对分布在不同节点的表、分区进行跨节点的关联操作,数据在不同节点间的分发方式对关联查询的性能至关重要。我们或许见到过这样的场景,一个关联查询在不同的 OBServer 上执行时性能差异很大,但是其执行计划除了分布式算子的差异外并无其他任何不同,其性能问题通常是因为是在 NLJ 和 SPF 中对被驱动表的 PX rescan 导致的。
适用版本
OceanBase 数据库所有版本。
名词解释
- NLJ:nested-loop join。
- SPF:subplan filter。
- PX:分布式并行访问。
- Rescan:在关联时,使用驱动表的每一行数据去扫描被驱动表,每一次扫描称为一次 rescan。
概述
在 NLJ(nested-loop join)和 SPF(subplan filter)中,需要针对驱动表中的每一行重新扫描被驱动表中的数据,如果被驱动表是分区表且分布在不同的机器上,在分布式重新扫描(Rescan)的过程中会包括资源释放、调度重启等过程,这些操作常常伴随着消息等待、同步以及网络传输操作,影响执行效率。
- 分布式交互包括释放上一次子协调节点执行资源的控制消息、调度子协调节点的控制消息等,这些若干控制消息的 RPC 交互会带来额外的开销。
- 调度重启的代价高是由于 Nested Loop Join 被驱动表可能是分布式复杂子计划,需要重新建立子计划调度关系,这个重新建立的过程会带来调度的开销。
PX Rescan 示例
EXPLAIN
SELECT *
FROM t1
LEFT JOIN t2
ON t1.v1 = t2.v1
AND t1.v2 = t2.v2
WHERE t1.v3 = 10;
|=====================================================
|ID|OPERATOR |NAME |EST. ROWS|COST|
-----------------------------------------------------
|0 |NESTED-LOOP OUTER JOIN | |99 |261 |
|1 | TABLE SCAN |T1 |10 |40 |
|2 | PX COORDINATOR | |99 |158 |
|3 | EXCHANGE OUT DISTR |:EX10000|99 |92 |
|4 | PX PARTITION ITERATOR| |99 |92 |
|5 | TABLE SCAN |T2 |99 |92 |
=====================================================
执行计划解读
这个关联查询 SQL 使用 NLJ 来完成 T1 LEFT JOIN T2 的操作:
- T1 表是驱动表,T2 表是被驱动表,T1 表和 T2 表根据 V1 和 V2 列值做左外连接。
- 对于 T1 表中的每一行中的 (V1,V2) 值,找出 T2 表中与之相等的行进行连接。
- 5 号算子中对 T2 表的 TABLE SCAN 根据 T1表中的 (V1,V2) 值进行扫描(可以是主键扫描),对于 T1 表的每一行数据,都需要重复上述重新扫描过程。
Rescan 代价评估
如果 T2 是非分区表,这个执行计划会具有比较好的性能,特别是当 T1 和 T2 表都在同一个 Server 上时,只需要简单的内存/IO 操作即可完成 Rescan 过程。 当 T2 表是分区表且分布在不同的 Server 上时,rescan 的效率较低,该执行计划性能可能会很差。此时,执行流程如下。
- 1 号算子从 T1 表中获得一行数据, 计算出 V1, V2 值。
- 2 号算子作为协调节点,重启 2-5 号子计划执行。
- 向所有 T2 表分区的子协调节点发出启动的控制消息,并且释放子协调节点上一次启动所占用的系统资源, 控制消息中包含 T1 表当前行 V1, V2 的参数值。
- 5 号算子在各子协调节点上拿到 T1 表当前行 V1, V2 的参数值,对 T2 表进行扫描。
- 3 号算子向 2 号算子返回匹配行。
- 0 号算子从 2 号算子拿到匹配行后进行连接操作。
- 重复以上流程直到结束。 Nested Loop Join 每读取完一行驱动表数据,都会重新执行第 2 步的流程。当 T1 表返回的记录数比较多时,整个分布式执行过程中 PX rescan 的性能消耗很高。
NLJ Rescan 优化
PX batch rescan 优化
对于上述 PX rescan 带来的性能问题,OceanBase 数据库从 V3.1.2 版本开始提供了 batch rescan 的优化,通过对驱动表中的记录进行分批,每一批数据 rescan 一次被驱动表,从而减小 rescan 的次数,提升性能。Batch rescan 的功能由隐藏参数控制,默认为开启。
| 参数名称 | 描述 | 默认值 |
|---|---|---|
_enable_px_batch_rescan |
控制在 NLJ 生成分布式 PX RESCAN 计划执行时是否使用BATCH RESCAN,可以获得更好的性能。 | True |
在 EXPLAN 结果中,PX batch rescan 的使用可以通过 px_batch_rescan=true 来识别,示例如下。
Outputs & filters:
-------------------------------------
0 - output ...
1 - output ... batch_join=false, px_batch_rescan=true
需要注意的是,到目前为止 PX batch rescan 对于使用 anti/semi Join 的 NLJ 是不支持的,也就是说对于 EXISTS、NOT EXISTS 子查询转换为 Join 的 NLJ 处理,可能无法使用该优化,需要把子查询改写成 LEFT JOIN 来使用 PX batch rescan 提升性能。
Group rescan 优化
与 PX batch rescan 相似,对于 NLJ 中右表为单分区表的场景,OceanBase 数据库提供了 group rescan 的优化,group rescan 是通过 batch join 来完成的。
- 一次从左表读取一批(默认为 1000 条)记录。
- 一次把 1000 条记录推到右表,进行 batch join。
在 EXPLAN 结果中,group rescan 的使用可以通过 batch_join=true 来识别,示例如下。
Outputs & filters:
-------------------------------------
0 - output ...
1 - output ... batch_join=true
Batch join 由隐藏参数控制,默认开启。
| 参数名称 | 描述 | 默认值 |
|---|---|---|
_NLJ_BATCHING_ENABLED |
控制在 NLJ 时是否使用 BATCH JOIN,可以获得更好的性能。 | True |
Semi Join 优化
对于 semi join,如果被驱动表是单分区表,可以使用 BC2HOST 来把驱动表的数据 broadcast 到被驱动表,从而避免分布式的 rescan,提升整个执行计划的性能。如果不是 semi join,或者 semi join 中 被驱动表是分区表,使用 BC2HOST 会带来正确性问题。
示例如下。
create table t1(c1 int primary key, c2 int);
create table t2(c1 int primary key, c2 int);
EXPLAIN
select /*+use_nl(t1 t2)*/ *
from t1
where exists (select 1 from t2 where t2.c1 = t1.c2);
# 改进前的执行计划:
==================================================
|ID|OPERATOR |NAME |EST. ROWS|COST|
--------------------------------------------------
|0 |NESTED-LOOP SEMI JOIN| |3 |257 |
|1 | TABLE SCAN |T1 |6 |37 |
|2 | PX COORDINATOR | |1 |37 |
|3 | EXCHANGE OUT DISTR |:EX10000|1 |36 |
|4 | TABLE GET |T2 |1 |36 |
==================================================
# 改进后的执行计划:
=============================================================
|ID|OPERATOR |NAME |EST. ROWS|COST|
-------------------------------------------------------------
|0 |PX COORDINATOR | |3 |260 |
|1 | EXCHANGE OUT DISTR |:EX10001|3 |259 |
|2 | NESTED-LOOP SEMI JOIN | |3 |259 |
|3 | EXCHANGE IN DISTR | |6 |40 |
|4 | EXCHANGE OUT DISTR (BC2HOST)|:EX10000|6 |37 |
|5 | TABLE SCAN |T1 |6 |37 |
|6 | TABLE GET |T2 |1 |36 |
=============================================================
对于以上 Semi Join,如果把 SQL 发送到 T2 所在的 leader 节点执行,对 T1 的 TABLE SCAN 是远程执行,对 T2 的 TABLE GET 是本地执行,事实上也会避免分布式 rescan 的代价,对该 SQL 的性能有大的提升。