From 3d01cbcd0e6217c7c39eb6bf9e52bdf5bfae51fb Mon Sep 17 00:00:00 2001 From: Alena Rybakina Date: Thu, 2 Jul 2026 15:53:09 +0300 Subject: [PATCH 1/3] orca: fall back on replicated CTE only when consumed in 2+ slices Only a replicated CTE consumed in 2+ different slices hangs, not any cross-slice consumer. Narrow the fallback check to that case. --- .../gporca/libgpopt/src/base/CUtils.cpp | 30 ++++- .../expected/qp_orca_fallback_optimizer.out | 43 +++--- src/test/regress/expected/shared_scan.out | 123 ++++++++++++++++-- .../expected/shared_scan_optimizer.out | 115 ++++++++++++++-- src/test/regress/sql/shared_scan.sql | 55 ++++++-- 5 files changed, 314 insertions(+), 52 deletions(-) diff --git a/src/backend/gporca/libgpopt/src/base/CUtils.cpp b/src/backend/gporca/libgpopt/src/base/CUtils.cpp index d114a639449..66d69cf3416 100644 --- a/src/backend/gporca/libgpopt/src/base/CUtils.cpp +++ b/src/backend/gporca/libgpopt/src/base/CUtils.cpp @@ -999,23 +999,39 @@ CUtils::FHasCrossSliceReplicatedCTEConsumer(CMemoryPool *mp, CExpression *pexpr) CollectCTESlices(mp, pexpr, 0 /*curSlice*/, &nextSlice, prodInfos, consInfos); + // Fall back only when one replicated CTE is consumed in two or more + // different slices (e.g. joined twice on different columns) -- that is the + // case that hangs. One consumer, or several in the same slice, stays on + // ORCA. BOOL cross = false; - for (ULONG ic = 0; ic < consInfos->Size(); ic++) + for (ULONG ip = 0; ip < prodInfos->Size() && !cross; ip++) { - SCTEInfo *cons = (*consInfos)[ic]; + SCTEInfo *prod = (*prodInfos)[ip]; - for (ULONG ip = 0; ip < prodInfos->Size(); ip++) + ULONG firstConsSlice = 0; + BOOL haveFirst = false; + + for (ULONG ic = 0; ic < consInfos->Size(); ic++) { - SCTEInfo *prod = (*prodInfos)[ip]; - if (prod->cteId == cons->cteId && prod->sliceId != cons->sliceId) + SCTEInfo *cons = (*consInfos)[ic]; + if (cons->cteId != prod->cteId) + { + continue; + } + + if (!haveFirst) + { + firstConsSlice = cons->sliceId; + haveFirst = true; + } + else if (cons->sliceId != firstConsSlice) { cross = true; - goto lExit; + break; } } } -lExit: prodInfos->Release(); consInfos->Release(); diff --git a/src/test/regress/expected/qp_orca_fallback_optimizer.out b/src/test/regress/expected/qp_orca_fallback_optimizer.out index c08b388756c..18bf2d8a196 100644 --- a/src/test/regress/expected/qp_orca_fallback_optimizer.out +++ b/src/test/regress/expected/qp_orca_fallback_optimizer.out @@ -409,26 +409,27 @@ t2 AS (SELECT id, refrcode FROM tbl2 WHERE REFERENCEID = 101991) JOIN t2 r1 ON p.isCalcDetail = r1.RefrCode LIMIT 1; -INFO: GPORCA failed to produce a plan, falling back to planner -DETAIL: Feature not supported: CTE Consumer placed on a different slice than its replicated Producer - QUERY PLAN ---------------------------------------------------------------------------------------------- - Limit (cost=0.34..83.24 rows=1 width=104) - -> Gather Motion 1:1 (slice1; segments: 1) (cost=0.34..581.83 rows=8 width=104) - -> Hash Join (cost=0.34..581.69 rows=8 width=104) - Hash Cond: ((tbl1.iscalcdetail)::text = (r1.refrcode)::text) - -> Hash Join (cost=0.17..580.97 rows=130 width=104) - Hash Cond: ((tbl1.iscalctrg)::text = (r.refrcode)::text) - -> Seq Scan on tbl1 (cost=0.00..344.00 rows=24400 width=104) - -> Hash (cost=0.11..0.11 rows=2 width=516) - -> Subquery Scan on r (cost=0.00..0.11 rows=6 width=516) - -> Seq Scan on tbl2 (cost=0.00..166.25 rows=6 width=548) - Filter: (referenceid = '101991'::numeric) - -> Hash (cost=0.11..0.11 rows=2 width=516) - -> Subquery Scan on r1 (cost=0.00..0.11 rows=6 width=516) - -> Seq Scan on tbl2 tbl2_1 (cost=0.00..166.25 rows=6 width=548) - Filter: (referenceid = '101991'::numeric) - Optimizer: Postgres query optimizer -(16 rows) + QUERY PLAN +---------------------------------------------------------------------------------------------------------------------- + Gather Motion 3:1 (slice3; segments: 3) (cost=0.00..1724.00 rows=1 width=24) + -> Sequence (cost=0.00..1724.00 rows=1 width=24) + -> Shared Scan (share slice:id 3:1) (cost=0.00..431.00 rows=1 width=1) + -> Materialize (cost=0.00..431.00 rows=1 width=1) + -> Seq Scan on tbl2 (cost=0.00..431.00 rows=1 width=263) + Filter: (referenceid = '101991'::numeric) + -> Redistribute Motion 1:3 (slice2) (cost=0.00..1293.00 rows=1 width=24) + -> Limit (cost=0.00..1293.00 rows=1 width=24) + -> Gather Motion 1:1 (slice1; segments: 1) (cost=0.00..1293.00 rows=1 width=24) + -> Hash Join (cost=0.00..1293.00 rows=1 width=24) + Hash Cond: ((tbl1.iscalctrg)::text = (share1_ref2.refrcode)::text) + -> Hash Join (cost=0.00..862.00 rows=1 width=24) + Hash Cond: ((tbl1.iscalcdetail)::text = (share1_ref3.refrcode)::text) + -> Seq Scan on tbl1 (cost=0.00..431.00 rows=1 width=24) + -> Hash (cost=431.00..431.00 rows=1 width=8) + -> Shared Scan (share slice:id 1:1) (cost=0.00..431.00 rows=1 width=8) + -> Hash (cost=431.00..431.00 rows=1 width=8) + -> Shared Scan (share slice:id 1:1) (cost=0.00..431.00 rows=1 width=8) + Optimizer: Pivotal Optimizer (GPORCA) +(19 rows) DROP TABLE tbl1, tbl2; diff --git a/src/test/regress/expected/shared_scan.out b/src/test/regress/expected/shared_scan.out index 3a4fddef1ce..58d408d28ee 100644 --- a/src/test/regress/expected/shared_scan.out +++ b/src/test/regress/expected/shared_scan.out @@ -250,11 +250,11 @@ where Optimizer: Postgres query optimizer (52 rows) --- ORCA should fallback when a CTE over a replicated table is referenced --- from multiple scalar subqueries. --- ss_t1 needs enough rows (40000) to push ORCA to the cross-slice plan; --- with fewer rows the bug does not manifest and the test would silently --- pass even without the fix. +-- A CTE over a replicated table referenced from several scalar subqueries. +-- Each CTE is consumed within a single slice, so it does not hit the +-- multi-slice replicated-CTE fallback: ORCA keeps the shared-scan plan and the +-- query must finish without hanging. +-- ss_t1 needs enough rows (40000) for ORCA to choose the shared-scan plan. CREATE TABLE ss_t1 AS SELECT generate_series(1, 40000) id DISTRIBUTED BY (id); @@ -264,6 +264,39 @@ CREATE TABLE ss_t2 AS ANALYZE ss_t1; ANALYZE ss_t2; SET statement_timeout = '15s'; +EXPLAIN (COSTS OFF) +WITH + cte1 AS (SELECT v FROM ss_t2 WHERE id = 1), + cte2 AS (SELECT v FROM ss_t2 WHERE id = 2) + SELECT (SELECT v FROM cte1) + (SELECT v FROM cte2) + + (SELECT v FROM cte1) + (SELECT v FROM cte2) AS result + FROM ss_t1 + LIMIT 1; + QUERY PLAN +-------------------------------------------------- + Limit + InitPlan 1 (returns $0) (slice7) + -> Gather Motion 1:1 (slice2; segments: 1) + -> Seq Scan on ss_t2 + Filter: (id = 1) + InitPlan 2 (returns $1) (slice8) + -> Gather Motion 1:1 (slice3; segments: 1) + -> Seq Scan on ss_t2 ss_t2_1 + Filter: (id = 2) + InitPlan 3 (returns $2) (slice9) + -> Gather Motion 1:1 (slice4; segments: 1) + -> Seq Scan on ss_t2 ss_t2_2 + Filter: (id = 1) + InitPlan 4 (returns $3) (slice6) + -> Gather Motion 1:1 (slice5; segments: 1) + -> Seq Scan on ss_t2 ss_t2_3 + Filter: (id = 2) + -> Gather Motion 3:1 (slice1; segments: 3) + -> Limit + -> Seq Scan on ss_t1 + Optimizer: Postgres query optimizer +(21 rows) + WITH cte1 AS (SELECT v FROM ss_t2 WHERE id = 1), cte2 AS (SELECT v FROM ss_t2 WHERE id = 2) @@ -278,6 +311,80 @@ WITH RESET statement_timeout; DROP TABLE ss_t1, ss_t2; +-- A replicated CTE with a UNION ALL body referenced from two correlated SubPlans +CREATE TABLE ss_c1 (id bigint, iscalctrg varchar(15) NOT NULL, iscalcdetail varchar(15)) + DISTRIBUTED REPLICATED; +CREATE TABLE ss_c2 (id numeric, refrcode varchar(255), referenceid numeric) + DISTRIBUTED REPLICATED; +INSERT INTO ss_c1 SELECT i, 'A'||(i%5), 'A'||(i%7) FROM generate_series(1, 50000) i; +INSERT INTO ss_c2 SELECT i, 'A'||(i%5), 101991 FROM generate_series(1, 50000) i; +ANALYZE ss_c1; +ANALYZE ss_c2; +SET statement_timeout = '15s'; +EXPLAIN (COSTS OFF) +WITH cte AS ( + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id < 25000 + UNION ALL + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id >= 25000 +) + SELECT (SELECT refrcode FROM cte WHERE refrcode = p.iscalctrg LIMIT 1) = 'A1' + AND (SELECT refrcode FROM cte WHERE refrcode = p.iscalcdetail LIMIT 1) = 'A1' AS ok + FROM ss_c1 p WHERE p.id = 1; + QUERY PLAN +------------------------------------------------------------------------------------------------------------------------------ + Gather Motion 1:1 (slice5; segments: 1) + -> Seq Scan on ss_c1 p + Filter: (id = 1) + SubPlan 1 (slice5; segments: 1) + -> Limit + -> Subquery Scan on cte + -> Append + -> Result + Filter: ((ss_c2.refrcode)::text = (p.iscalctrg)::text) + -> Materialize + -> Gather Motion 1:1 (slice1; segments: 1) + -> Seq Scan on ss_c2 + Filter: ((id < '25000'::numeric) AND (referenceid = '101991'::numeric)) + -> Result + Filter: ((ss_c2_1.refrcode)::text = (p.iscalctrg)::text) + -> Materialize + -> Gather Motion 1:1 (slice2; segments: 1) + -> Seq Scan on ss_c2 ss_c2_1 + Filter: ((id >= '25000'::numeric) AND (referenceid = '101991'::numeric)) + SubPlan 2 (slice5; segments: 1) + -> Limit + -> Subquery Scan on cte_1 + -> Append + -> Result + Filter: ((ss_c2_2.refrcode)::text = (p.iscalcdetail)::text) + -> Materialize + -> Gather Motion 1:1 (slice3; segments: 1) + -> Seq Scan on ss_c2 ss_c2_2 + Filter: ((id < '25000'::numeric) AND (referenceid = '101991'::numeric)) + -> Result + Filter: ((ss_c2_3.refrcode)::text = (p.iscalcdetail)::text) + -> Materialize + -> Gather Motion 1:1 (slice4; segments: 1) + -> Seq Scan on ss_c2 ss_c2_3 + Filter: ((id >= '25000'::numeric) AND (referenceid = '101991'::numeric)) + Optimizer: Postgres query optimizer +(36 rows) + +WITH cte AS ( + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id < 25000 + UNION ALL + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id >= 25000 +) + SELECT (SELECT refrcode FROM cte WHERE refrcode = p.iscalctrg LIMIT 1) = 'A1' + AND (SELECT refrcode FROM cte WHERE refrcode = p.iscalcdetail LIMIT 1) = 'A1' AS ok + FROM ss_c1 p WHERE p.id = 1; + ok +---- + t +(1 row) + +RESET statement_timeout; +DROP TABLE ss_c1, ss_c2; -- Test the scenario which already opened many fds -- start_ignore RESET search_path; @@ -322,7 +429,7 @@ create table t1 (a int, b int, c int) distributed by (a); explain (costs off) with cte1 as ( select max(c) as c from t1 -), +), cte2 as ( select d as c from generate_series( @@ -354,7 +461,7 @@ left join t1 u on l.c = u.c; with cte1 as ( select max(c) as c from t1 -), +), cte2 as ( select d as c from generate_series( @@ -385,7 +492,7 @@ from gp_segment_configuration where role = 'p' and content = -1; select gp_inject_fault('material_pre_tuplestore_flush', 'sleep', '', '', '', 1, 1, 5, dbid) -from gp_segment_configuration where role = 'p' and content = -1; +from gp_segment_configuration where role = 'p' and content = -1; gp_inject_fault ----------------- Success: diff --git a/src/test/regress/expected/shared_scan_optimizer.out b/src/test/regress/expected/shared_scan_optimizer.out index a1371f86f6d..a8e2247d0a0 100644 --- a/src/test/regress/expected/shared_scan_optimizer.out +++ b/src/test/regress/expected/shared_scan_optimizer.out @@ -263,11 +263,11 @@ where Optimizer: Postgres query optimizer (52 rows) --- ORCA should fallback when a CTE over a replicated table is referenced --- from multiple scalar subqueries. --- ss_t1 needs enough rows (40000) to push ORCA to the cross-slice plan; --- with fewer rows the bug does not manifest and the test would silently --- pass even without the fix. +-- A CTE over a replicated table referenced from several scalar subqueries. +-- Each CTE is consumed within a single slice, so it does not hit the +-- multi-slice replicated-CTE fallback: ORCA keeps the shared-scan plan and the +-- query must finish without hanging. +-- ss_t1 needs enough rows (40000) for ORCA to choose the shared-scan plan. CREATE TABLE ss_t1 AS SELECT generate_series(1, 40000) id DISTRIBUTED BY (id); @@ -277,6 +277,44 @@ CREATE TABLE ss_t2 AS ANALYZE ss_t1; ANALYZE ss_t2; SET statement_timeout = '15s'; +EXPLAIN (COSTS OFF) +WITH + cte1 AS (SELECT v FROM ss_t2 WHERE id = 1), + cte2 AS (SELECT v FROM ss_t2 WHERE id = 2) + SELECT (SELECT v FROM cte1) + (SELECT v FROM cte2) + + (SELECT v FROM cte1) + (SELECT v FROM cte2) AS result + FROM ss_t1 + LIMIT 1; + QUERY PLAN +------------------------------------------------------------------------------------ + Gather Motion 3:1 (slice3; segments: 3) + -> Sequence + -> Shared Scan (share slice:id 3:0) + -> Materialize + -> Seq Scan on ss_t2 ss_t2_1 + Filter: (id = 1) + -> Sequence + -> Shared Scan (share slice:id 3:1) + -> Materialize + -> Seq Scan on ss_t2 + Filter: (id = 2) + -> Redistribute Motion 1:3 (slice2) + -> Limit + -> Gather Motion 3:1 (slice1; segments: 3) + -> Limit + -> Result + -> Seq Scan on ss_t1 + SubPlan 1 (slice1; segments: 3) + -> Shared Scan (share slice:id 1:0) + SubPlan 2 (slice1; segments: 3) + -> Shared Scan (share slice:id 1:1) + SubPlan 3 (slice1; segments: 3) + -> Shared Scan (share slice:id 1:0) + SubPlan 4 (slice1; segments: 3) + -> Shared Scan (share slice:id 1:1) + Optimizer: Pivotal Optimizer (GPORCA) +(26 rows) + WITH cte1 AS (SELECT v FROM ss_t2 WHERE id = 1), cte2 AS (SELECT v FROM ss_t2 WHERE id = 2) @@ -291,6 +329,67 @@ WITH RESET statement_timeout; DROP TABLE ss_t1, ss_t2; +-- A replicated CTE with a UNION ALL body referenced from two correlated SubPlans +CREATE TABLE ss_c1 (id bigint, iscalctrg varchar(15) NOT NULL, iscalcdetail varchar(15)) + DISTRIBUTED REPLICATED; +CREATE TABLE ss_c2 (id numeric, refrcode varchar(255), referenceid numeric) + DISTRIBUTED REPLICATED; +INSERT INTO ss_c1 SELECT i, 'A'||(i%5), 'A'||(i%7) FROM generate_series(1, 50000) i; +INSERT INTO ss_c2 SELECT i, 'A'||(i%5), 101991 FROM generate_series(1, 50000) i; +ANALYZE ss_c1; +ANALYZE ss_c2; +SET statement_timeout = '15s'; +EXPLAIN (COSTS OFF) +WITH cte AS ( + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id < 25000 + UNION ALL + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id >= 25000 +) + SELECT (SELECT refrcode FROM cte WHERE refrcode = p.iscalctrg LIMIT 1) = 'A1' + AND (SELECT refrcode FROM cte WHERE refrcode = p.iscalcdetail LIMIT 1) = 'A1' AS ok + FROM ss_c1 p WHERE p.id = 1; + QUERY PLAN +---------------------------------------------------------------------------------------------------------- + Gather Motion 1:1 (slice1; segments: 1) + -> Sequence + -> Shared Scan (share slice:id 1:0) + -> Materialize + -> Append + -> Seq Scan on ss_c2 + Filter: ((referenceid = '101991'::numeric) AND (id < '25000'::numeric)) + -> Seq Scan on ss_c2 ss_c2_1 + Filter: ((referenceid = '101991'::numeric) AND (id >= '25000'::numeric)) + -> Result + -> Seq Scan on ss_c1 + Filter: (id = 1) + SubPlan 1 (slice1; segments: 1) + -> Limit + -> Result + Filter: ((share0_ref2.refrcode)::text = (ss_c1.iscalctrg)::text) + -> Shared Scan (share slice:id 1:0) + SubPlan 2 (slice1; segments: 1) + -> Limit + -> Result + Filter: ((share0_ref3.refrcode)::text = (ss_c1.iscalcdetail)::text) + -> Shared Scan (share slice:id 1:0) + Optimizer: Pivotal Optimizer (GPORCA) +(23 rows) + +WITH cte AS ( + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id < 25000 + UNION ALL + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id >= 25000 +) + SELECT (SELECT refrcode FROM cte WHERE refrcode = p.iscalctrg LIMIT 1) = 'A1' + AND (SELECT refrcode FROM cte WHERE refrcode = p.iscalcdetail LIMIT 1) = 'A1' AS ok + FROM ss_c1 p WHERE p.id = 1; + ok +---- + t +(1 row) + +RESET statement_timeout; +DROP TABLE ss_c1, ss_c2; -- Test the scenario which already opened many fds -- start_ignore RESET search_path; @@ -335,7 +434,7 @@ create table t1 (a int, b int, c int) distributed by (a); explain (costs off) with cte1 as ( select max(c) as c from t1 -), +), cte2 as ( select d as c from generate_series( @@ -391,7 +490,7 @@ left join t1 u on l.c = u.c; with cte1 as ( select max(c) as c from t1 -), +), cte2 as ( select d as c from generate_series( @@ -422,7 +521,7 @@ from gp_segment_configuration where role = 'p' and content = -1; select gp_inject_fault('material_pre_tuplestore_flush', 'sleep', '', '', '', 1, 1, 5, dbid) -from gp_segment_configuration where role = 'p' and content = -1; +from gp_segment_configuration where role = 'p' and content = -1; gp_inject_fault ----------------- Success: diff --git a/src/test/regress/sql/shared_scan.sql b/src/test/regress/sql/shared_scan.sql index f7a4a06d186..95f5ea88ee8 100644 --- a/src/test/regress/sql/shared_scan.sql +++ b/src/test/regress/sql/shared_scan.sql @@ -125,11 +125,11 @@ where and (stat.schema_name || '.' ||stat.table_name not in (select table_nm_onl_act from tbls_w_onl_actl_data)) or (stat.schema_name || '.' ||stat.table_name in (select table_nm_onl_act from tbls_w_onl_actl_data)); --- ORCA should fallback when a CTE over a replicated table is referenced --- from multiple scalar subqueries. --- ss_t1 needs enough rows (40000) to push ORCA to the cross-slice plan; --- with fewer rows the bug does not manifest and the test would silently --- pass even without the fix. +-- A CTE over a replicated table referenced from several scalar subqueries. +-- Each CTE is consumed within a single slice, so it does not hit the +-- multi-slice replicated-CTE fallback: ORCA keeps the shared-scan plan and the +-- query must finish without hanging. +-- ss_t1 needs enough rows (40000) for ORCA to choose the shared-scan plan. CREATE TABLE ss_t1 AS SELECT generate_series(1, 40000) id DISTRIBUTED BY (id); @@ -140,6 +140,14 @@ ANALYZE ss_t1; ANALYZE ss_t2; SET statement_timeout = '15s'; +EXPLAIN (COSTS OFF) +WITH + cte1 AS (SELECT v FROM ss_t2 WHERE id = 1), + cte2 AS (SELECT v FROM ss_t2 WHERE id = 2) + SELECT (SELECT v FROM cte1) + (SELECT v FROM cte2) + + (SELECT v FROM cte1) + (SELECT v FROM cte2) AS result + FROM ss_t1 + LIMIT 1; WITH cte1 AS (SELECT v FROM ss_t2 WHERE id = 1), cte2 AS (SELECT v FROM ss_t2 WHERE id = 2) @@ -150,6 +158,37 @@ WITH RESET statement_timeout; DROP TABLE ss_t1, ss_t2; +-- A replicated CTE with a UNION ALL body referenced from two correlated SubPlans +CREATE TABLE ss_c1 (id bigint, iscalctrg varchar(15) NOT NULL, iscalcdetail varchar(15)) + DISTRIBUTED REPLICATED; +CREATE TABLE ss_c2 (id numeric, refrcode varchar(255), referenceid numeric) + DISTRIBUTED REPLICATED; +INSERT INTO ss_c1 SELECT i, 'A'||(i%5), 'A'||(i%7) FROM generate_series(1, 50000) i; +INSERT INTO ss_c2 SELECT i, 'A'||(i%5), 101991 FROM generate_series(1, 50000) i; +ANALYZE ss_c1; +ANALYZE ss_c2; + +SET statement_timeout = '15s'; +EXPLAIN (COSTS OFF) +WITH cte AS ( + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id < 25000 + UNION ALL + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id >= 25000 +) + SELECT (SELECT refrcode FROM cte WHERE refrcode = p.iscalctrg LIMIT 1) = 'A1' + AND (SELECT refrcode FROM cte WHERE refrcode = p.iscalcdetail LIMIT 1) = 'A1' AS ok + FROM ss_c1 p WHERE p.id = 1; +WITH cte AS ( + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id < 25000 + UNION ALL + SELECT id, refrcode FROM ss_c2 WHERE referenceid = 101991 AND id >= 25000 +) + SELECT (SELECT refrcode FROM cte WHERE refrcode = p.iscalctrg LIMIT 1) = 'A1' + AND (SELECT refrcode FROM cte WHERE refrcode = p.iscalcdetail LIMIT 1) = 'A1' AS ok + FROM ss_c1 p WHERE p.id = 1; +RESET statement_timeout; +DROP TABLE ss_c1, ss_c2; + -- Test the scenario which already opened many fds -- start_ignore RESET search_path; @@ -173,7 +212,7 @@ create table t1 (a int, b int, c int) distributed by (a); explain (costs off) with cte1 as ( select max(c) as c from t1 -), +), cte2 as ( select d as c from generate_series( @@ -185,7 +224,7 @@ left join t1 u on l.c = u.c; with cte1 as ( select max(c) as c from t1 -), +), cte2 as ( select d as c from generate_series( @@ -209,7 +248,7 @@ select gp_inject_fault('material_pre_tuplestore_flush', 'reset', dbid) from gp_segment_configuration where role = 'p' and content = -1; select gp_inject_fault('material_pre_tuplestore_flush', 'sleep', '', '', '', 1, 1, 5, dbid) -from gp_segment_configuration where role = 'p' and content = -1; +from gp_segment_configuration where role = 'p' and content = -1; set optimizer_parallel_union = on; explain (costs off) From 496b249fdf9d81b86bdfb6bf59f14eb4a316f72f Mon Sep 17 00:00:00 2001 From: Alena Rybakina Date: Mon, 13 Jul 2026 02:47:50 +0300 Subject: [PATCH 2/3] orca: remove expr-level cross-slice replicated-CTE walker FHasCrossSliceReplicatedCTEConsumer modeled slice assignment on the ORCA CExpression tree, before DXL/Plan translation. --- .../libgpopt/include/gpopt/base/CUtils.h | 3 - .../gporca/libgpopt/src/base/CUtils.cpp | 138 ------------------ .../src/translate/CTranslatorExprToDXL.cpp | 14 -- 3 files changed, 155 deletions(-) diff --git a/src/backend/gporca/libgpopt/include/gpopt/base/CUtils.h b/src/backend/gporca/libgpopt/include/gpopt/base/CUtils.h index 7d335a41f09..de3cdf964f7 100644 --- a/src/backend/gporca/libgpopt/include/gpopt/base/CUtils.h +++ b/src/backend/gporca/libgpopt/include/gpopt/base/CUtils.h @@ -1024,9 +1024,6 @@ class CUtils static BOOL FScalarConstOrBinaryCoercible(CExpression *pexpr); - static BOOL FHasCrossSliceReplicatedCTEConsumer(CMemoryPool *mp, - CExpression *pexpr); - static CExpression *ReplaceColrefWithProjectExpr(CMemoryPool *mp, CExpression *pexpr, CColRef *pcolref, diff --git a/src/backend/gporca/libgpopt/src/base/CUtils.cpp b/src/backend/gporca/libgpopt/src/base/CUtils.cpp index 66d69cf3416..1b85d6ab4c3 100644 --- a/src/backend/gporca/libgpopt/src/base/CUtils.cpp +++ b/src/backend/gporca/libgpopt/src/base/CUtils.cpp @@ -901,144 +901,6 @@ CUtils::FHasCTEAnchor(CExpression *pexpr) return false; } -// True if the distribution is replicated-like. -static BOOL -FReplicatedLikeDistribution(CDistributionSpec::EDistributionType edt) -{ - return (CDistributionSpec::EdtStrictReplicated == edt || - CDistributionSpec::EdtTaintedReplicated == edt || - CDistributionSpec::EdtUniversal == edt); -} - -struct SCTEInfo -{ - ULONG cteId; - ULONG sliceId; - - SCTEInfo(ULONG cte_id, ULONG slice_id) : cteId(cte_id), sliceId(slice_id) - { - } -}; - -typedef CDynamicPtrArray > CTEInfoArray; - -// Walk the physical tree, recording the slice id of every replicated -// CTE Producer and every CTE Consumer. Slices are delimited by Motion -// nodes: each non-scalar child of a Motion lives in a fresh slice -- -// same motId-stack idea as in apply_shareinput_xslice. -static void -CollectCTESlices(CMemoryPool *mp, CExpression *pexpr, ULONG curSlice, - ULONG *pNextSlice, CTEInfoArray *prodInfos, - CTEInfoArray *consInfos) -{ - COperator *pop = pexpr->Pop(); - - if (COperator::EopPhysicalCTEProducer == pop->Eopid()) - { - // Producer's distribution comes from its only child -- inspect - // it there. Skip non-replicated Producers; they cannot trigger - // the cross-slice issue we are checking for. - GPOS_ASSERT(1 == pexpr->Arity()); - CExpression *pexprChild = (*pexpr)[0]; - CDrvdPropPlan *pdpplan = - CDrvdPropPlan::Pdpplan(pexprChild->PdpDerive()); - - if (FReplicatedLikeDistribution(pdpplan->Pds()->Edt())) - { - prodInfos->Append(GPOS_NEW(mp) SCTEInfo( - CPhysicalCTEProducer::PopConvert(pop)->UlCTEId(), curSlice)); - } - } - else if (COperator::EopPhysicalCTEConsumer == pop->Eopid()) - { - // Consumer is a leaf -- record (cteId, curSlice) and let the - // caller decide later, once the whole tree has been walked. - consInfos->Append(GPOS_NEW(mp) SCTEInfo( - CPhysicalCTEConsumer::PopConvert(pop)->UlCTEId(), curSlice)); - } - - BOOL isMotion = CUtils::FPhysicalMotion(pop); - - for (ULONG ul = 0; ul < pexpr->Arity(); ul++) - { - CExpression *pexprChild = (*pexpr)[ul]; - - // Scalar subtrees (predicates, project elements) never run as - // separate executor groups, so they cannot host a slice. - if (pexprChild->Pop()->FScalar()) - { - continue; - } - - // Allocate a fresh slice id for each non-scalar child of a - // Motion; otherwise the child stays in the parent's slice. - ULONG childSlice = curSlice; - if (isMotion) - { - (*pNextSlice)++; - childSlice = *pNextSlice; - } - - CollectCTESlices(mp, pexprChild, childSlice, pNextSlice, prodInfos, - consInfos); - } -} - -BOOL -CUtils::FHasCrossSliceReplicatedCTEConsumer(CMemoryPool *mp, CExpression *pexpr) -{ - if (NULL == pexpr) - { - return false; - } - - CTEInfoArray *prodInfos = GPOS_NEW(mp) CTEInfoArray(mp); - CTEInfoArray *consInfos = GPOS_NEW(mp) CTEInfoArray(mp); - ULONG nextSlice = 0; - - CollectCTESlices(mp, pexpr, 0 /*curSlice*/, &nextSlice, prodInfos, - consInfos); - - // Fall back only when one replicated CTE is consumed in two or more - // different slices (e.g. joined twice on different columns) -- that is the - // case that hangs. One consumer, or several in the same slice, stays on - // ORCA. - BOOL cross = false; - - for (ULONG ip = 0; ip < prodInfos->Size() && !cross; ip++) - { - SCTEInfo *prod = (*prodInfos)[ip]; - - ULONG firstConsSlice = 0; - BOOL haveFirst = false; - - for (ULONG ic = 0; ic < consInfos->Size(); ic++) - { - SCTEInfo *cons = (*consInfos)[ic]; - if (cons->cteId != prod->cteId) - { - continue; - } - - if (!haveFirst) - { - firstConsSlice = cons->sliceId; - haveFirst = true; - } - else if (cons->sliceId != firstConsSlice) - { - cross = true; - break; - } - } - } - - prodInfos->Release(); - consInfos->Release(); - - return cross; -} - //--------------------------------------------------------------------------- // @class: // CUtils::FHasSubqueryOrApply diff --git a/src/backend/gporca/libgpopt/src/translate/CTranslatorExprToDXL.cpp b/src/backend/gporca/libgpopt/src/translate/CTranslatorExprToDXL.cpp index 00d49d57d0c..216024b16e4 100644 --- a/src/backend/gporca/libgpopt/src/translate/CTranslatorExprToDXL.cpp +++ b/src/backend/gporca/libgpopt/src/translate/CTranslatorExprToDXL.cpp @@ -351,20 +351,6 @@ CTranslatorExprToDXL::PdxlnTranslate(CExpression *pexpr, GPOS_ASSERT(NULL == m_pdpplan); - // Walk the physical tree and detect a CTE Consumer placed on a - // different slice than its Producer when the Producer's output is - // replicated-like (StrictReplicated/TaintedReplicated/Universal). - // Fall back to the Postgres optimizer if it is detected because - // it breaks Producer-Consumer locality and can hang the - // query at execution. - if (CUtils::FHasCrossSliceReplicatedCTEConsumer(m_mp, pexpr)) - { - GPOS_RAISE( - gpdxl::ExmaDXL, gpdxl::ExmiExpr2DXLUnsupportedFeature, - GPOS_WSZ_LIT( - "CTE Consumer placed on a different slice than its replicated Producer")); - } - m_pdpplan = CDrvdPropPlan::Pdpplan(pexpr->PdpDerive()); m_pdpplan->AddRef(); From 58c9933340befd16e2b60f1684e8eb0e61d33cbb Mon Sep 17 00:00:00 2001 From: Alena Rybakina Date: Mon, 13 Jul 2026 02:55:29 +0300 Subject: [PATCH 3/3] orca: fall back on cross-slice shared scan with mismatched segment coverage A cross-slice ShareInputScan only executes correctly when its producer and consumer gangs run on the same set of segments: the writer on each producing segment waits for an ack from a reader on that same segment (nodeShareInputScan.c). When the two slices cover different segment sets -- e.g. a DISTRIBUTED REPLICATED CTE materialized on all segments but consumed in a single-segment slice under a Gather 1:1 -- a writer waits for an ack that never arrives and the query hangs to statement_timeout. Detect this on the final plan, where per-slice segment coverage is actually known. apply_shareinput_xslice() already walks the plan with a motId stack; carry a parallel stack of each slice's coverage key (SINGLETON segindex, or "all segments", mirroring how FillSliceTable() sizes an ORCA plan's gangs). When a share is marked cross-slice, compare the producer's and consumer's coverage; if they differ (and the share is not a QD share, which pass 4 relocates onto the QD), set a hazard flag. optimize_query() then discards the ORCA plan and lets the Postgres planner plan the query. --- src/backend/cdb/cdbmutate.c | 108 ++++++++++++++++++ src/backend/optimizer/plan/orca.c | 16 +++ src/include/nodes/relation.h | 9 ++ .../expected/qp_orca_fallback_optimizer.out | 43 ++++--- 4 files changed, 154 insertions(+), 22 deletions(-) diff --git a/src/backend/cdb/cdbmutate.c b/src/backend/cdb/cdbmutate.c index b7e7fbf14ad..1483c5223aa 100644 --- a/src/backend/cdb/cdbmutate.c +++ b/src/backend/cdb/cdbmutate.c @@ -2212,6 +2212,71 @@ shareinput_peekmot(ApplyShareInputContext *ctxt) return linitial_int(ctxt->motStack); } +/* + * Segment-coverage key of a slice, used to detect cross-slice shared scans + * whose producer and consumer gangs run on different segment sets (which + * deadlocks the writer/reader rendezvous in nodeShareInputScan.c). + * + * The key mirrors how FillSliceTable() sizes an ORCA plan's gangs: a slice + * runs on a single process when the flow at the top of the slice is + * FLOW_SINGLETON (segindex -1 = QD, >= 0 = one segment); otherwise ORCA + * dispatches it to all segments. Two slices share the same segment set iff + * their keys are equal. + */ +#define SLICE_COVERAGE_ALLSEG (INT_MIN) /* dispatched to all segments */ +#define SLICE_COVERAGE_UNSET (INT_MIN + 1) /* producer not seen yet */ + +static int +shareinput_slice_coverage(Flow *flow) +{ + if (flow != NULL && flow->flotype == FLOW_SINGLETON) + return flow->segindex; + + return SLICE_COVERAGE_ALLSEG; +} + +/* + * Coverage key of the top slice. A plain SELECT's top slice runs on the QD, + * but when the query is part of a CTAS/COPY/REFRESH, or is a DML on a + * distributed relation, FillSliceTable() turns the top slice into a writer + * gang dispatched to all segments instead. + */ +static int +shareinput_topslice_coverage(PlannerInfo *root) +{ + Query *parse = root->parse; + + if (parse->parentStmtType != PARENTSTMTTYPE_NONE) + return SLICE_COVERAGE_ALLSEG; + + if (parse->commandType != CMD_SELECT && parse->resultRelation > 0) + { + RangeTblEntry *rte = rt_fetch(parse->resultRelation, parse->rtable); + + if (GpPolicyFetch(rte->relid)->ptype != POLICYTYPE_ENTRY) + return SLICE_COVERAGE_ALLSEG; + } + + return -1; /* top slice runs on the QD */ +} + +/* Parallel stack carrying the coverage key of each enclosing slice. */ +static void +shareinput_pushcov(ApplyShareInputContext *ctxt, int cov) +{ + ctxt->covStack = lcons_int(cov, ctxt->covStack); +} +static void +shareinput_popcov(ApplyShareInputContext *ctxt) +{ + ctxt->covStack = list_delete_first(ctxt->covStack); +} +static int +shareinput_peekcov(ApplyShareInputContext *ctxt) +{ + return linitial_int(ctxt->covStack); +} + /* * Replace the target list of ShareInputScan nodes, with references @@ -2373,7 +2438,10 @@ shareinput_mutator_xslice_1(Node *node, PlannerInfo *root, bool fPop) if (fPop) { if (IsA(plan, Motion)) + { shareinput_popmot(ctxt); + shareinput_popcov(ctxt); + } return false; } @@ -2382,6 +2450,10 @@ shareinput_mutator_xslice_1(Node *node, PlannerInfo *root, bool fPop) Motion *motion = (Motion *) plan; shareinput_pushmot(ctxt, motion->motionID); + shareinput_pushcov(ctxt, + shareinput_slice_coverage(motion->plan.lefttree ? + motion->plan.lefttree->flow : + NULL)); return true; } @@ -2411,6 +2483,7 @@ shareinput_mutator_xslice_1(Node *node, PlannerInfo *root, bool fPop) */ ctxt->producers[sisc->share_id] = sisc; ctxt->sliceMarks[sisc->share_id] = motId; + ctxt->producerCoverage[sisc->share_id] = shareinput_peekcov(ctxt); } } @@ -2431,7 +2504,10 @@ shareinput_mutator_xslice_2(Node *node, PlannerInfo *root, bool fPop) if (fPop) { if (IsA(plan, Motion)) + { shareinput_popmot(ctxt); + shareinput_popcov(ctxt); + } return false; } @@ -2440,6 +2516,10 @@ shareinput_mutator_xslice_2(Node *node, PlannerInfo *root, bool fPop) Motion *motion = (Motion *) plan; shareinput_pushmot(ctxt, motion->motionID); + shareinput_pushcov(ctxt, + shareinput_slice_coverage(motion->plan.lefttree ? + motion->plan.lefttree->flow : + NULL)); return true; } @@ -2469,6 +2549,27 @@ shareinput_mutator_xslice_2(Node *node, PlannerInfo *root, bool fPop) incr_plan_nsharer_xslice(plan_slicemark.plan); sisc->driver_slice = motId; + + /* + * A cross-slice share only rendezvouses correctly when the + * producer and this consumer run on the same set of segments: + * the writer on each producing segment waits for an ack from a + * reader on that same segment (nodeShareInputScan.c). If the + * two slices cover different segment sets, a writer waits for + * an ack that never comes, or a reader waits for a producer + * that never runs -- the query hangs. + * + * QD shares are exempt: pass 4 relocates every slice touching + * such a share onto the QD, so they end up consistent. + */ + if (!list_member_int(ctxt->qdShares, sisc->share_id)) + { + int prodCov = ctxt->producerCoverage[sisc->share_id]; + int consCov = shareinput_peekcov(ctxt); + + if (prodCov != SLICE_COVERAGE_UNSET && prodCov != consCov) + ctxt->crossSliceCoverageHazard = true; + } } } } @@ -2592,15 +2693,22 @@ apply_shareinput_xslice(Plan *plan, PlannerInfo *root) PlannerGlobal *glob = root->glob; ApplyShareInputContext *ctxt = &glob->share; ShareInputContext walker_ctxt; + int i; ctxt->motStack = NULL; + ctxt->covStack = NULL; ctxt->qdShares = NULL; ctxt->qdSlices = NULL; ctxt->nextPlanId = 0; + ctxt->crossSliceCoverageHazard = false; ctxt->sliceMarks = palloc0(ctxt->producer_count * sizeof(int)); + ctxt->producerCoverage = palloc(ctxt->producer_count * sizeof(int)); + for (i = 0; i < ctxt->producer_count; i++) + ctxt->producerCoverage[i] = SLICE_COVERAGE_UNSET; shareinput_pushmot(ctxt, 0); + shareinput_pushcov(ctxt, shareinput_topslice_coverage(root)); walker_ctxt.base.node = (Node *) root; diff --git a/src/backend/optimizer/plan/orca.c b/src/backend/optimizer/plan/orca.c index e8cc1950923..78a2f7bb046 100644 --- a/src/backend/optimizer/plan/orca.c +++ b/src/backend/optimizer/plan/orca.c @@ -184,6 +184,22 @@ optimize_query(Query *parse, ParamListInfo boundParams) /* Post-process ShareInputScan nodes */ (void) apply_shareinput_xslice(result->planTree, root); + /* + * apply_shareinput_xslice() flags a cross-slice shared scan whose producer + * and consumer gangs run on different segment sets (e.g. a replicated CTE + * materialized on all segments but consumed in a single-segment slice). + * Such a plan deadlocks at execution in the ShareInputScan writer/reader + * rendezvous, so discard it and let the Postgres planner plan the query. + */ + if (glob->share.crossSliceCoverageHazard) + { + ereport(optimizer_trace_fallback ? INFO : DEBUG1, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("GPORCA failed to produce a plan, falling back to planner"), + errdetail("GPORCA produced a cross-slice shared scan with mismatched segment coverage."))); + return NULL; + } + /* * Fix ShareInputScans for EXPLAIN, like in standard_planner(). For all * subplans first, and then for the main plan tree. diff --git a/src/include/nodes/relation.h b/src/include/nodes/relation.h index b609b3406c2..d3f63dcb315 100644 --- a/src/include/nodes/relation.h +++ b/src/include/nodes/relation.h @@ -88,14 +88,23 @@ typedef struct ApplyShareInputContext int *share_refcounts; int share_refcounts_sz; /* allocated sized of 'share_refcounts' */ List *motStack; + List *covStack; /* parallel to motStack: segment-coverage + * key of each enclosing slice */ List *qdShares; List *qdSlices; int nextPlanId; ShareInputScan **producers; int *sliceMarks; /* one for each producer */ + int *producerCoverage; /* one per producer: coverage key of the + * slice its producer runs in */ int producer_count; + bool crossSliceCoverageHazard; /* a cross-slice shared scan whose + * producer and consumer slices run + * on different segment sets was + * found; ORCA plan must fall back */ + } ApplyShareInputContext; diff --git a/src/test/regress/expected/qp_orca_fallback_optimizer.out b/src/test/regress/expected/qp_orca_fallback_optimizer.out index 18bf2d8a196..47df837089a 100644 --- a/src/test/regress/expected/qp_orca_fallback_optimizer.out +++ b/src/test/regress/expected/qp_orca_fallback_optimizer.out @@ -409,27 +409,26 @@ t2 AS (SELECT id, refrcode FROM tbl2 WHERE REFERENCEID = 101991) JOIN t2 r1 ON p.isCalcDetail = r1.RefrCode LIMIT 1; - QUERY PLAN ----------------------------------------------------------------------------------------------------------------------- - Gather Motion 3:1 (slice3; segments: 3) (cost=0.00..1724.00 rows=1 width=24) - -> Sequence (cost=0.00..1724.00 rows=1 width=24) - -> Shared Scan (share slice:id 3:1) (cost=0.00..431.00 rows=1 width=1) - -> Materialize (cost=0.00..431.00 rows=1 width=1) - -> Seq Scan on tbl2 (cost=0.00..431.00 rows=1 width=263) - Filter: (referenceid = '101991'::numeric) - -> Redistribute Motion 1:3 (slice2) (cost=0.00..1293.00 rows=1 width=24) - -> Limit (cost=0.00..1293.00 rows=1 width=24) - -> Gather Motion 1:1 (slice1; segments: 1) (cost=0.00..1293.00 rows=1 width=24) - -> Hash Join (cost=0.00..1293.00 rows=1 width=24) - Hash Cond: ((tbl1.iscalctrg)::text = (share1_ref2.refrcode)::text) - -> Hash Join (cost=0.00..862.00 rows=1 width=24) - Hash Cond: ((tbl1.iscalcdetail)::text = (share1_ref3.refrcode)::text) - -> Seq Scan on tbl1 (cost=0.00..431.00 rows=1 width=24) - -> Hash (cost=431.00..431.00 rows=1 width=8) - -> Shared Scan (share slice:id 1:1) (cost=0.00..431.00 rows=1 width=8) - -> Hash (cost=431.00..431.00 rows=1 width=8) - -> Shared Scan (share slice:id 1:1) (cost=0.00..431.00 rows=1 width=8) - Optimizer: Pivotal Optimizer (GPORCA) -(19 rows) +INFO: GPORCA failed to produce a plan, falling back to planner +DETAIL: GPORCA produced a cross-slice shared scan with mismatched segment coverage. + QUERY PLAN +--------------------------------------------------------------------------------------------- + Limit (cost=0.34..83.24 rows=1 width=104) + -> Gather Motion 1:1 (slice1; segments: 1) (cost=0.34..581.83 rows=8 width=104) + -> Hash Join (cost=0.34..581.69 rows=8 width=104) + Hash Cond: ((tbl1.iscalcdetail)::text = (r1.refrcode)::text) + -> Hash Join (cost=0.17..580.97 rows=130 width=104) + Hash Cond: ((tbl1.iscalctrg)::text = (r.refrcode)::text) + -> Seq Scan on tbl1 (cost=0.00..344.00 rows=24400 width=104) + -> Hash (cost=0.11..0.11 rows=2 width=516) + -> Subquery Scan on r (cost=0.00..0.11 rows=6 width=516) + -> Seq Scan on tbl2 (cost=0.00..166.25 rows=6 width=548) + Filter: (referenceid = '101991'::numeric) + -> Hash (cost=0.11..0.11 rows=2 width=516) + -> Subquery Scan on r1 (cost=0.00..0.11 rows=6 width=516) + -> Seq Scan on tbl2 tbl2_1 (cost=0.00..166.25 rows=6 width=548) + Filter: (referenceid = '101991'::numeric) + Optimizer: Postgres query optimizer +(16 rows) DROP TABLE tbl1, tbl2;