From 40a36826b654ea2377e77c18c4d06cc4917f62df Mon Sep 17 00:00:00 2001 From: Alexander Korotkov Date: Mon, 30 Mar 2026 01:31:27 +0300 Subject: [PATCH v16 3/3] MergeAppend should support Async Foreign Scan subplans This commit makes the MergeAppend node async-capable, similar to the existing async support for Append nodes. When the planner chooses MergeAppend for partitioned tables with foreign partitions, asynchronous execution is now possible, providing significant performance improvements. A new GUC enable_async_merge_append controls this feature (default on). Unlike Append, which can return async results in any order, MergeAppend must return results in sort order. To handle this, ExecMergeAppendAsyncGetNext requests and caches results per-subplan, only returning them when the binary heap merge requires that specific subplan's tuple. The postgres_fdw is updated to work generically with both Append and MergeAppend requestors via GetAppendEventSet()/GetNeedRequest() helpers. Discussion: https://postgr.es/m/59be194c5a409fb9fc9f2031581b8a44%40postgrespro.ru Author: Alexander Pyhalov Reviewed-by: Matheus Alcantara Reviewed-by: Alena Rybakina --- .../postgres_fdw/expected/postgres_fdw.out | 288 ++++++++++++++++++ contrib/postgres_fdw/postgres_fdw.c | 10 +- contrib/postgres_fdw/sql/postgres_fdw.sql | 87 ++++++ doc/src/sgml/config.sgml | 14 + src/backend/executor/execAsync.c | 4 + src/backend/executor/nodeMergeAppend.c | 279 ++++++++++++++++- src/backend/optimizer/path/costsize.c | 1 + src/backend/optimizer/plan/createplan.c | 9 + .../utils/activity/wait_event_names.txt | 2 +- src/backend/utils/misc/guc_parameters.dat | 8 + src/backend/utils/misc/postgresql.conf.sample | 1 + src/include/executor/nodeMergeAppend.h | 1 + src/include/nodes/execnodes.h | 27 ++ src/include/optimizer/cost.h | 1 + src/test/regress/expected/sysviews.out | 3 +- 15 files changed, 728 insertions(+), 7 deletions(-) diff --git a/contrib/postgres_fdw/expected/postgres_fdw.out b/contrib/postgres_fdw/expected/postgres_fdw.out index 0f5271d476e..e73db70ca69 100644 --- a/contrib/postgres_fdw/expected/postgres_fdw.out +++ b/contrib/postgres_fdw/expected/postgres_fdw.out @@ -11575,6 +11575,46 @@ SELECT * FROM result_tbl ORDER BY a; (2 rows) DELETE FROM result_tbl; +-- Test Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM async_pt WHERE b % 100 = 0 ORDER BY b, a; + QUERY PLAN +------------------------------------------------------------------------------------------------------------------------------ + Merge Append + Sort Key: async_pt.b, async_pt.a + -> Async Foreign Scan on public.async_p1 async_pt_1 + Output: async_pt_1.a, async_pt_1.b, async_pt_1.c + Remote SQL: SELECT a, b, c FROM public.base_tbl1 WHERE (((b % 100) = 0)) ORDER BY b ASC NULLS LAST, a ASC NULLS LAST + -> Async Foreign Scan on public.async_p2 async_pt_2 + Output: async_pt_2.a, async_pt_2.b, async_pt_2.c + Remote SQL: SELECT a, b, c FROM public.base_tbl2 WHERE (((b % 100) = 0)) ORDER BY b ASC NULLS LAST, a ASC NULLS LAST +(8 rows) + +SELECT * FROM async_pt WHERE b % 100 = 0 ORDER BY b, a; + a | b | c +------+-----+------ + 1000 | 0 | 0000 + 2000 | 0 | 0000 + 1100 | 100 | 0100 + 2100 | 100 | 0100 + 1200 | 200 | 0200 + 2200 | 200 | 0200 + 1300 | 300 | 0300 + 2300 | 300 | 0300 + 1400 | 400 | 0400 + 2400 | 400 | 0400 + 1500 | 500 | 0500 + 2500 | 500 | 0500 + 1600 | 600 | 0600 + 2600 | 600 | 0600 + 1700 | 700 | 0700 + 2700 | 700 | 0700 + 1800 | 800 | 0800 + 2800 | 800 | 0800 + 1900 | 900 | 0900 + 2900 | 900 | 0900 +(20 rows) + -- Test error handling, if accessing one of the foreign partitions errors out CREATE FOREIGN TABLE async_p_broken PARTITION OF async_pt FOR VALUES FROM (10000) TO (10001) SERVER loopback OPTIONS (table_name 'non_existent_table'); @@ -11623,6 +11663,76 @@ COPY async_pt TO stdout; --error ERROR: cannot copy from foreign table "async_p1" DETAIL: Partition "async_p1" is a foreign table in partitioned table "async_pt" HINT: Try the COPY (SELECT ...) TO variant. +-- Test Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; + QUERY PLAN +------------------------------------------------------------------------------------------------------ + Merge Append + Sort Key: async_pt.b, async_pt.a + -> Async Foreign Scan on public.async_p1 async_pt_1 + Output: async_pt_1.a, async_pt_1.b, async_pt_1.c + Filter: (async_pt_1.b === 505) + Remote SQL: SELECT a, b, c FROM public.base_tbl1 ORDER BY b ASC NULLS LAST, a ASC NULLS LAST + -> Async Foreign Scan on public.async_p2 async_pt_2 + Output: async_pt_2.a, async_pt_2.b, async_pt_2.c + Filter: (async_pt_2.b === 505) + Remote SQL: SELECT a, b, c FROM public.base_tbl2 ORDER BY b ASC NULLS LAST, a ASC NULLS LAST + -> Async Foreign Scan on public.async_p3 async_pt_3 + Output: async_pt_3.a, async_pt_3.b, async_pt_3.c + Filter: (async_pt_3.b === 505) + Remote SQL: SELECT a, b, c FROM public.base_tbl3 ORDER BY b ASC NULLS LAST, a ASC NULLS LAST +(14 rows) + +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; + a | b | c +------+-----+------ + 1505 | 505 | 0505 + 2505 | 505 | 0505 + 3505 | 505 | 0505 +(3 rows) + +-- Test async Merge Append rescan +EXPLAIN (VERBOSE, COSTS OFF) +SELECT + ARRAY(SELECT f.i FROM (SELECT b + g.i FROM async_pt WHERE a > g.i ORDER BY b) f(i) ORDER BY f.i LIMIT 10) +FROM generate_series(1, 3) g(i); + QUERY PLAN +---------------------------------------------------------------------------------------------------------------------------------- + Function Scan on pg_catalog.generate_series g + Output: ARRAY(SubPlan array_1) + Function Call: generate_series(1, 3) + SubPlan array_1 + -> Limit + Output: f.i + -> Sort + Output: f.i + Sort Key: f.i + -> Subquery Scan on f + Output: f.i + -> Merge Append + Sort Key: async_pt.b + -> Async Foreign Scan on public.async_p1 async_pt_1 + Output: (async_pt_1.b + g.i), async_pt_1.b + Remote SQL: SELECT b FROM public.base_tbl1 WHERE ((a > $1::integer)) ORDER BY b ASC NULLS LAST + -> Async Foreign Scan on public.async_p2 async_pt_2 + Output: (async_pt_2.b + g.i), async_pt_2.b + Remote SQL: SELECT b FROM public.base_tbl2 WHERE ((a > $1::integer)) ORDER BY b ASC NULLS LAST + -> Async Foreign Scan on public.async_p3 async_pt_3 + Output: (async_pt_3.b + g.i), async_pt_3.b + Remote SQL: SELECT b FROM public.base_tbl3 WHERE ((a > $1::integer)) ORDER BY b ASC NULLS LAST +(22 rows) + +SELECT + ARRAY(SELECT f.i FROM (SELECT b + g.i FROM async_pt WHERE a > g.i ORDER BY b) f(i) ORDER BY f.i LIMIT 10) +FROM generate_series(1, 3) g(i); + array +--------------------------- + {1,1,1,6,6,6,11,11,11,16} + {2,2,2,7,7,7,12,12,12,17} + {3,3,3,8,8,8,13,13,13,18} +(3 rows) + DROP FOREIGN TABLE async_p3; DROP TABLE base_tbl3; -- Check case where the partitioned table has local/remote partitions @@ -11658,6 +11768,37 @@ SELECT * FROM result_tbl ORDER BY a; (3 rows) DELETE FROM result_tbl; +-- Test Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; + QUERY PLAN +------------------------------------------------------------------------------------------------------ + Merge Append + Sort Key: async_pt.b, async_pt.a + -> Async Foreign Scan on public.async_p1 async_pt_1 + Output: async_pt_1.a, async_pt_1.b, async_pt_1.c + Filter: (async_pt_1.b === 505) + Remote SQL: SELECT a, b, c FROM public.base_tbl1 ORDER BY b ASC NULLS LAST, a ASC NULLS LAST + -> Async Foreign Scan on public.async_p2 async_pt_2 + Output: async_pt_2.a, async_pt_2.b, async_pt_2.c + Filter: (async_pt_2.b === 505) + Remote SQL: SELECT a, b, c FROM public.base_tbl2 ORDER BY b ASC NULLS LAST, a ASC NULLS LAST + -> Sort + Output: async_pt_3.a, async_pt_3.b, async_pt_3.c + Sort Key: async_pt_3.b, async_pt_3.a + -> Seq Scan on public.async_p3 async_pt_3 + Output: async_pt_3.a, async_pt_3.b, async_pt_3.c + Filter: (async_pt_3.b === 505) +(16 rows) + +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; + a | b | c +------+-----+------ + 1505 | 505 | 0505 + 2505 | 505 | 0505 + 3505 | 505 | 0505 +(3 rows) + -- partitionwise joins SET enable_partitionwise_join TO true; CREATE TABLE join_tbl (a1 int, b1 int, c1 text, a2 int, b2 int, c2 text); @@ -12440,6 +12581,153 @@ SELECT a FROM base_tbl WHERE (a, random() > 0) IN (SELECT a, random() > 0 FROM f DROP FOREIGN TABLE foreign_tbl CASCADE; NOTICE: drop cascades to foreign table foreign_tbl2 DROP TABLE base_tbl; +-- Test async Merge Append +CREATE TABLE distr1 (i int, j int, k text) PARTITION BY HASH (i); +CREATE TABLE base1 (i int, j int, k text); +CREATE TABLE base2 (i int, j int, k text); +CREATE FOREIGN TABLE distr1_p1 PARTITION OF distr1 FOR VALUES WITH (MODULUS 2, REMAINDER 0) +SERVER loopback OPTIONS (table_name 'base1'); +CREATE FOREIGN TABLE distr1_p2 PARTITION OF distr1 FOR VALUES WITH (MODULUS 2, REMAINDER 1) +SERVER loopback OPTIONS (table_name 'base2'); +CREATE TABLE distr2 (i int, j int, k text) PARTITION BY HASH (i); +CREATE TABLE base3 (i int, j int, k text); +CREATE TABLE base4 (i int, j int, k text); +CREATE FOREIGN TABLE distr2_p1 PARTITION OF distr2 FOR VALUES WITH (MODULUS 2, REMAINDER 0) +SERVER loopback OPTIONS (table_name 'base3'); +CREATE FOREIGN TABLE distr2_p2 PARTITION OF distr2 FOR VALUES WITH (MODULUS 2, REMAINDER 1) +SERVER loopback OPTIONS (table_name 'base4'); +INSERT INTO distr1 +SELECT i, i*10, 'data_' || i FROM generate_series(1, 1000) i; +INSERT INTO distr2 +SELECT i, i*10, 'data_' || i FROM generate_series(1, 100) i; +ANALYZE distr1_p1; +ANALYZE distr1_p2; +ANALYZE distr2_p1; +ANALYZE distr2_p2; +SET enable_partitionwise_join TO ON; +-- Test joins with async Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM distr1, distr2 WHERE distr1.i=distr2.i AND distr2.j > 90 and distr2.k like 'data%' +ORDER BY distr2.i LIMIT 10; + QUERY PLAN +------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- + Limit + Output: distr1.i, distr1.j, distr1.k, distr2.i, distr2.j, distr2.k + -> Merge Append + Sort Key: distr1.i + -> Async Foreign Scan + Output: distr1_1.i, distr1_1.j, distr1_1.k, distr2_1.i, distr2_1.j, distr2_1.k + Relations: (public.distr1_p1 distr1_1) INNER JOIN (public.distr2_p1 distr2_1) + Remote SQL: SELECT r3.i, r3.j, r3.k, r5.i, r5.j, r5.k FROM (public.base1 r3 INNER JOIN public.base3 r5 ON (((r3.i = r5.i)) AND ((r5.j > 90)) AND ((r5.k ~~ 'data%')))) ORDER BY r3.i ASC NULLS LAST + -> Async Foreign Scan + Output: distr1_2.i, distr1_2.j, distr1_2.k, distr2_2.i, distr2_2.j, distr2_2.k + Relations: (public.distr1_p2 distr1_2) INNER JOIN (public.distr2_p2 distr2_2) + Remote SQL: SELECT r4.i, r4.j, r4.k, r6.i, r6.j, r6.k FROM (public.base2 r4 INNER JOIN public.base4 r6 ON (((r4.i = r6.i)) AND ((r6.j > 90)) AND ((r6.k ~~ 'data%')))) ORDER BY r4.i ASC NULLS LAST +(12 rows) + +SELECT * FROM distr1, distr2 WHERE distr1.i=distr2.i AND distr2.j > 90 and distr2.k like 'data%' +ORDER BY distr2.i LIMIT 10; + i | j | k | i | j | k +----+-----+---------+----+-----+--------- + 10 | 100 | data_10 | 10 | 100 | data_10 + 11 | 110 | data_11 | 11 | 110 | data_11 + 12 | 120 | data_12 | 12 | 120 | data_12 + 13 | 130 | data_13 | 13 | 130 | data_13 + 14 | 140 | data_14 | 14 | 140 | data_14 + 15 | 150 | data_15 | 15 | 150 | data_15 + 16 | 160 | data_16 | 16 | 160 | data_16 + 17 | 170 | data_17 | 17 | 170 | data_17 + 18 | 180 | data_18 | 18 | 180 | data_18 + 19 | 190 | data_19 | 19 | 190 | data_19 +(10 rows) + +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM distr1 LEFT JOIN distr2 ON distr1.i=distr2.i AND distr2.k like 'data%' WHERE distr1.i > 90 +ORDER BY distr1.i LIMIT 20; + QUERY PLAN +-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- + Limit + Output: distr1.i, distr1.j, distr1.k, distr2.i, distr2.j, distr2.k + -> Merge Append + Sort Key: distr1.i + -> Async Foreign Scan + Output: distr1_1.i, distr1_1.j, distr1_1.k, distr2_1.i, distr2_1.j, distr2_1.k + Relations: (public.distr1_p1 distr1_1) LEFT JOIN (public.distr2_p1 distr2_1) + Remote SQL: SELECT r4.i, r4.j, r4.k, r6.i, r6.j, r6.k FROM (public.base1 r4 LEFT JOIN public.base3 r6 ON (((r4.i = r6.i)) AND ((r6.k ~~ 'data%')))) WHERE ((r4.i > 90)) ORDER BY r4.i ASC NULLS LAST + -> Async Foreign Scan + Output: distr1_2.i, distr1_2.j, distr1_2.k, distr2_2.i, distr2_2.j, distr2_2.k + Relations: (public.distr1_p2 distr1_2) LEFT JOIN (public.distr2_p2 distr2_2) + Remote SQL: SELECT r5.i, r5.j, r5.k, r7.i, r7.j, r7.k FROM (public.base2 r5 LEFT JOIN public.base4 r7 ON (((r5.i = r7.i)) AND ((r7.k ~~ 'data%')))) WHERE ((r5.i > 90)) ORDER BY r5.i ASC NULLS LAST +(12 rows) + +SELECT * FROM distr1 LEFT JOIN distr2 ON distr1.i=distr2.i AND distr2.k like 'data%' WHERE distr1.i > 90 +ORDER BY distr1.i LIMIT 20; + i | j | k | i | j | k +-----+------+----------+-----+------+---------- + 91 | 910 | data_91 | 91 | 910 | data_91 + 92 | 920 | data_92 | 92 | 920 | data_92 + 93 | 930 | data_93 | 93 | 930 | data_93 + 94 | 940 | data_94 | 94 | 940 | data_94 + 95 | 950 | data_95 | 95 | 950 | data_95 + 96 | 960 | data_96 | 96 | 960 | data_96 + 97 | 970 | data_97 | 97 | 970 | data_97 + 98 | 980 | data_98 | 98 | 980 | data_98 + 99 | 990 | data_99 | 99 | 990 | data_99 + 100 | 1000 | data_100 | 100 | 1000 | data_100 + 101 | 1010 | data_101 | | | + 102 | 1020 | data_102 | | | + 103 | 1030 | data_103 | | | + 104 | 1040 | data_104 | | | + 105 | 1050 | data_105 | | | + 106 | 1060 | data_106 | | | + 107 | 1070 | data_107 | | | + 108 | 1080 | data_108 | | | + 109 | 1090 | data_109 | | | + 110 | 1100 | data_110 | | | +(20 rows) + +-- Test pruning with async Merge Append +DELETE FROM distr2; +INSERT INTO distr2 +SELECT i%10, i*10, 'data_' || i FROM generate_series(1, 1000) i; +DEALLOCATE ALL; +SET plan_cache_mode TO force_generic_plan; +PREPARE async_pt_query (int, int) AS + SELECT * FROM distr2 WHERE i = ANY(ARRAY[$1, $2]) + ORDER BY i,j + LIMIT 10; +EXPLAIN (VERBOSE, COSTS OFF) + EXECUTE async_pt_query(1, 1); + QUERY PLAN +------------------------------------------------------------------------------------------------------------------------------------------------------------ + Limit + Output: distr2.i, distr2.j, distr2.k + -> Merge Append + Sort Key: distr2.i, distr2.j + Subplans Removed: 1 + -> Async Foreign Scan on public.distr2_p1 distr2_1 + Output: distr2_1.i, distr2_1.j, distr2_1.k + Remote SQL: SELECT i, j, k FROM public.base3 WHERE ((i = ANY (ARRAY[$1::integer, $2::integer]))) ORDER BY i ASC NULLS LAST, j ASC NULLS LAST +(8 rows) + +EXECUTE async_pt_query(1, 1); + i | j | k +---+-----+--------- + 1 | 10 | data_1 + 1 | 110 | data_11 + 1 | 210 | data_21 + 1 | 310 | data_31 + 1 | 410 | data_41 + 1 | 510 | data_51 + 1 | 610 | data_61 + 1 | 710 | data_71 + 1 | 810 | data_81 + 1 | 910 | data_91 +(10 rows) + +RESET plan_cache_mode; +RESET enable_partitionwise_join; +DROP TABLE distr1, distr2, base1, base2, base3, base4; ALTER SERVER loopback OPTIONS (DROP async_capable); ALTER SERVER loopback2 OPTIONS (DROP async_capable); -- =================================================================== diff --git a/contrib/postgres_fdw/postgres_fdw.c b/contrib/postgres_fdw/postgres_fdw.c index 7416d09c7e2..b7444fdb22e 100644 --- a/contrib/postgres_fdw/postgres_fdw.c +++ b/contrib/postgres_fdw/postgres_fdw.c @@ -7214,12 +7214,16 @@ postgresForeignAsyncConfigureWait(AsyncRequest *areq) ForeignScanState *node = (ForeignScanState *) areq->requestee; PgFdwScanState *fsstate = (PgFdwScanState *) node->fdw_state; AsyncRequest *pendingAreq = fsstate->conn_state->pendingAreq; - AppendState *requestor = (AppendState *) areq->requestor; - WaitEventSet *set = requestor->as.eventset; + PlanState *requestor = areq->requestor; + WaitEventSet *set; + Bitmapset *needrequest; /* This should not be called unless callback_pending */ Assert(areq->callback_pending); + set = GetAppendEventSet(requestor); + needrequest = GetNeedRequest(requestor); + /* * If process_pending_request() has been invoked on the given request * before we get here, we might have some tuples already; in which case @@ -7257,7 +7261,7 @@ postgresForeignAsyncConfigureWait(AsyncRequest *areq) * below, because we might otherwise end up with no configured events * other than the postmaster death event. */ - if (!bms_is_empty(requestor->as.needrequest)) + if (!bms_is_empty(needrequest)) return; if (GetNumRegisteredWaitEvents(set) > 1) return; diff --git a/contrib/postgres_fdw/sql/postgres_fdw.sql b/contrib/postgres_fdw/sql/postgres_fdw.sql index 49ed797e8ef..73fc5317ad5 100644 --- a/contrib/postgres_fdw/sql/postgres_fdw.sql +++ b/contrib/postgres_fdw/sql/postgres_fdw.sql @@ -3943,6 +3943,11 @@ INSERT INTO result_tbl SELECT a, b, 'AAA' || c FROM async_pt WHERE b === 505; SELECT * FROM result_tbl ORDER BY a; DELETE FROM result_tbl; +-- Test Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM async_pt WHERE b % 100 = 0 ORDER BY b, a; +SELECT * FROM async_pt WHERE b % 100 = 0 ORDER BY b, a; + -- Test error handling, if accessing one of the foreign partitions errors out CREATE FOREIGN TABLE async_p_broken PARTITION OF async_pt FOR VALUES FROM (10000) TO (10001) SERVER loopback OPTIONS (table_name 'non_existent_table'); @@ -3966,6 +3971,20 @@ DELETE FROM result_tbl; -- Test COPY TO when foreign table is partition COPY async_pt TO stdout; --error +-- Test Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; + +-- Test async Merge Append rescan +EXPLAIN (VERBOSE, COSTS OFF) +SELECT + ARRAY(SELECT f.i FROM (SELECT b + g.i FROM async_pt WHERE a > g.i ORDER BY b) f(i) ORDER BY f.i LIMIT 10) +FROM generate_series(1, 3) g(i); +SELECT + ARRAY(SELECT f.i FROM (SELECT b + g.i FROM async_pt WHERE a > g.i ORDER BY b) f(i) ORDER BY f.i LIMIT 10) +FROM generate_series(1, 3) g(i); + DROP FOREIGN TABLE async_p3; DROP TABLE base_tbl3; @@ -3981,6 +4000,11 @@ INSERT INTO result_tbl SELECT * FROM async_pt WHERE b === 505; SELECT * FROM result_tbl ORDER BY a; DELETE FROM result_tbl; +-- Test Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; +SELECT * FROM async_pt WHERE b === 505 ORDER BY b, a; + -- partitionwise joins SET enable_partitionwise_join TO true; @@ -4219,6 +4243,69 @@ SELECT a FROM base_tbl WHERE (a, random() > 0) IN (SELECT a, random() > 0 FROM f DROP FOREIGN TABLE foreign_tbl CASCADE; DROP TABLE base_tbl; +-- Test async Merge Append +CREATE TABLE distr1 (i int, j int, k text) PARTITION BY HASH (i); +CREATE TABLE base1 (i int, j int, k text); +CREATE TABLE base2 (i int, j int, k text); +CREATE FOREIGN TABLE distr1_p1 PARTITION OF distr1 FOR VALUES WITH (MODULUS 2, REMAINDER 0) +SERVER loopback OPTIONS (table_name 'base1'); +CREATE FOREIGN TABLE distr1_p2 PARTITION OF distr1 FOR VALUES WITH (MODULUS 2, REMAINDER 1) +SERVER loopback OPTIONS (table_name 'base2'); + +CREATE TABLE distr2 (i int, j int, k text) PARTITION BY HASH (i); +CREATE TABLE base3 (i int, j int, k text); +CREATE TABLE base4 (i int, j int, k text); +CREATE FOREIGN TABLE distr2_p1 PARTITION OF distr2 FOR VALUES WITH (MODULUS 2, REMAINDER 0) +SERVER loopback OPTIONS (table_name 'base3'); +CREATE FOREIGN TABLE distr2_p2 PARTITION OF distr2 FOR VALUES WITH (MODULUS 2, REMAINDER 1) +SERVER loopback OPTIONS (table_name 'base4'); + +INSERT INTO distr1 +SELECT i, i*10, 'data_' || i FROM generate_series(1, 1000) i; + +INSERT INTO distr2 +SELECT i, i*10, 'data_' || i FROM generate_series(1, 100) i; + +ANALYZE distr1_p1; +ANALYZE distr1_p2; +ANALYZE distr2_p1; +ANALYZE distr2_p2; + +SET enable_partitionwise_join TO ON; + +-- Test joins with async Merge Append +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM distr1, distr2 WHERE distr1.i=distr2.i AND distr2.j > 90 and distr2.k like 'data%' +ORDER BY distr2.i LIMIT 10; +SELECT * FROM distr1, distr2 WHERE distr1.i=distr2.i AND distr2.j > 90 and distr2.k like 'data%' +ORDER BY distr2.i LIMIT 10; + +EXPLAIN (VERBOSE, COSTS OFF) +SELECT * FROM distr1 LEFT JOIN distr2 ON distr1.i=distr2.i AND distr2.k like 'data%' WHERE distr1.i > 90 +ORDER BY distr1.i LIMIT 20; +SELECT * FROM distr1 LEFT JOIN distr2 ON distr1.i=distr2.i AND distr2.k like 'data%' WHERE distr1.i > 90 +ORDER BY distr1.i LIMIT 20; + +-- Test pruning with async Merge Append +DELETE FROM distr2; +INSERT INTO distr2 +SELECT i%10, i*10, 'data_' || i FROM generate_series(1, 1000) i; + +DEALLOCATE ALL; +SET plan_cache_mode TO force_generic_plan; +PREPARE async_pt_query (int, int) AS + SELECT * FROM distr2 WHERE i = ANY(ARRAY[$1, $2]) + ORDER BY i,j + LIMIT 10; +EXPLAIN (VERBOSE, COSTS OFF) + EXECUTE async_pt_query(1, 1); +EXECUTE async_pt_query(1, 1); +RESET plan_cache_mode; + +RESET enable_partitionwise_join; + +DROP TABLE distr1, distr2, base1, base2, base3, base4; + ALTER SERVER loopback OPTIONS (DROP async_capable); ALTER SERVER loopback2 OPTIONS (DROP async_capable); diff --git a/doc/src/sgml/config.sgml b/doc/src/sgml/config.sgml index 229f41353eb..dc20fe770df 100644 --- a/doc/src/sgml/config.sgml +++ b/doc/src/sgml/config.sgml @@ -5556,6 +5556,20 @@ ANY num_sync ( + enable_async_merge_append (boolean) + + enable_async_merge_append configuration parameter + + + + + Enables or disables the query planner's use of async-aware + merge append plan types. The default is on. + + + + enable_bitmapscan (boolean) diff --git a/src/backend/executor/execAsync.c b/src/backend/executor/execAsync.c index cf7ddbb01f4..f839f5f255c 100644 --- a/src/backend/executor/execAsync.c +++ b/src/backend/executor/execAsync.c @@ -18,6 +18,7 @@ #include "executor/executor.h" #include "executor/instrument.h" #include "executor/nodeAppend.h" +#include "executor/nodeMergeAppend.h" #include "executor/nodeForeignscan.h" /* @@ -122,6 +123,9 @@ ExecAsyncResponse(AsyncRequest *areq) case T_AppendState: ExecAsyncAppendResponse(areq); break; + case T_MergeAppendState: + ExecAsyncMergeAppendResponse(areq); + break; default: /* If the node doesn't support async, caller messed up. */ elog(ERROR, "unrecognized node type: %d", diff --git a/src/backend/executor/nodeMergeAppend.c b/src/backend/executor/nodeMergeAppend.c index cd03b2bc7f8..3f0d7fb3a45 100644 --- a/src/backend/executor/nodeMergeAppend.c +++ b/src/backend/executor/nodeMergeAppend.c @@ -40,11 +40,16 @@ #include "executor/execAppend.h" #include "executor/executor.h" -#include "executor/execPartition.h" +#include "executor/execAsync.h" #include "executor/nodeMergeAppend.h" +#include "executor/execPartition.h" #include "lib/binaryheap.h" #include "miscadmin.h" +#include "storage/latch.h" #include "utils/sortsupport.h" +#include "utils/wait_event.h" + +#define EVENT_BUFFER_SIZE 16 /* * We have one slot for each item in the heap array. We use SlotNumber @@ -56,6 +61,12 @@ typedef int32 SlotNumber; static TupleTableSlot *ExecMergeAppend(PlanState *pstate); static int heap_compare_slots(Datum a, Datum b, void *arg); +static void classify_matching_subplans(MergeAppendState *node); +static void ExecMergeAppendAsyncBegin(MergeAppendState *node); +static void ExecMergeAppendAsyncGetNext(MergeAppendState *node, int mplan); +static bool ExecMergeAppendAsyncRequest(MergeAppendState *node, int mplan); +static void ExecMergeAppendAsyncEventWait(MergeAppendState *node); + /* ---------------------------------------------------------------- * ExecInitMergeAppend @@ -87,10 +98,16 @@ ExecInitMergeAppend(MergeAppend *node, EState *estate, int eflags) -1, NULL); + if (mergestate->ms.nasyncplans > 0 && mergestate->ms.valid_subplans_identified) + classify_matching_subplans(mergestate); + mergestate->ms_slots = palloc0_array(TupleTableSlot *, mergestate->ms.nplans); mergestate->ms_heap = binaryheap_allocate(mergestate->ms.nplans, heap_compare_slots, mergestate); + mergestate->ms_has_asyncresults = NULL; + mergestate->ms_asyncremain = NULL; + /* * initialize sort-key information */ @@ -153,8 +170,13 @@ ExecMergeAppend(PlanState *pstate) node->ms.valid_subplans = ExecFindMatchingSubPlans(node->ms.prune_state, false, NULL); node->ms.valid_subplans_identified = true; + classify_matching_subplans(node); } + /* If there are any async subplans, begin executing them. */ + if (node->ms.nasyncplans > 0) + ExecMergeAppendAsyncBegin(node); + /* * First time through: pull the first tuple from each valid subplan, * and set up the heap. @@ -166,6 +188,16 @@ ExecMergeAppend(PlanState *pstate) if (!TupIsNull(node->ms_slots[i])) binaryheap_add_unordered(node->ms_heap, Int32GetDatum(i)); } + + /* Look at valid async subplans */ + i = -1; + while ((i = bms_next_member(node->ms.valid_asyncplans, i)) >= 0) + { + ExecMergeAppendAsyncGetNext(node, i); + if (!TupIsNull(node->ms_slots[i])) + binaryheap_add_unordered(node->ms_heap, Int32GetDatum(i)); + } + binaryheap_build(node->ms_heap); node->ms_initialized = true; } @@ -180,7 +212,13 @@ ExecMergeAppend(PlanState *pstate) * to not pull tuples until necessary.) */ i = DatumGetInt32(binaryheap_first(node->ms_heap)); - node->ms_slots[i] = ExecProcNode(node->ms.plans[i]); + if (bms_is_member(i, node->ms.asyncplans)) + ExecMergeAppendAsyncGetNext(node, i); + else + { + Assert(bms_is_member(i, node->ms.valid_subplans)); + node->ms_slots[i] = ExecProcNode(node->ms.plans[i]); + } if (!TupIsNull(node->ms_slots[i])) binaryheap_replace_first(node->ms_heap, Int32GetDatum(i)); else @@ -196,6 +234,8 @@ ExecMergeAppend(PlanState *pstate) { i = DatumGetInt32(binaryheap_first(node->ms_heap)); result = node->ms_slots[i]; + /* For async plan record that we can get the next tuple */ + node->ms_has_asyncresults = bms_del_member(node->ms_has_asyncresults, i); } return result; @@ -260,8 +300,243 @@ ExecEndMergeAppend(MergeAppendState *node) void ExecReScanMergeAppend(MergeAppendState *node) { + int nasyncplans = node->ms.nasyncplans; + ExecReScanAppender(&node->ms); + /* Reset specific merge append async state */ + if (nasyncplans > 0) + { + bms_free(node->ms_asyncremain); + node->ms_asyncremain = NULL; + bms_free(node->ms_has_asyncresults); + node->ms_has_asyncresults = NULL; + } binaryheap_reset(node->ms_heap); node->ms_initialized = false; } + +/* ---------------------------------------------------------------- + * classify_matching_subplans + * + * Classify the node's ms_valid_subplans into sync ones and + * async ones, adjust it to contain sync ones only, and save + * async ones in the node's ms_valid_asyncplans. + * ---------------------------------------------------------------- + */ +static void +classify_matching_subplans(MergeAppendState *node) +{ + Assert(node->ms.valid_subplans_identified); + + /* Nothing to do if there are no valid subplans. */ + if (bms_is_empty(node->ms.valid_subplans)) + { + node->ms_asyncremain = NULL; + return; + } + + /* No valid async subplans identified. */ + if (!classify_matching_subplans_common( + &node->ms.valid_subplans, + node->ms.asyncplans, + &node->ms.valid_asyncplans)) + node->ms_asyncremain = NULL; +} + +/* ---------------------------------------------------------------- + * ExecMergeAppendAsyncBegin + * + * Begin executing designed async-capable subplans. + * ---------------------------------------------------------------- + */ +static void +ExecMergeAppendAsyncBegin(MergeAppendState *node) +{ + /* ExecMergeAppend() identifies valid subplans */ + Assert(node->ms.valid_subplans_identified); + + /* Initialize state variables. */ + node->ms_asyncremain = bms_copy(node->ms.valid_asyncplans); + + /* Nothing to do if there are no valid async subplans. */ + if (bms_is_empty(node->ms_asyncremain)) + return; + + ExecAppenderAsyncBegin(&node->ms); +} + +/* ---------------------------------------------------------------- + * ExecMergeAppendAsyncGetNext + * + * Get the next tuple from specified asynchronous subplan. + * ---------------------------------------------------------------- + */ +static void +ExecMergeAppendAsyncGetNext(MergeAppendState *node, int mplan) +{ + node->ms_slots[mplan] = NULL; + + /* Request a tuple asynchronously. */ + if (ExecMergeAppendAsyncRequest(node, mplan)) + return; + + /* + * node->ms_asyncremain can be NULL if we have fetched tuples, but haven't + * returned them yet. In this case ExecMergeAppendAsyncRequest() above + * just returns tuples without performing a request. + */ + while (bms_is_member(mplan, node->ms_asyncremain)) + { + CHECK_FOR_INTERRUPTS(); + + /* Wait or poll for async events. */ + ExecMergeAppendAsyncEventWait(node); + + /* Request a tuple asynchronously. */ + if (ExecMergeAppendAsyncRequest(node, mplan)) + return; + + /* + * Waiting until there's no async requests pending or we got some + * tuples from our request + */ + } + + /* No tuples */ + return; +} + +/* ---------------------------------------------------------------- + * ExecMergeAppendAsyncRequest + * + * Request a tuple asynchronously. + * ---------------------------------------------------------------- + */ +static bool +ExecMergeAppendAsyncRequest(MergeAppendState *node, int mplan) +{ + Bitmapset *needrequest; + int i; + + /* + * If we've already fetched necessary data, just return it + */ + if (bms_is_member(mplan, node->ms_has_asyncresults)) + { + node->ms_slots[mplan] = node->ms.asyncresults[mplan]; + return true; + } + + /* + * Get a list of members which can process request and don't have data + * ready. + */ + needrequest = NULL; + i = -1; + while ((i = bms_next_member(node->ms.needrequest, i)) >= 0) + { + if (!bms_is_member(i, node->ms_has_asyncresults)) + needrequest = bms_add_member(needrequest, i); + } + + /* + * If there's no members, which still need request, no need to send it. + */ + if (bms_is_empty(needrequest)) + return false; + + /* Clear ms_needrequest flag, as we are going to send requests now */ + node->ms.needrequest = bms_del_members(node->ms.needrequest, needrequest); + + /* Make a new request for each of the async subplans that need it. */ + i = -1; + while ((i = bms_next_member(needrequest, i)) >= 0) + { + AsyncRequest *areq = node->ms.asyncrequests[i]; + + /* + * We've just checked that subplan doesn't already have some fetched + * data + */ + Assert(!bms_is_member(i, node->ms_has_asyncresults)); + + /* Do the actual work. */ + ExecAsyncRequest(areq); + } + bms_free(needrequest); + + /* Return needed asynchronously-generated results if any. */ + if (bms_is_member(mplan, node->ms_has_asyncresults)) + { + node->ms_slots[mplan] = node->ms.asyncresults[mplan]; + return true; + } + + return false; +} + +/* ---------------------------------------------------------------- + * ExecAsyncMergeAppendResponse + * + * Receive a response from an asynchronous request we made. + * ---------------------------------------------------------------- + */ +void +ExecAsyncMergeAppendResponse(AsyncRequest *areq) +{ + MergeAppendState *node = (MergeAppendState *) areq->requestor; + TupleTableSlot *slot = areq->result; + + /* The result should be a TupleTableSlot or NULL. */ + Assert(slot == NULL || IsA(slot, TupleTableSlot)); + /* We should handle previous async result prior to getting new one */ + Assert(!bms_is_member(areq->request_index, node->ms_has_asyncresults)); + + node->ms.asyncresults[areq->request_index] = NULL; + /* Nothing to do if the request is pending. */ + if (!areq->request_complete) + { + /* The request would have been pending for a callback. */ + Assert(areq->callback_pending); + return; + } + + /* If the result is NULL or an empty slot, there's nothing more to do. */ + if (TupIsNull(slot)) + { + /* The ending subplan wouldn't have been pending for a callback. */ + Assert(!areq->callback_pending); + node->ms_asyncremain = bms_del_member(node->ms_asyncremain, + areq->request_index); + return; + } + + /* Mark that the async request has a result */ + node->ms_has_asyncresults = bms_add_member(node->ms_has_asyncresults, + areq->request_index); + /* Save result so we can return it. */ + node->ms.asyncresults[areq->request_index] = slot; + + /* + * Mark the subplan that returned a result as ready for a new request. We + * don't launch another one here immediately because it might complete. + */ + node->ms.needrequest = bms_add_member(node->ms.needrequest, + areq->request_index); +} + +/* ---------------------------------------------------------------- + * ExecMergeAppendAsyncEventWait + * + * Wait or poll for file descriptor events and fire callbacks. + * ---------------------------------------------------------------- + */ +static void +ExecMergeAppendAsyncEventWait(MergeAppendState *node) +{ + /* We should never be called when there are no valid async subplans. */ + Assert(bms_num_members(node->ms_asyncremain) > 0); + + ExecAppenderAsyncEventWait(&node->ms, -1 /* no timeout */ , WAIT_EVENT_APPEND_READY); +} diff --git a/src/backend/optimizer/path/costsize.c b/src/backend/optimizer/path/costsize.c index 1c575e56ff6..b6109e5b91e 100644 --- a/src/backend/optimizer/path/costsize.c +++ b/src/backend/optimizer/path/costsize.c @@ -164,6 +164,7 @@ bool enable_parallel_hash = true; bool enable_partition_pruning = true; bool enable_presorted_aggregate = true; bool enable_async_append = true; +bool enable_async_merge_append = true; typedef struct { diff --git a/src/backend/optimizer/plan/createplan.c b/src/backend/optimizer/plan/createplan.c index b23643fedb8..b340d1dfdd6 100644 --- a/src/backend/optimizer/plan/createplan.c +++ b/src/backend/optimizer/plan/createplan.c @@ -1466,6 +1466,7 @@ create_merge_append_plan(PlannerInfo *root, MergeAppendPath *best_path, List *subplans = NIL; ListCell *subpaths; RelOptInfo *rel = best_path->path.parent; + bool consider_async = false; /* * We don't have the actual creation of the MergeAppend node split out @@ -1481,6 +1482,10 @@ create_merge_append_plan(PlannerInfo *root, MergeAppendPath *best_path, node->ap.apprelids = rel->relids; node->ap.child_append_relid_sets = best_path->child_append_relid_sets; + consider_async = (enable_async_merge_append && + !best_path->path.parallel_safe && + list_length(best_path->subpaths) > 1); + /* * Compute sort column info, and adjust MergeAppend's tlist as needed. * Because we pass adjust_tlist_in_place = true, we may ignore the @@ -1581,6 +1586,10 @@ create_merge_append_plan(PlannerInfo *root, MergeAppendPath *best_path, subplan = sort_plan; } + /* If needed, check to see if subplan can be executed asynchronously */ + if (consider_async) + mark_async_capable_plan(subplan, subpath); + subplans = lappend(subplans, subplan); } diff --git a/src/backend/utils/activity/wait_event_names.txt b/src/backend/utils/activity/wait_event_names.txt index 6be80d2daad..141157ce51d 100644 --- a/src/backend/utils/activity/wait_event_names.txt +++ b/src/backend/utils/activity/wait_event_names.txt @@ -106,7 +106,7 @@ ABI_compatibility: Section: ClassName - WaitEventIPC -APPEND_READY "Waiting for subplan nodes of an Append plan node to be ready." +APPEND_READY "Waiting for subplan nodes of an Append or a MergeAppend plan node to be ready." ARCHIVE_CLEANUP_COMMAND "Waiting for to complete." ARCHIVE_COMMAND "Waiting for to complete." BACKEND_TERMINATION "Waiting for the termination of another backend." diff --git a/src/backend/utils/misc/guc_parameters.dat b/src/backend/utils/misc/guc_parameters.dat index 0a862693fcd..9848964b024 100644 --- a/src/backend/utils/misc/guc_parameters.dat +++ b/src/backend/utils/misc/guc_parameters.dat @@ -861,6 +861,14 @@ boot_val => 'true', }, +{ name => 'enable_async_merge_append', type => 'bool', context => 'PGC_USERSET', group => 'QUERY_TUNING_METHOD', + short_desc => 'Enables the planner\'s use of async merge append plans.', + flags => 'GUC_EXPLAIN', + variable => 'enable_async_merge_append', + boot_val => 'true', +}, + + { name => 'enable_bitmapscan', type => 'bool', context => 'PGC_USERSET', group => 'QUERY_TUNING_METHOD', short_desc => 'Enables the planner\'s use of bitmap-scan plans.', flags => 'GUC_EXPLAIN', diff --git a/src/backend/utils/misc/postgresql.conf.sample b/src/backend/utils/misc/postgresql.conf.sample index cf15597385b..9b8de8581bc 100644 --- a/src/backend/utils/misc/postgresql.conf.sample +++ b/src/backend/utils/misc/postgresql.conf.sample @@ -413,6 +413,7 @@ # - Planner Method Configuration - #enable_async_append = on +#enable_async_merge_append = on #enable_bitmapscan = on #enable_gathermerge = on #enable_hashagg = on diff --git a/src/include/executor/nodeMergeAppend.h b/src/include/executor/nodeMergeAppend.h index dfcf45099ba..2255cc68b21 100644 --- a/src/include/executor/nodeMergeAppend.h +++ b/src/include/executor/nodeMergeAppend.h @@ -19,5 +19,6 @@ extern MergeAppendState *ExecInitMergeAppend(MergeAppend *node, EState *estate, int eflags); extern void ExecEndMergeAppend(MergeAppendState *node); extern void ExecReScanMergeAppend(MergeAppendState *node); +extern void ExecAsyncMergeAppendResponse(AsyncRequest *areq); #endif /* NODEMERGEAPPEND_H */ diff --git a/src/include/nodes/execnodes.h b/src/include/nodes/execnodes.h index 0b93c004727..9e311a9b79c 100644 --- a/src/include/nodes/execnodes.h +++ b/src/include/nodes/execnodes.h @@ -1575,8 +1575,35 @@ typedef struct MergeAppendState TupleTableSlot **ms_slots; /* array of length ms_nplans */ struct binaryheap *ms_heap; /* binary heap of slot indices */ bool ms_initialized; /* are subplans started? */ + + /* Merge-specific async tracking */ + Bitmapset *ms_has_asyncresults; /* plans which have async results */ + Bitmapset *ms_asyncremain; /* remaining asynchronous plans */ } MergeAppendState; +/* Getters for AppendState and MergeAppendState */ +static inline struct WaitEventSet * +GetAppendEventSet(PlanState *ps) +{ + Assert(IsA(ps, AppendState) || IsA(ps, MergeAppendState)); + + if (IsA(ps, AppendState)) + return ((AppendState *) ps)->as.eventset; + else + return ((MergeAppendState *) ps)->ms.eventset; +} + +static inline Bitmapset * +GetNeedRequest(PlanState *ps) +{ + Assert(IsA(ps, AppendState) || IsA(ps, MergeAppendState)); + + if (IsA(ps, AppendState)) + return ((AppendState *) ps)->as.needrequest; + else + return ((MergeAppendState *) ps)->ms.needrequest; +} + /* ---------------- * RecursiveUnionState information * diff --git a/src/include/optimizer/cost.h b/src/include/optimizer/cost.h index f2fd5d31507..798af1fcd5c 100644 --- a/src/include/optimizer/cost.h +++ b/src/include/optimizer/cost.h @@ -70,6 +70,7 @@ extern PGDLLIMPORT bool enable_parallel_hash; extern PGDLLIMPORT bool enable_partition_pruning; extern PGDLLIMPORT bool enable_presorted_aggregate; extern PGDLLIMPORT bool enable_async_append; +extern PGDLLIMPORT bool enable_async_merge_append; extern PGDLLIMPORT int constraint_exclusion; extern double index_pages_fetched(double tuples_fetched, BlockNumber pages, diff --git a/src/test/regress/expected/sysviews.out b/src/test/regress/expected/sysviews.out index 132b56a5864..422ca8b7d1f 100644 --- a/src/test/regress/expected/sysviews.out +++ b/src/test/regress/expected/sysviews.out @@ -156,6 +156,7 @@ select name, setting from pg_settings where name like 'enable%'; name | setting --------------------------------+--------- enable_async_append | on + enable_async_merge_append | on enable_bitmapscan | on enable_distinct_reordering | on enable_eager_aggregate | on @@ -180,7 +181,7 @@ select name, setting from pg_settings where name like 'enable%'; enable_seqscan | on enable_sort | on enable_tidscan | on -(25 rows) +(26 rows) -- There are always wait event descriptions for various types. InjectionPoint -- may be present or absent, depending on history since last postmaster start. -- 2.39.5 (Apple Git-154)