首批通过分布式安全可靠测评,为关键业务系统打造
EXCHANGE
更新时间:2025-11-06 19:52:48
EXCHANGE 算子用于线程间进行数据交互的算子。
EXCHANGE 算子适用于在分布式场景,一般都是成对出现的,数据源端有一个 OUT 算子,目的端会有一个 IN 算子。
EXCH-IN/OUT
EXCH-IN/OUT 即 EXCHANGE IN/ EXCHANGE OUT 用于将多个分区上的数据汇聚到一起,发送到查询所在的主节点上。
如下例所示,下面的查询中访问了 5 个分区(p0-p4)的数据,其中 1 号算子接受 2 号算子产生的输出,并将数据传出;0 号算子接收多个分区上 1 号算子产生的输出,并将结果汇总输出。
obclient> CREATE TABLE t15 (c1 INT, c2 INT) PARTITION BY HASH(c1) PARTITIONS 5;
Query OK, 0 rows affected
obclient> EXPLAIN SELECT * FROM t15;
+---------------------------------------------------------------------------+
| Query Plan |
+---------------------------------------------------------------------------+
| ============================================================= |
| |ID|OPERATOR |NAME |EST.ROWS|EST.TIME(us)| |
| ------------------------------------------------------------- |
| |0 |PX COORDINATOR | |1 |20 | |
| |1 |└─EXCHANGE OUT DISTR |:EX10000|1 |20 | |
| |2 | └─PX PARTITION ITERATOR| |1 |19 | |
| |3 | └─TABLE FULL SCAN |t15 |1 |19 | |
| ============================================================= |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([INTERNAL_FUNCTION(t15.c1, t15.c2)]), filter(nil), rowset=16 |
| 1 - output([INTERNAL_FUNCTION(t15.c1, t15.c2)]), filter(nil), rowset=16 |
| dop=1 |
| 2 - output([t15.c1], [t15.c2]), filter(nil), rowset=16 |
| force partition granule |
| 3 - output([t15.c1], [t15.c2]), filter(nil), rowset=16 |
| access([t15.c1], [t15.c2]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, |
| range_key([t15.__pk_increment]), range(MIN ; MAX)always true |
+---------------------------------------------------------------------------+
上述示例的执行计划展示中的 outputs & filters 详细列出了 EXCH-IN/OUT 算子的输出信息如下:
| 信息名称 | 含义 |
|---|---|
| PX COORDINATOR | 一种特殊的 EXCHANGE IN 算子,除了能够拉回远程的数据外,还负责调度子计划的执行。 |
| Output | 该算子输出的表达式。 |
| filter | 该算子上的过滤条件。由于示例中 EXCH-IN/OUT 算子没有设置 filter,所以为 nil。 |
EXCH-IN/OUT (REMOTE)
EXCH-IN/OUT (REMOTE) 算子用于将远程的数据(单个分区的数据)拉回本地。
如下例所示,在 A 机器上创建了一张非分区表,在 B 机器上执行查询,读取该表的数据。此时,由于待读取的数据在远程,执行计划中分配了 0 号算子和 1 号算子来拉取远程的数据。其中,1 号算子在 A 机器上执行,读取 t14 表的数据,并将数据传出;0 号算子在 B 机器上执行,接收 1 号算子产生的输出。
obclient> CREATE TABLE t14 (c1 INT, c2 INT);
Query OK, 0 rows affected
obclient> EXPLAIN SELECT * FROM t14;
+--------------------------------------------------------------------+
| Query Plan |
+--------------------------------------------------------------------+
| ===================================================== |
| |ID|OPERATOR |NAME|EST.ROWS|EST.TIME(us)| |
| ----------------------------------------------------- |
| |0 |EXCHANGE IN REMOTE | |1 |5 | |
| |1 |└─EXCHANGE OUT REMOTE| |1 |4 | |
| |2 | └─TABLE FULL SCAN |t14 |1 |3 | |
| ===================================================== |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([t14.C1], [t14.C2]), filter(nil) |
| 1 - output([t14.C1], [t14.C2]), filter(nil) |
| 2 - output([t14.C1], [t14.C2]), filter(nil), rowset=16 |
| access([t14.C1], [t14.C2]), partitions(p0) |
| is_index_back=false, is_global_index=false, |
| range_key([t14.__pk_increment]), range(MIN ; MAX)always true |
+--------------------------------------------------------------------+
上述示例的执行计划展示中的 outputs & filters 详细列出了 EXCH-IN/OUT (REMOTE) 算子的输出信息,字段的含义与 EXCH-IN/OUT 算子相同。
EXCH-IN/OUT (PKEY)
EXCH-IN/OUT (PKEY) 算子用于数据重分区。它通常用于二元算子中,将一侧孩子节点的数据按照另外一侧孩子节点的分区方式进行重分区。
如下示例中,该查询是对两个分区表的数据进行联接,执行计划将 s15 表的数据按照 t13 表的分区方式进行重分区,4 号算子的输入是 s15 表扫描的结果,对于 s15 表的每一行,该算子会根据 t13 表的数据分区,以及根据查询的联接条件,确定一行数据应该发送到哪个节点。
此外,可以看到 3 号算子是一个 EXCHANGE IN DISTR,它是一个特殊的 EXCHANGE IN 算子,它用于在汇总多个分区的数据时,会进行一定的归并排序,在这个执行计划中,3 号算子接收到的每个分区的数据都是按照 c1 列有序排列的,它会对每个接收到的数据进行归并排序,从而保证结果输出结果也是按照 c1 列有序排列的。
obclient> CREATE TABLE t13 (c1 INT, c2 INT) PARTITION BY HASH(c1) PARTITIONS 5;
Query OK, 0 rows affected
obclient> CREATE TABLE s15 (c1 INT PRIMARY KEY, c2 INT) PARTITION BY HASH(c1) PARTITIONS 4;
Query OK, 0 rows affected
obclient> EXPLAIN SELECT * FROM s15, t13 WHERE s15.c1 = t13.c1;
+-------------------------------------------------------------------------------------------+
| Query Plan |
+-------------------------------------------------------------------------------------------+
| ===================================================================== |
| |ID|OPERATOR |NAME |EST.ROWS|EST.TIME(us)| |
| --------------------------------------------------------------------- |
| |0 |PX COORDINATOR | |1 |38 | |
| |1 |└─EXCHANGE OUT DISTR |:EX10001|1 |37 | |
| |2 | └─HASH JOIN | |1 |36 | |
| |3 | ├─EXCHANGE IN DISTR | |1 |17 | |
| |4 | │ └─EXCHANGE OUT DISTR (PKEY)|:EX10000|1 |16 | |
| |5 | │ └─PX PARTITION ITERATOR | |1 |16 | |
| |6 | │ └─TABLE FULL SCAN |s15 |1 |16 | |
| |7 | └─PX PARTITION ITERATOR | |1 |19 | |
| |8 | └─TABLE FULL SCAN |t13 |1 |19 | |
| ===================================================================== |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([INTERNAL_FUNCTION(s15.c1, s15.c2, t13.c1, t13.c2)]), filter(nil), rowset=16 |
| 1 - output([INTERNAL_FUNCTION(s15.c1, s15.c2, t13.c1, t13.c2)]), filter(nil), rowset=16 |
| dop=1 |
| 2 - output([s15.c1], [t13.c1], [s15.c2], [t13.c2]), filter(nil), rowset=16 |
| equal_conds([s15.c1 = t13.c1]), other_conds(nil) |
| 3 - output([s15.c1], [s15.c2]), filter(nil), rowset=16 |
| 4 - output([s15.c1], [s15.c2]), filter(nil), rowset=16 |
| (#keys=1, [s15.c1]), dop=1 |
| 5 - output([s15.c1], [s15.c2]), filter(nil), rowset=16 |
| force partition granule |
| 6 - output([s15.c1], [s15.c2]), filter(nil), rowset=16 |
| access([s15.c1], [s15.c2]), partitions(p[0-3]) |
| is_index_back=false, is_global_index=false, |
| range_key([s15.c1]), range(MIN ; MAX)always true |
| 7 - output([t13.c1], [t13.c2]), filter(nil), rowset=16 |
| affinitize, force partition granule |
| 8 - output([t13.c1], [t13.c2]), filter(nil), rowset=16 |
| access([t13.c1], [t13.c2]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, |
| range_key([t13.__pk_increment]), range(MIN ; MAX)always true |
+-------------------------------------------------------------------------------------------+
上述示例的执行计划展示中的 outputs & filters 详细列出了 EXCH-IN/OUT (PKEY) 算子的输出信息如下:
| 信息名称 | 含义 |
|---|---|
| PX COORDINATOR | 一种特殊的 EXCHANGE IN 算子,除了能够拉回远程的数据外,还负责调度子计划的执行。 |
| Output | 该算子输出的表达式。 |
| filter | 该算子上的过滤条件。由于示例中 EXCH-IN/OUT(PKEY) 算子没有设置 filter,所以为 nil。 |
| pkey | 按照哪一列进行重分区。例如,#keys=1, [s15.c1] 表示按照 c1 列重分区。 |
EXCH-IN/OUT (HASH)
EXCH-IN/OUT (HASH) 算子用于对数据使用一组哈希函数进行重分区。
如下例所示的执行计划中,3-5 号以及 7-8 号是两组使用哈希重分区的 EXCHANGE 算子。这两组算子的作用是把 t12 表和 s14 表的数据按照一组新的哈希函数打散成多份,示例中的哈希列为 t12.c2 和 s14.c2,这保证了 c2 列取值相同的行会被分发到同一份中。基于重分区之后的数据,2 号算子 HASH JOIN 会对每一份数据按照 t12.c2= s14.c2 进行联接。
此外,由于查询中执行了并行度为 2,计划中展示了 dop = 2 (DOP 是 Degree of Parallelism 的缩写)。
obclient> CREATE TABLE t12 (c1 INT, c2 INT) PARTITION BY HASH(c1) PARTITIONS 4;
Query OK, 0 rows affected
obclient> CREATE TABLE s14 (c1 INT, c2 INT) PARTITION BY HASH(c1) PARTITIONS 4;
Query OK, 0 rows affected
obclient> EXPLAIN SELECT /*+PARALLEL(2) LEADING(t12 s14) USE_HASH(s14) PQ_DISTRIBUTE(s14 HASH HASH)*/ * FROM t12, s14 WHERE t12.c2 = s14.c2;
+-------------------------------------------------------------------------------------------+
| Query Plan |
+-------------------------------------------------------------------------------------------+
| ===================================================================== |
| |ID|OPERATOR |NAME |EST.ROWS|EST.TIME(us)| |
| --------------------------------------------------------------------- |
| |0 |PX COORDINATOR | |1 |18 | |
| |1 |└─EXCHANGE OUT DISTR |:EX10002|1 |17 | |
| |2 | └─HASH JOIN | |1 |17 | |
| |3 | ├─EXCHANGE IN DISTR | |1 |9 | |
| |4 | │ └─EXCHANGE OUT DISTR (HASH)|:EX10000|1 |8 | |
| |5 | │ └─PX BLOCK ITERATOR | |1 |8 | |
| |6 | │ └─TABLE FULL SCAN |t12 |1 |8 | |
| |7 | └─EXCHANGE IN DISTR | |1 |9 | |
| |8 | └─EXCHANGE OUT DISTR (HASH)|:EX10001|1 |8 | |
| |9 | └─PX BLOCK ITERATOR | |1 |8 | |
| |10| └─TABLE FULL SCAN |s14 |1 |8 | |
| ===================================================================== |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([INTERNAL_FUNCTION(t12.c1, t12.c2, s14.c1, s14.c2)]), filter(nil), rowset=16 |
| 1 - output([INTERNAL_FUNCTION(t12.c1, t12.c2, s14.c1, s14.c2)]), filter(nil), rowset=16 |
| dop=2 |
| 2 - output([t12.c2], [s14.c2], [t12.c1], [s14.c1]), filter(nil), rowset=16 |
| equal_conds([t12.c2 = s14.c2]), other_conds(nil) |
| 3 - output([t12.c2], [t12.c1]), filter(nil), rowset=16 |
| 4 - output([t12.c2], [t12.c1]), filter(nil), rowset=16 |
| (#keys=1, [t12.c2]), dop=2 |
| 5 - output([t12.c1], [t12.c2]), filter(nil), rowset=16 |
| 6 - output([t12.c1], [t12.c2]), filter(nil), rowset=16 |
| access([t12.c1], [t12.c2]), partitions(p[0-3]) |
| is_index_back=false, is_global_index=false, |
| range_key([t12.__pk_increment]), range(MIN ; MAX)always true |
| 7 - output([s14.c2], [s14.c1]), filter(nil), rowset=16 |
| 8 - output([s14.c2], [s14.c1]), filter(nil), rowset=16 |
| (#keys=1, [s14.c2]), dop=2 |
| 9 - output([s14.c1], [s14.c2]), filter(nil), rowset=16 |
| 10 - output([s14.c1], [s14.c2]), filter(nil), rowset=16 |
| access([s14.c1], [s14.c2]), partitions(p[0-3]) |
| is_index_back=false, is_global_index=false, |
| range_key([s14.__pk_increment]), range(MIN ; MAX)always true |
+-------------------------------------------------------------------------------------------+
其中,PX PARTITION ITERATOR 算子用于按照分区粒度迭代数据,详细信息请参见 GI。
上述示例的执行计划展示中的 outputs & filters 详细列出了 EXCH-IN/OUT (HASH) 算子的输出信息如下:
| 信息名称 | 含义 |
|---|---|
| PX COORDINATOR | 一种特殊的 EXCHANGE IN 算子,除了能够拉回远程的数据外,还负责调度子计划的执行。 |
| Output | 该算子输出的表达式。 |
| filter | 该算子上的过滤条件。由于示例中 EXCH-IN/OUT (HASH) 算子没有设置 filter,所以为 nil。 |
| pkey | 按照哪一列进行哈希重分区。例如,#keys=1, [s14.c2] 表示按照 c2 列进行哈希重分区。 |
EXCH-IN/OUT (BC2HOST)
EXCH-IN/OUT (BC2HOST) 算子用于对输入数据使用 BC2HOST 的方法进行重分区,它会将数据广播到其他线程上。
如下示例的执行计划中,3-4 号是一组使用 BC2HOST 重分区方式的EXCHANGE 算子。它会将 t11 表的数据广播到每个线程上,s13表每个分区的数据都会尝试和被广播的 t11 表数据进行联接。
obclient> CREATE TABLE t11 (c1 INT, c2 INT) PARTITION BY HASH(c1) PARTITIONS 4;
Query OK, 0 rows affected
obclient> CREATE TABLE s13 (c1 INT, c2 INT) PARTITION BY HASH(c1) PARTITIONS 4;
Query OK, 0 rows affected
obclient> INSERT INTO s VALUES (1, 1), (2, 2), (3, 3), (4, 4);
Query OK, 1 rows affected
obclient> EXPALIN SELECT /*+PARALLEL(2) */ * FROM t11, s13 WHERE t11.c2 = s13.c2;
+-------------------------------------------------------------------------------------------+
| Query Plan |
+-------------------------------------------------------------------------------------------+
| ======================================================================== |
| |ID|OPERATOR |NAME |EST.ROWS|EST.TIME(us)| |
| ------------------------------------------------------------------------ |
| |0 |PX COORDINATOR | |1 |18 | |
| |1 |└─EXCHANGE OUT DISTR |:EX10001|1 |17 | |
| |2 | └─SHARED HASH JOIN | |1 |17 | |
| |3 | ├─EXCHANGE IN DISTR | |1 |9 | |
| |4 | │ └─EXCHANGE OUT DISTR (BC2HOST)|:EX10000|1 |8 | |
| |5 | │ └─PX BLOCK ITERATOR | |1 |8 | |
| |6 | │ └─TABLE FULL SCAN |t11 |1 |8 | |
| |7 | └─PX BLOCK ITERATOR | |4 |8 | |
| |8 | └─TABLE FULL SCAN |s13 |4 |8 | |
| ======================================================================== |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([INTERNAL_FUNCTION(t11.c1, t11.c2, s13.c1, s13.c2)]), filter(nil), rowset=16 |
| 1 - output([INTERNAL_FUNCTION(t11.c1, t11.c2, s13.c1, s13.c2)]), filter(nil), rowset=16 |
| dop=2 |
| 2 - output([t11.c2], [s13.c2], [t11.c1], [s13.c1]), filter(nil), rowset=16 |
| equal_conds([t11.c2 = s13.c2]), other_conds(nil) |
| 3 - output([t11.c2], [t11.c1]), filter(nil), rowset=16 |
| 4 - output([t11.c2], [t11.c1]), filter(nil), rowset=16 |
| dop=2 |
| 5 - output([t11.c1], [t11.c2]), filter(nil), rowset=16 |
| 6 - output([t11.c1], [t11.c2]), filter(nil), rowset=16 |
| access([t11.c1], [t11.c2]), partitions(p[0-3]) |
| is_index_back=false, is_global_index=false, |
| range_key([t11.__pk_increment]), range(MIN ; MAX)always true |
| 7 - output([s13.c1], [s13.c2]), filter(nil), rowset=16 |
| 8 - output([s13.c1], [s13.c2]), filter(nil), rowset=16 |
| access([s13.c1], [s13.c2]), partitions(p[0-3]) |
| is_index_back=false, is_global_index=false, |
| range_key([s13.__pk_increment]), range(MIN ; MAX)always true |
+-------------------------------------------------------------------------------------------+
上述示例的执行计划展示中的 outputs & filters 详细列出了 EXCH-IN/OUT (BC2HOST) 算子的信息,字段的含义与 EXCH-IN/OUT 算子相同。