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]

Reply via email to