
# SMP场景下的Full Partition-wise Join
SMP场景下的Partition-wise Join计划是基于代价选择的，在路径生成的过程中，会对比Partition-wise Join和非Partition-wise Join路径的估算代价，选择代价较低的路径。其开关参数为GUC参数enable_smp_partitionwise。
SMP场景下的Partition-wise Join分为Full Partition-wise Join和Partial Partition-wise Join两种方式：
- Full Partition-wise Join是指相互Join的两张表为分区策略完全相同的两张分区表。
- Full Partition-wise Join路径生成条件是两张表的分区键是一对相互匹配的Join key。
#### 使用规格
SMP场景下的Full Partition-wise Join的使用规格：
- 支持一级HASH分区表和一级RANGE分区表。
- Hash分区表的分区策略完全相同是指分区键类型相同、分区数相同。
- Range分区表的分区策略完全相同是指分区键类型相同、分区数相同、分区键数量相同、每个分区的边界值相同。
- 非透明多写特性下，仅支持Stream计划。
- 非透明多写特性下，仅支持Join算子在单DN内完成计算，即Join算子的数据不跨节点。
- 支持Hash Join和Nest Loop Join。
- 支持Seqscan、Indexscan、Indexonlyscan、Imcvscan。其中，对于Indexscan和Indexonlyscan，只支持分区Local索引，且索引类型为BTREE或UBTREE。
- 相关规格继承SMP规格，不支持SMP场景下的IUD操作。
- 需要开启SMP功能，且设置query_dop的值大于1。
- 其他规格继承Stream计划和SMP场景的规格。
 
#### 约束
- 不支持对range分布的中间结果表的Full Partition-wise Join。
- 剪枝到单分区场景下不支持Partition-wise Join。
- 仅支持在同一个Node Group下的Partition-wise Join。
- 不支持bucket表、临时表下的Partition-wise Join。
- 对分区指定DN的场景下，不支持Partition-wise Join。
- 不支持二级分区的Partition-wise Join。
- 仅支持Stream 计划，不支持light-proxy、FQS和PGXC计划。
- 在数据倾斜、分区剪枝等场景下，可能导致执行性能劣化。
- 不支持计算倾斜下的Partition-wise Join。
 
#### 示例
非透明多写特性下：
```
--创建Hash分区表，相同分布分区键。
gaussdb=# CREATE TABLE hash_part
(
    a INTEGER,
    b INTEGER,
    c INTEGER
)
DISTRIBUTE BY HASH(a)
PARTITION BY HASH(a)
(
    PARTITION p1,
    PARTITION p2,
    PARTITION p3,
    PARTITION p4,
    PARTITION p5
);
CREATE TABLE
--创建Hash分区表，不同分布分区键。
gaussdb=# CREATE TABLE hash_part_diff_key
(
    a INTEGER,
    b INTEGER,
    c INTEGER
)
DISTRIBUTE BY HASH(a)
PARTITION BY HASH(b)
(
    PARTITION p1,
    PARTITION p2,
    PARTITION p3,
    PARTITION p4,
    PARTITION p5
);
CREATE TABLE
--创建Range分区表，不同分布分区键。
gaussdb=# CREATE TABLE range_part_diff_key
(
    a INTEGER,
    b INTEGER,
    c INTEGER,
    d INTEGER,
    e INTEGER
) DISTRIBUTE BY RANGE(a) (
    SLICE p1 VALUES LESS THAN (100),
    SLICE p2 VALUES LESS THAN (200),
    SLICE p3 VALUES LESS THAN (300),
    SLICE p4 VALUES LESS THAN (400),
    SLICE p5 VALUES LESS THAN (500)
)
PARTITION BY RANGE(b) (
    PARTITION p1 VALUES LESS THAN (100),
    PARTITION p2 VALUES LESS THAN (200),
    PARTITION p3 VALUES LESS THAN (300),
    PARTITION p4 VALUES LESS THAN (400),
    PARTITION p5 VALUES LESS THAN (500)
);
CREATE TABLE
--使用Stream计划。
gaussdb=# SET enable_fast_query_shipping = off;
SET
gaussdb=# SET enable_stream_operator = on;
SET
--debug模式下，设置query_dop为5，开启SMP。
gaussdb=# SET query_dop = 5;
SET
--不走nest loop计划。
gaussdb=# SET enable_material=off;
SET
gaussdb=# SET enable_broadcast=off;
SET
gaussdb=# SET enable_nestloop=off;
SET
-- 插入10万条基础测试数据（a/b/c为随机整数，保证哈希分布的随机性）。
gaussdb=# INSERT INTO hash_part (a, b, c) 
SELECT 
  FLOOR(RANDOM() * 1000000)::INTEGER AS a,
  FLOOR(RANDOM() * 500000)::INTEGER AS b,
  FLOOR(RANDOM() * 1000)::INTEGER AS c
FROM GENERATE_SERIES(1, 100000);
INSERT 0 100000
-- 循环5次，每次将现有数据翻倍插入（最终总数据量 = 10万 * 2^5 = 320万）。
gaussdb=# DO $$
BEGIN
  FOR i IN 1..5 LOOP
    INSERT INTO hash_part (a, b, c) 
    SELECT a, b, c FROM hash_part;
  END LOOP;
END 
$$;
ANONYMOUS BLOCK EXECUTE
gaussdb=# INSERT INTO hash_part_diff_key(a, b, c) 
SELECT 
  FLOOR(RANDOM() * 1000000)::INTEGER AS a,
  FLOOR(RANDOM() * 500000)::INTEGER AS b,
  FLOOR(RANDOM() * 1000)::INTEGER AS c
FROM GENERATE_SERIES(1, 100000);
INSERT 0 100000
gaussdb=# DO $$
BEGIN
  FOR i IN 1..5 LOOP
    INSERT INTO hash_part_diff_key (a, b, c) 
    SELECT a, b, c FROM hash_part_diff_key;
  END LOOP;
END 
$$;
ANONYMOUS BLOCK EXECUTE
gaussdb=# INSERT INTO range_part_diff_key(a, b, c) 
SELECT 
  FLOOR(RANDOM() * 500)::INTEGER AS a,
  FLOOR(RANDOM() * 500)::INTEGER AS b,
  FLOOR(RANDOM() * 500)::INTEGER AS c
FROM GENERATE_SERIES(1, 100000);
INSERT 0 100000
gaussdb=# DO $$
BEGIN
  FOR i IN 1..5 LOOP
    INSERT INTO range_part_diff_key(a, b, c) 
    SELECT a, b, c FROM range_part_diff_key;
  END LOOP;
END 
$$;
ANONYMOUS BLOCK EXECUTE
--analyze数据。
gaussdb=# ANALYZE;
ANALYZE
--关闭SMP场景下的Partition-wise Join开关。
gaussdb=# SET enable_smp_partitionwise = off;
SET
--查看非Partition-wise Join的计划。从计划中可以看出来，在通过Partition Iterator+Partitioned Seq Scan两层算子完成数据扫描之后，通过Streaming(type: LOCAL REDISTRIBUTE)算子对数据进行了一次重分布，用于保证Join算子中数据能够相互匹配。
gaussdb=# EXPLAIN (COSTS OFF) SELECT * FROM hash_part t1, hash_part t2 WHERE t1.a = t2.a;
                                QUERY PLAN
--------------------------------------------------------------------------
 Streaming (type: GATHER)
   Node/s: All datanodes
   ->  Streaming(type: LOCAL GATHER dop: 1/5)
         Spawn on: All datanodes
         ->  Hash Join
               Hash Cond: (t1.a = t2.a)
               ->  Streaming(type: LOCAL REDISTRIBUTE dop: 5/5)
                     Spawn on: All datanodes
                     ->  Partition Iterator
                           Iterations: 5
                           ->  Partitioned Seq Scan on hash_part t1
                                 Selected Partitions:  1..5
               ->  Hash
                     ->  Streaming(type: LOCAL REDISTRIBUTE dop: 5/5)
                           Spawn on: All datanodes
                           ->  Partition Iterator
                                 Iterations: 5
                                 ->  Partitioned Seq Scan on hash_part t2
                                       Selected Partitions:  1..5
(19 rows)
--打开SMP场景下的Partition-wise Join开关。
gaussdb=# SET enable_smp_partitionwise = on;
SET
--查看Partition-wise Join的执行计划。从计划中可以看出，Partition-wise Join计划消除掉了Streaming算子，即数据不再需要在线程之间重新分布，减少了数据搬运的开销，提升了Join操作的性能。
gaussdb=# EXPLAIN (COSTS OFF) SELECT * FROM hash_part t1, hash_part t2 WHERE t1.a = t2.a;
                             QUERY PLAN
--------------------------------------------------------------------
 Streaming (type: GATHER)
   Node/s: All datanodes
   ->  Streaming(type: LOCAL GATHER dop: 1/5)
         Spawn on: All datanodes
         ->  Hash Join (Partition-wise Join)
               Hash Cond: (t1.a = t2.a)
               ->  Partition Iterator
                     Iterations: 5
                     ->  Partitioned Seq Scan on hash_part t1
                           Selected Partitions:  1..5
               ->  Hash
                     ->  Partition Iterator
                           Iterations: 5
                           ->  Partitioned Seq Scan on hash_part t2
                                 Selected Partitions:  1..5
(15 rows)
--查看Partition-wise Join的执行计划。从计划中可以看出，Partition-wise Join计划消除掉了Streaming算子，即数据不再需要在线程之间重新分布，减少了数据搬运的开销，提升了Join操作的性能。hash_part_diff_key表的分布键和分区键都是Join Clause的子集，且一一对应，可以做Full Partition-wise Join。
gaussdb=# EXPLAIN (COSTS OFF) SELECT * FROM hash_part_diff_key t1, hash_part_diff_key t2 WHERE t1.a = t2.a AND t1.b = t2.b;
                                 QUERY PLAN
-----------------------------------------------------------------------------
 Streaming (type: GATHER)
   Node/s: All datanodes
   ->  Streaming(type: LOCAL GATHER dop: 1/5)
         Spawn on: All datanodes
         ->  Hash Join (Partition-wise Join)
               Hash Cond: ((t1.a = t2.a) AND (t1.b = t2.b))
               ->  Partition Iterator
                     Iterations: 5
                     ->  Partitioned Seq Scan on hash_part_diff_key t1
                           Selected Partitions:  1..5
               ->  Hash
                     ->  Partition Iterator
                           Iterations: 5
                           ->  Partitioned Seq Scan on hash_part_diff_key t2
                                 Selected Partitions:  1..5
(15 rows)
--查看Partition-wise Join的执行计划。从计划中可以看出，Partition-wise Join计划消除掉了Streaming算子，即数据不再需要在线程之间重新分布，减少了数据搬运的开销，提升了Join操作的性能。range_part_diff_key表的分布键和分区键都是Join Clause的子集，且一一对应，且range中每个slice和partition的范围都一致，可以做Full Partition-wise Join。
gaussdb=# EXPLAIN (COSTS OFF) SELECT * FROM range_part_diff_key t1, range_part_diff_key t2 WHERE t1.a = t2.a AND t1.b = t2.b;
                                  QUERY PLAN
------------------------------------------------------------------------------
 Streaming (type: GATHER)
   Node/s: All datanodes
   ->  Streaming(type: LOCAL GATHER dop: 1/5)
         Spawn on: All datanodes
         ->  Hash Join (Partition-wise Join)
               Hash Cond: ((t1.a = t2.a) AND (t1.b = t2.b))
               ->  Partition Iterator
                     Iterations: 5
                     ->  Partitioned Seq Scan on range_part_diff_key t1
                           Selected Partitions:  1..5
               ->  Hash
                     ->  Partition Iterator
                           Iterations: 5
                           ->  Partitioned Seq Scan on range_part_diff_key t2
                                 Selected Partitions:  1..5
(15 rows)
-- range分布的中间结果表不再支持Partition-wise Join。
gaussdb=# EXPLAIN (COSTS OFF) SELECT * FROM range_part_diff_key t1, range_part_diff_key t2 , range_part_diff_key t3 WHERE t1.a = t2.a AND t1.b = t2.b AND  t1.a = t3.a AND t1.b = t3.b;
                                        QUERY PLAN
------------------------------------------------------------------------------------------
 Streaming (type: GATHER)
   Node/s: All datanodes
   ->  Streaming(type: LOCAL GATHER dop: 1/5)
         Spawn on: All datanodes
         ->  Hash Join
               Hash Cond: ((t1.a = t3.a) AND (t1.b = t3.b))
               ->  Streaming(type: LOCAL REDISTRIBUTE dop: 5/5)
                     Spawn on: All datanodes
                     ->  Hash Join (Partition-wise Join)
                           Hash Cond: ((t1.a = t2.a) AND (t1.b = t2.b))
                           ->  Partition Iterator
                                 Iterations: 5
                                 ->  Partitioned Seq Scan on range_part_diff_key t1
                                       Selected Partitions:  1..5
                           ->  Hash
                                 ->  Partition Iterator
                                       Iterations: 5
                                       ->  Partitioned Seq Scan on range_part_diff_key t2
                                             Selected Partitions:  1..5
               ->  Hash
                     ->  Streaming(type: LOCAL REDISTRIBUTE dop: 5/5)
                           Spawn on: All datanodes
                           ->  Partition Iterator
                                 Iterations: 5
                                 ->  Partitioned Seq Scan on range_part_diff_key t3
                                       Selected Partitions:  1..5
(26 rows)
-- 删除分区表。
gaussdb=# DROP TABLE hash_part;
gaussdb=# DROP TABLE hash_part_diff_key;
gaussdb=# DROP TABLE range_part_diff_key;
```
![](https://support.huaweicloud.com/distributed-devg-v10-gaussdb/public_sys-resources/note_3.0-zh-cn.png)
- 仅在SMP场景下的Partition-wise Join计划中，Join算子右侧会有Partition-wise Join提示信息。非SMP场景无该提示信息。
- Partition-wise Join计划的执行性能取决于数据量最大的分区的执行性能。所以在分区间存在严重的数据倾斜或存在分区剪枝的场景下，Partition-wise Join可能导致性能劣化。
 
