This is an automated email from the ASF dual-hosted git repository.
leborchuk pushed a commit to branch REL_2_STABLE
in repository https://gitbox.apache.org/repos/asf/cloudberry.git
The following commit(s) were added to refs/heads/REL_2_STABLE by this push:
new 205a7a942e0 Fix ndistinct-by-segments for partitioned tables (#2027)
205a7a942e0 is described below
commit 205a7a942e0b923c4036d18e71cfb953beecac3a
Author: Alena Rybakina <[email protected]>
AuthorDate: Mon Sep 28 17:09:34 2026 +0300
Fix ndistinct-by-segments for partitioned tables (#2027)
When ANALYZE builds statistics for a partitioned table, it adds up the
per-segment ndistinct of all partitions. If the same values appear in
every partition, they are counted once per partition, so the result is
too big: 8 partitions give 8 times the real value.
ORCA uses this number to estimate how many rows a partial aggregate
returns. With the inflated number it thinks the partial aggregate
removes almost no rows, and picks a one-stage aggregate that sends all
rows through the motion.
A segment cannot have more distinct values than the whole table, so
cap the value at the table's ndistinct times the number of segments.
---
src/backend/commands/analyze.c | 44 ++++++++++++-
src/backend/commands/analyzeutils.c | 16 ++++-
src/include/commands/analyzeutils.h | 3 +-
src/test/regress/expected/incremental_analyze.out | 78 +++++++++++++++++++++++
src/test/regress/sql/incremental_analyze.sql | 48 ++++++++++++++
5 files changed, 186 insertions(+), 3 deletions(-)
diff --git a/src/backend/commands/analyze.c b/src/backend/commands/analyze.c
index fbd26f2ee9c..2630c4943f8 100644
--- a/src/backend/commands/analyze.c
+++ b/src/backend/commands/analyze.c
@@ -120,6 +120,7 @@
#include "utils/timestamp.h"
#include "access/appendonlywriter.h"
+#include "catalog/gp_distribution_policy.h"
#include "catalog/heap.h"
#include "catalog/pg_am.h"
#include "cdb/cdbappendonlyam.h"
@@ -4987,13 +4988,54 @@ merge_leaf_stats(VacAttrStatsP stats,
old_context = MemoryContextSwitchTo(stats->anl_context);
bool valid;
double ndinstinct_by_segs = 0;
+ double leaf_ndistinct_sum = 0;
Datum *ndvbs;
valid = aggregate_leaf_partition_ndvbs(
- numPartitions, heaptupleStats, relTuples,
&ndinstinct_by_segs);
+ numPartitions, heaptupleStats, relTuples,
&ndinstinct_by_segs,
+ &leaf_ndistinct_sum);
if (valid)
{
+ /*
+ * The leaves' values are summed, which is only right
when leaves hold
+ * disjoint values, e.g. for the partitioning key. A
value repeated in
+ * every partition is counted once per partition, so
the sum grows
+ * with the number of partitions and ORCA overestimates
the output of
+ * a local aggregate.
+ *
+ * This statistic counts a distinct value once per
segment holding
+ * it, so it is the column's ndistinct times the
average number of
+ * segments a value sits on. That average does not
depend on the
+ * partitioning, so take it from the leaves, where the
sums are free
+ * of the double counting, and apply it to the root's
ndistinct.
+ * For a partitioning key, whose values belong to one
leaf each, the
+ * result is the plain sum as before.
+ *
+ * The average is between one and the number of
segments. Keeping
+ * the sum as an upper bound guards against the root's
ndistinct and
+ * the leaves' one disagreeing, as they are estimated
separately.
+ */
+ double root_ndistinct = stats->stadistinct < 0
?
+ -stats->stadistinct * totalTuples :
stats->stadistinct;
+
+ if (root_ndistinct > 0)
+ {
+ GpPolicy *policy =
GpPolicyFetch(stats->attr->attrelid);
+ int numsegments =
policy->numsegments > 0 ?
+ policy->numsegments :
getgpsegmentCount();
+ double segs_per_value;
+
+ pfree(policy);
+
+ segs_per_value = leaf_ndistinct_sum > 0 ?
+ ndinstinct_by_segs / leaf_ndistinct_sum
: numsegments;
+ segs_per_value = Min(Max(segs_per_value, 1.0),
numsegments);
+
+ ndinstinct_by_segs = Min(ndinstinct_by_segs,
+
root_ndistinct * segs_per_value);
+ }
+
ndvbs = (Datum *) palloc(sizeof(Datum));
ndvbs[0] = Float8GetDatum(ndinstinct_by_segs);
diff --git a/src/backend/commands/analyzeutils.c
b/src/backend/commands/analyzeutils.c
index c4f9fa5dedb..42335cf920f 100644
--- a/src/backend/commands/analyzeutils.c
+++ b/src/backend/commands/analyzeutils.c
@@ -1340,11 +1340,13 @@ bool
aggregate_leaf_partition_ndvbs(int nParts,
HeapTuple
*heaptupleStats,
float4 *relTuples,
- float8 *result)
+ float8 *result,
+ float8
*ndistinct_sum)
{
bool valid;
Assert(nParts > 0);
Assert(result);
+ Assert(ndistinct_sum);
AttStatsSlot **ndvbsSlots = (AttStatsSlot **) palloc0((nParts) *
sizeof(AttStatsSlot *));
valid = getNdvBySegHeapTuple(ndvbsSlots, heaptupleStats, relTuples,
nParts);
@@ -1352,7 +1354,19 @@ aggregate_leaf_partition_ndvbs(int nParts,
for (int i = 0; i < nParts; i++)
{
if (ndvbsSlots[i]) {
+ Form_pg_statistic stat;
+
*result +=
DatumGetFloat8(ndvbsSlots[i]->values[0]);
+
+ /*
+ * Sum the leaves' own ndistinct as well. Both
sums are over
+ * the same leaves and on the same (absolute)
scale, so their
+ * ratio tells how many segments a distinct
value of this
+ * column sits on, on average. See
merge_leaf_stats().
+ */
+ stat = (Form_pg_statistic)
GETSTRUCT(heaptupleStats[i]);
+ *ndistinct_sum += stat->stadistinct < 0 ?
+ -stat->stadistinct * relTuples[i] :
stat->stadistinct;
}
}
}
diff --git a/src/include/commands/analyzeutils.h
b/src/include/commands/analyzeutils.h
index 33ca85d458a..8f936b9b6fe 100644
--- a/src/include/commands/analyzeutils.h
+++ b/src/include/commands/analyzeutils.h
@@ -64,5 +64,6 @@ extern bool leaf_parts_analyzed(Oid attrelid, Oid
relid_exclude, List *va_cols,
extern bool aggregate_leaf_partition_ndvbs(int nParts,
HeapTuple
*heaptupleStats,
float4
*relTuples,
- float8 *result);
+ float8 *result,
+ float8
*ndistinct_sum);
#endif /* ANALYZEUTILS_H */
diff --git a/src/test/regress/expected/incremental_analyze.out
b/src/test/regress/expected/incremental_analyze.out
index a2caf59e0e9..8577c5dacc4 100644
--- a/src/test/regress/expected/incremental_analyze.out
+++ b/src/test/regress/expected/incremental_analyze.out
@@ -2046,3 +2046,81 @@ INFO: analyzing "public.foo_1_prt_20210201"
INFO: Executing SQL: select pg_catalog.gp_acquire_sample_rows(65903, 400,
'f');
INFO: analyzing "public.foo" inheritance tree
rollback;
+-- ndistinct-by-segments of the root is merged from the leaves. For a column
+-- whose values repeat in every partition it must not grow with the number of
+-- partitions, otherwise ORCA overestimates the output of a local aggregate
+-- and gives up the multi-stage plan.
+set default_statistics_target = 100;
+drop table if exists ndvbs_part;
+NOTICE: table "ndvbs_part" does not exist, skipping
+create table ndvbs_part (id int, pk int, a int, c int) distributed by (id)
+ partition by range (pk) (start (1) end (9) every (1));
+insert into ndvbs_part select g, (g % 8) + 1, g % 53, g % 3 from
generate_series(1, 24000) g;
+analyze ndvbs_part;
+select c.relname, a.attname,
+ case 8 when s.stakind1 then s.stavalues1::text when s.stakind2 then
s.stavalues2::text
+ when s.stakind3 then s.stavalues3::text when s.stakind4 then
s.stavalues4::text
+ when s.stakind5 then s.stavalues5::text end as ndv_by_segments
+ from pg_statistic s
+ join pg_class c on c.oid = s.starelid
+ join pg_attribute a on a.attrelid = s.starelid and a.attnum = s.staattnum
+ where c.relname in ('ndvbs_part', 'ndvbs_part_1_prt_1') and a.attname in
('a', 'c')
+ order by 1, 2;
+ relname | attname | ndv_by_segments
+--------------------+---------+-----------------
+ ndvbs_part | a | {159}
+ ndvbs_part | c | {9}
+ ndvbs_part_1_prt_1 | a | {159}
+ ndvbs_part_1_prt_1 | c | {9}
+(4 rows)
+
+set optimizer = on;
+explain (costs off) select a, c, count(*) from ndvbs_part group by a, c;
+ QUERY PLAN
+----------------------------------------------------------------------
+ Gather Motion 3:1 (slice1; segments: 3)
+ -> Finalize HashAggregate
+ Group Key: a, c
+ -> Redistribute Motion 3:3 (slice2; segments: 3)
+ Hash Key: a, c
+ -> Streaming Partial HashAggregate
+ Group Key: a, c
+ -> Dynamic Seq Scan on ndvbs_part
+ Number of partitions to scan: 8 (out of 8)
+ Optimizer: GPORCA
+(10 rows)
+
+reset optimizer;
+reset default_statistics_target;
+drop table ndvbs_part;
+-- A value of a column collocated with the distribution key sits on a single
+-- segment, so its ndistinct-by-segments equals its ndistinct, while a value of
+-- an unrelated column sits on every segment. Neither depends on the number of
+-- partitions repeating the values.
+set default_statistics_target = 100;
+drop table if exists ndvbs_part_seg;
+NOTICE: table "ndvbs_part_seg" does not exist, skipping
+create table ndvbs_part_seg (id int, pk int, dk int, wide int) distributed by
(dk)
+ partition by range (pk) (start (1) end (9) every (1));
+insert into ndvbs_part_seg select g, (g % 8) + 1, g % 53, g % 51 from
generate_series(1, 24000) g;
+analyze ndvbs_part_seg;
+select c.relname, a.attname,
+ case 8 when s.stakind1 then s.stavalues1::text when s.stakind2 then
s.stavalues2::text
+ when s.stakind3 then s.stavalues3::text when s.stakind4 then
s.stavalues4::text
+ when s.stakind5 then s.stavalues5::text end as ndv_by_segments
+ from pg_statistic s
+ join pg_class c on c.oid = s.starelid
+ join pg_attribute a on a.attrelid = s.starelid and a.attnum = s.staattnum
+ where c.relname in ('ndvbs_part_seg', 'ndvbs_part_seg_1_prt_1')
+ and a.attname in ('dk', 'wide')
+ order by 2, 1;
+ relname | attname | ndv_by_segments
+------------------------+---------+-----------------
+ ndvbs_part_seg | dk | {53}
+ ndvbs_part_seg_1_prt_1 | dk | {53}
+ ndvbs_part_seg | wide | {153}
+ ndvbs_part_seg_1_prt_1 | wide | {153}
+(4 rows)
+
+reset default_statistics_target;
+drop table ndvbs_part_seg;
diff --git a/src/test/regress/sql/incremental_analyze.sql
b/src/test/regress/sql/incremental_analyze.sql
index ec418b4d693..c326ca2d893 100644
--- a/src/test/regress/sql/incremental_analyze.sql
+++ b/src/test/regress/sql/incremental_analyze.sql
@@ -883,3 +883,51 @@ truncate foo_1_prt_20210201;
insert into foo select a, '20210101'::date+a from (select
generate_series(31,40) a) t1;
analyze verbose foo_1_prt_20210201;
rollback;
+
+-- ndistinct-by-segments of the root is merged from the leaves. For a column
+-- whose values repeat in every partition it must not grow with the number of
+-- partitions, otherwise ORCA overestimates the output of a local aggregate
+-- and gives up the multi-stage plan.
+set default_statistics_target = 100;
+drop table if exists ndvbs_part;
+create table ndvbs_part (id int, pk int, a int, c int) distributed by (id)
+ partition by range (pk) (start (1) end (9) every (1));
+insert into ndvbs_part select g, (g % 8) + 1, g % 53, g % 3 from
generate_series(1, 24000) g;
+analyze ndvbs_part;
+select c.relname, a.attname,
+ case 8 when s.stakind1 then s.stavalues1::text when s.stakind2 then
s.stavalues2::text
+ when s.stakind3 then s.stavalues3::text when s.stakind4 then
s.stavalues4::text
+ when s.stakind5 then s.stavalues5::text end as ndv_by_segments
+ from pg_statistic s
+ join pg_class c on c.oid = s.starelid
+ join pg_attribute a on a.attrelid = s.starelid and a.attnum = s.staattnum
+ where c.relname in ('ndvbs_part', 'ndvbs_part_1_prt_1') and a.attname in
('a', 'c')
+ order by 1, 2;
+set optimizer = on;
+explain (costs off) select a, c, count(*) from ndvbs_part group by a, c;
+reset optimizer;
+reset default_statistics_target;
+drop table ndvbs_part;
+
+-- A value of a column collocated with the distribution key sits on a single
+-- segment, so its ndistinct-by-segments equals its ndistinct, while a value of
+-- an unrelated column sits on every segment. Neither depends on the number of
+-- partitions repeating the values.
+set default_statistics_target = 100;
+drop table if exists ndvbs_part_seg;
+create table ndvbs_part_seg (id int, pk int, dk int, wide int) distributed by
(dk)
+ partition by range (pk) (start (1) end (9) every (1));
+insert into ndvbs_part_seg select g, (g % 8) + 1, g % 53, g % 51 from
generate_series(1, 24000) g;
+analyze ndvbs_part_seg;
+select c.relname, a.attname,
+ case 8 when s.stakind1 then s.stavalues1::text when s.stakind2 then
s.stavalues2::text
+ when s.stakind3 then s.stavalues3::text when s.stakind4 then
s.stavalues4::text
+ when s.stakind5 then s.stavalues5::text end as ndv_by_segments
+ from pg_statistic s
+ join pg_class c on c.oid = s.starelid
+ join pg_attribute a on a.attrelid = s.starelid and a.attnum = s.staattnum
+ where c.relname in ('ndvbs_part_seg', 'ndvbs_part_seg_1_prt_1')
+ and a.attname in ('dk', 'wide')
+ order by 2, 1;
+reset default_statistics_target;
+drop table ndvbs_part_seg;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]