From 5ad5a4a15b8a977b0753eaaf48ee8be6a8fc6b9f Mon Sep 17 00:00:00 2001 From: coldWater Date: Tue, 4 Aug 2026 12:31:04 +0800 Subject: [PATCH 1/5] feat(query): support cascading hierarchical grouping sets --- src/query/settings/src/settings_default.rs | 7 + .../settings/src/settings_getter_setter.rs | 4 + .../sql/src/planner/optimizer/optimizer.rs | 3 + .../optimizer/optimizers/cascades/cascade.rs | 49 +--- .../distributed/materialized_cte.rs | 151 ++++++++++ .../optimizer/optimizers/distributed/mod.rs | 2 + .../rule_hierarchical_grouping_sets.rs | 42 +-- .../optimizer/hierarchical_grouping_sets.rs | 125 ++++++++ .../optimizer/hierarchical_grouping_sets.txt | 270 ++++++++++++++++++ src/query/sql/tests/it/optimizer/mod.rs | 1 + .../group_by_grouping_sets_union_all.test | 52 ++++ 11 files changed, 639 insertions(+), 67 deletions(-) create mode 100644 src/query/sql/src/planner/optimizer/optimizers/distributed/materialized_cte.rs create mode 100644 src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.rs create mode 100644 src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.txt diff --git a/src/query/settings/src/settings_default.rs b/src/query/settings/src/settings_default.rs index b5b6e3a3997..696809b34b6 100644 --- a/src/query/settings/src/settings_default.rs +++ b/src/query/settings/src/settings_default.rs @@ -666,6 +666,13 @@ impl DefaultSettings { scope: SettingScope::Both, range: Some(SettingRange::Numeric(0..=1)), }), + ("enable_cascading_grouping_sets", DefaultSettingValue { + value: UserSettingValue::UInt64(0), + desc: "Builds hierarchical grouping sets from the closest available parent grouping.", + mode: SettingMode::Both, + scope: SettingScope::Both, + range: Some(SettingRange::Numeric(0..=1)), + }), ("storage_fetch_part_num", DefaultSettingValue { value: UserSettingValue::UInt64(2), desc: "Sets the number of partitions that are fetched in parallel from storage during query execution.", diff --git a/src/query/settings/src/settings_getter_setter.rs b/src/query/settings/src/settings_getter_setter.rs index 193ea6de6a1..519fb1014b2 100644 --- a/src/query/settings/src/settings_getter_setter.rs +++ b/src/query/settings/src/settings_getter_setter.rs @@ -623,6 +623,10 @@ impl Settings { Ok(self.try_get_u64("grouping_sets_to_union")? == 1) } + pub fn get_enable_cascading_grouping_sets(&self) -> Result { + Ok(self.try_get_u64("enable_cascading_grouping_sets")? == 1) + } + pub fn get_lazy_read_threshold(&self) -> Result { self.try_get_u64("lazy_read_threshold") } diff --git a/src/query/sql/src/planner/optimizer/optimizer.rs b/src/query/sql/src/planner/optimizer/optimizer.rs index 26eb554d9c4..864c78b3d11 100644 --- a/src/query/sql/src/planner/optimizer/optimizer.rs +++ b/src/query/sql/src/planner/optimizer/optimizer.rs @@ -36,6 +36,7 @@ use crate::optimizer::optimizers::CommonSubexpressionOptimizer; use crate::optimizer::optimizers::DPhpyOptimizer; use crate::optimizer::optimizers::EliminateSelfJoinOptimizer; use crate::optimizer::optimizers::distributed::BroadcastToShuffleOptimizer; +use crate::optimizer::optimizers::distributed::MaterializedCTEDistributionOptimizer; use crate::optimizer::optimizers::operator::CleanupUnusedCTEOptimizer; use crate::optimizer::optimizers::operator::DeduplicateJoinConditionOptimizer; use crate::optimizer::optimizers::operator::FinalizeSpatialJoinOptimizer; @@ -293,6 +294,8 @@ pub async fn optimize_query(opt_ctx: Arc, s_expr: SExpr) -> Re ) // Cascades optimizer may fail due to timeout, fallback to heuristic optimizer in this case. .add(CascadesOptimizer::new(opt_ctx.clone())?) + // Normalize distributed MaterializedCTE producers after physical properties are settled. + .add(MaterializedCTEDistributionOptimizer::new(opt_ctx.clone())) // Eliminate unnecessary scalar calculations to clean up the final plan .add_if( !opt_ctx.get_planning_agg_index(), diff --git a/src/query/sql/src/planner/optimizer/optimizers/cascades/cascade.rs b/src/query/sql/src/planner/optimizer/optimizers/cascades/cascade.rs index 20c8eae11dd..dd4b6ee4352 100644 --- a/src/query/sql/src/planner/optimizer/optimizers/cascades/cascade.rs +++ b/src/query/sql/src/planner/optimizer/optimizers/cascades/cascade.rs @@ -25,7 +25,6 @@ use crate::optimizer::OptimizerContext; use crate::optimizer::cost::CostModel; use crate::optimizer::ir::Distribution; use crate::optimizer::ir::Memo; -use crate::optimizer::ir::RelExpr; use crate::optimizer::ir::RequiredProperty; use crate::optimizer::ir::SExpr; use crate::optimizer::optimizers::cascades::cost::DefaultCostModel; @@ -38,7 +37,6 @@ use crate::optimizer::optimizers::distributed::DistributedOptimizer; use crate::optimizer::optimizers::distributed::SortAndLimitPushDownOptimizer; use crate::optimizer::optimizers::rule::RuleSet; use crate::optimizer::optimizers::rule::TransformResult; -use crate::plans::RelOperator; /// A cascades-style search engine to enumerate possible alternations of a relational expression and /// find the optimal one. @@ -93,7 +91,7 @@ impl CascadesOptimizer { let result = self.optimize_internal(s_expr.clone()); // Process different cases based on the result - let mut optimized_expr = match result { + let optimized_expr = match result { Ok(expr) => { // After successful optimization, apply sort and limit push down if distributed optimization is enabled if opt_ctx.get_enable_distributed_optimization() { @@ -122,54 +120,9 @@ impl CascadesOptimizer { } }; - optimized_expr = Self::remove_exchanges_for_serial_sequence(optimized_expr)?; - Ok(optimized_expr) } - fn remove_exchanges_for_serial_sequence(s_expr: SExpr) -> Result { - if Self::has_sequence_with_serial_left_child(&s_expr)? { - Self::remove_all_exchanges(s_expr) - } else { - Ok(s_expr) - } - } - - fn has_sequence_with_serial_left_child(s_expr: &SExpr) -> Result { - if let RelOperator::Sequence(_) = s_expr.plan.as_ref() { - let left_child = s_expr.left_child(); - let rel_expr = RelExpr::with_s_expr(left_child); - let physical_prop = rel_expr.derive_physical_prop()?; - - if physical_prop.distribution == Distribution::Serial { - return Ok(true); - } - } - - for child in s_expr.children() { - if Self::has_sequence_with_serial_left_child(child)? { - return Ok(true); - } - } - - Ok(false) - } - - fn remove_all_exchanges(s_expr: SExpr) -> Result { - if let RelOperator::Exchange(_) = s_expr.plan.as_ref() { - return Self::remove_all_exchanges(s_expr.unary_child().clone()); - } - - let mut new_children = Vec::new(); - for child in s_expr.children() { - let processed_child = Self::remove_all_exchanges(child.clone())?; - new_children.push(Arc::new(processed_child)); - } - - let result = s_expr.replace_children(new_children); - Ok(result) - } - fn optimize_internal(&mut self, s_expr: SExpr) -> Result { // Update rule set based on current flags // This ensures we use the most up-to-date flag values, regardless of when the optimizer was created diff --git a/src/query/sql/src/planner/optimizer/optimizers/distributed/materialized_cte.rs b/src/query/sql/src/planner/optimizer/optimizers/distributed/materialized_cte.rs new file mode 100644 index 00000000000..a8a0090f925 --- /dev/null +++ b/src/query/sql/src/planner/optimizer/optimizers/distributed/materialized_cte.rs @@ -0,0 +1,151 @@ +// Copyright 2021 Datafuse Labs +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use databend_common_exception::ErrorCode; +use databend_common_exception::Result; +use databend_common_expression::Scalar; +use databend_common_expression::types::NumberScalar; + +use crate::optimizer::Optimizer; +use crate::optimizer::OptimizerContext; +use crate::optimizer::ir::Distribution; +use crate::optimizer::ir::RelExpr; +use crate::optimizer::ir::SExpr; +use crate::optimizer::ir::SExprVisitor; +use crate::optimizer::ir::VisitAction; +use crate::plans::ConstantExpr; +use crate::plans::Exchange; +use crate::plans::RelOperator; +use crate::plans::ScalarExpr; + +pub struct MaterializedCTEDistributionOptimizer { + ctx: Arc, +} + +impl MaterializedCTEDistributionOptimizer { + pub fn new(ctx: Arc) -> Self { + Self { ctx } + } + + pub fn optimize_sync(&self, s_expr: &SExpr) -> Result { + let mut result = if self.ctx.get_enable_distributed_optimization() { + s_expr + .accept(&mut SerialProducerRedistributor)? + .unwrap_or_else(|| s_expr.clone()) + } else { + s_expr.clone() + }; + + let mut finder = SerialSequenceFinder::default(); + result.accept(&mut finder)?; + if finder.found { + result = result.accept(&mut ExchangeRemover)?.unwrap_or(result); + } + + Ok(result) + } +} + +/// A Sequence producer must participate in the distributed fragment graph. When +/// its CTE definition ends in a scalar operator, redistribute the result after +/// that operator instead of forcing the entire query to run without Exchanges. +/// Hashing by a constant preserves the producer's rows without broadcasting or +/// changing the scalar operator's empty-input behavior. +struct SerialProducerRedistributor; + +impl SExprVisitor for SerialProducerRedistributor { + fn visit(&mut self, _expr: &SExpr) -> Result { + Ok(VisitAction::Continue) + } + + fn post_visit(&mut self, expr: &SExpr) -> Result { + if !matches!(expr.plan(), RelOperator::Sequence(_)) { + return Ok(VisitAction::Continue); + } + + let left = expr.left_child(); + let physical_prop = RelExpr::with_s_expr(left).derive_physical_prop()?; + if physical_prop.distribution != Distribution::Serial { + return Ok(VisitAction::Continue); + } + if !matches!(left.plan(), RelOperator::MaterializedCTE(_)) { + return Err(ErrorCode::Internal( + "Sequence left child is expected to be MaterializedCTE".to_string(), + )); + } + + let hash_key = ScalarExpr::ConstantExpr(ConstantExpr { + value: Scalar::Number(NumberScalar::UInt32(0)), + span: None, + }); + let exchange = left + .unary_child_arc() + .ref_build_unary(Exchange::GlobalHash(vec![hash_key])); + let left = left.replace_children([Arc::new(exchange)]); + Ok(VisitAction::Replace(expr.replace_left_child(left))) + } +} + +#[derive(Default)] +struct SerialSequenceFinder { + found: bool, +} + +impl SExprVisitor for SerialSequenceFinder { + fn visit(&mut self, expr: &SExpr) -> Result { + if self.found { + return Ok(VisitAction::SkipChildren); + } + + if matches!(expr.plan(), RelOperator::Sequence(_)) { + let left = expr.left_child(); + let physical_prop = RelExpr::with_s_expr(left).derive_physical_prop()?; + if physical_prop.distribution == Distribution::Serial { + self.found = true; + return Ok(VisitAction::SkipChildren); + } + } + + Ok(VisitAction::Continue) + } +} + +struct ExchangeRemover; + +impl SExprVisitor for ExchangeRemover { + fn visit(&mut self, _expr: &SExpr) -> Result { + Ok(VisitAction::Continue) + } + + fn post_visit(&mut self, expr: &SExpr) -> Result { + if matches!(expr.plan(), RelOperator::Exchange(_)) { + Ok(VisitAction::Replace(expr.unary_child().clone())) + } else { + Ok(VisitAction::Continue) + } + } +} + +#[async_trait::async_trait] +impl Optimizer for MaterializedCTEDistributionOptimizer { + fn name(&self) -> String { + "MaterializedCTEDistributionOptimizer".to_string() + } + + async fn optimize(&mut self, s_expr: &SExpr) -> Result { + self.optimize_sync(s_expr) + } +} diff --git a/src/query/sql/src/planner/optimizer/optimizers/distributed/mod.rs b/src/query/sql/src/planner/optimizer/optimizers/distributed/mod.rs index 81c6da2b752..b8728e840ba 100644 --- a/src/query/sql/src/planner/optimizer/optimizers/distributed/mod.rs +++ b/src/query/sql/src/planner/optimizer/optimizers/distributed/mod.rs @@ -15,8 +15,10 @@ #[allow(clippy::module_inception)] mod distributed; mod distributed_merge; +mod materialized_cte; mod sort_and_limit; pub use distributed::DistributedOptimizer; pub use distributed_merge::BroadcastToShuffleOptimizer; +pub use materialized_cte::MaterializedCTEDistributionOptimizer; pub use sort_and_limit::SortAndLimitPushDownOptimizer; diff --git a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs index 6087062f56f..89537506806 100644 --- a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs +++ b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs @@ -51,8 +51,6 @@ use crate::plans::UnionAll; use crate::plans::VisitorMut; use crate::plans::walk_expr_mut; -const ID: RuleID = RuleID::HierarchicalGroupingSetsToUnion; - /// True hierarchical optimization for GROUPING SETS with multi-layer dependency analysis. /// /// This implements genuine hierarchical aggregation where higher-level groupings @@ -65,20 +63,17 @@ const ID: RuleID = RuleID::HierarchicalGroupingSetsToUnion; /// Level 3: CTE_level_1_* -> GROUP BY () -> CTE_level_0 /// Final: UNION ALL of all results pub struct RuleHierarchicalGroupingSetsToUnion { - id: RuleID, matchers: Vec, cte_channel_size: usize, + enable_cascading: bool, } impl RuleHierarchicalGroupingSetsToUnion { pub fn new(ctx: Arc) -> Self { - let cte_channel_size = ctx - .get_table_ctx() - .get_settings() - .get_grouping_sets_channel_size() - .unwrap(); + let settings = ctx.get_table_ctx().get_settings(); + let cte_channel_size = settings.get_grouping_sets_channel_size().unwrap(); + let enable_cascading = settings.get_enable_cascading_grouping_sets().unwrap(); Self { - id: ID, matchers: vec![Matcher::MatchOp { op_type: RelOp::EvalScalar, children: vec![Matcher::MatchOp { @@ -87,6 +82,7 @@ impl RuleHierarchicalGroupingSetsToUnion { }], }], cte_channel_size: cte_channel_size as usize, + enable_cascading, } } @@ -206,18 +202,26 @@ impl RuleHierarchicalGroupingSetsToUnion { } } - /// Optimize the hierarchy to minimize intermediate CTEs while maximizing reuse + /// Choose between the existing shared-base hierarchy and a nearest-parent cascade. fn optimize_hierarchy(&self, levels: &mut [GroupingLevel]) { - // For each level, choose the most detailed parent that can generate it for i in 0..levels.len() { if !levels[i].possible_parents.is_empty() { - // Choose the parent with maximum columns (most detailed) - // This ensures we reuse the most specific aggregation available - let best_parent_level_idx = *levels[i] - .possible_parents - .iter() - .max_by_key(|&&parent_idx| levels[parent_idx].level) - .unwrap(); + let best_parent_level_idx = if self.enable_cascading { + // The closest strict superset minimizes the rows re-aggregated by a + // ROLLUP chain, at the cost of a longer dependency path. + *levels[i] + .possible_parents + .iter() + .min_by_key(|&&parent_idx| levels[parent_idx].level) + .unwrap() + } else { + // Preserve the existing shared-base shape for controlled A/B tests. + *levels[i] + .possible_parents + .iter() + .max_by_key(|&&parent_idx| levels[parent_idx].level) + .unwrap() + }; // Store the set_index of the chosen parent, not the level index levels[i].chosen_parent = Some(levels[best_parent_level_idx].set_index); @@ -878,7 +882,7 @@ impl RuleHierarchicalGroupingSetsToUnion { impl Rule for RuleHierarchicalGroupingSetsToUnion { fn id(&self) -> RuleID { - self.id + RuleID::HierarchicalGroupingSetsToUnion } fn apply(&self, s_expr: &SExpr, state: &mut TransformResult) -> Result<()> { diff --git a/src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.rs b/src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.rs new file mode 100644 index 00000000000..7b5fa99ac65 --- /dev/null +++ b/src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.rs @@ -0,0 +1,125 @@ +// Copyright 2021 Datafuse Labs +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use databend_common_catalog::table_context::TableContextSettings; +use databend_common_exception::Result; +use databend_common_sql::optimizer::ir::Distribution; +use databend_common_sql::optimizer::ir::RelExpr; +use databend_common_sql::optimizer::ir::SExpr; +use databend_common_sql::optimizer::ir::SExprVisitor; +use databend_common_sql::optimizer::ir::VisitAction; +use databend_common_sql::plans::Operator; +use databend_common_sql::plans::Plan; +use databend_common_sql::plans::RelOp; + +use crate::framework::LiteTableContext; +use crate::framework::golden::SqlTestCase; +use crate::framework::golden::open_golden_file; +use crate::framework::golden::write_case_header; + +async fn write_optimized_case( + file: &mut impl std::io::Write, + case: &SqlTestCase, + enable_cascading: bool, +) -> Result<()> { + let ctx = LiteTableContext::create().await?; + ctx.set_cluster_node_num(2); + ctx.set_table_warehouse_distribution(true); + for setup_sql in case.setup_sqls { + ctx.register_setup_sql(setup_sql).await?; + } + ctx.get_settings() + .set_setting("grouping_sets_to_union".to_string(), "1".to_string())?; + ctx.get_settings().set_setting( + "enable_cascading_grouping_sets".to_string(), + u8::from(enable_cascading).to_string(), + )?; + + let raw_plan = ctx.bind_sql(case.sql).await?; + let optimized_plan = ctx.optimize_plan(raw_plan.clone()).await?; + assert_no_serial_sequence_producer(&optimized_plan)?; + + write_case_header(file, case)?; + writeln!(file, "raw_plan:")?; + writeln!(file, "{}", raw_plan.format_indent(Default::default())?)?; + writeln!(file, "optimized_plan:")?; + writeln!( + file, + "{}", + optimized_plan.format_indent(Default::default())? + )?; + writeln!(file)?; + + Ok(()) +} + +fn assert_no_serial_sequence_producer(plan: &Plan) -> Result<()> { + let Plan::Query { s_expr, .. } = plan else { + unreachable!("the ROLLUP test query should bind to Plan::Query") + }; + + struct SequenceDistributionChecker; + + impl SExprVisitor for SequenceDistributionChecker { + fn visit(&mut self, expr: &SExpr) -> Result { + if expr.plan().rel_op() == RelOp::Sequence { + let left_prop = RelExpr::with_s_expr(expr.left_child()).derive_physical_prop()?; + assert_ne!( + left_prop.distribution, + Distribution::Serial, + "a Serial Sequence producer removes every Exchange from the query" + ); + } + Ok(VisitAction::Continue) + } + } + + s_expr.accept(&mut SequenceDistributionChecker)?; + Ok(()) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +async fn test_hierarchical_grouping_sets_optimizer_outcomes() -> Result<()> { + let mut file = open_golden_file("optimizer", "hierarchical_grouping_sets.txt")?; + + let shared_base = SqlTestCase { + name: "rollup_shared_base", + description: "The default hierarchy derives every lower ROLLUP level from the most detailed grouping.", + setup_sqls: &[ROLLUP_TABLE], + sql: ROLLUP_QUERY, + }; + write_optimized_case(&mut file, &shared_base, false).await?; + + let cascading = SqlTestCase { + name: "rollup_nearest_parent_cascade", + description: "The optional cascade derives every lower ROLLUP level from its closest strict superset.", + setup_sqls: &[ROLLUP_TABLE], + sql: ROLLUP_QUERY, + }; + write_optimized_case(&mut file, &cascading, true).await?; + + Ok(()) +} + +const ROLLUP_TABLE: &str = "CREATE TABLE t +( + a UInt64, + b UInt64, + c UInt64, + v UInt64 +)"; + +const ROLLUP_QUERY: &str = "SELECT a, b, c, sum(v), count(v), min(v), max(v) +FROM t +GROUP BY ROLLUP(a, b, c)"; diff --git a/src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.txt b/src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.txt new file mode 100644 index 00000000000..77c5e05de27 --- /dev/null +++ b/src/query/sql/tests/it/optimizer/hierarchical_grouping_sets.txt @@ -0,0 +1,270 @@ +=== rollup_shared_base === +description: The default hierarchy derives every lower ROLLUP level from the most detailed grouping. +sql: SELECT a, b, c, sum(v), count(v), min(v), max(v) +FROM t +GROUP BY ROLLUP(a, b, c) +raw_plan: +EvalScalar +├── scalars: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] +└── Aggregate(Initial) + ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2), _grouping_id (#7) AS (#7)] + ├── aggregate functions: [sum(t.v (#3)) AS (#8), count(t.v (#3)) AS (#9), min(t.v (#3)) AS (#10), max(t.v (#3)) AS (#11)] + └── EvalScalar + ├── scalars: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2), t.v (#3) AS (#3)] + └── Scan + ├── table: default.t (#0) + ├── filters: [] + ├── order by: [] + └── limit: NONE + +optimized_plan: +Exchange(Merge) +└── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ ├── ref_count: 4 + │ ├── channel_size: Some(2) + │ └── Aggregate(Final) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2)] + │ ├── aggregate functions: [sum(t.v (#3)) AS (#8), count(t.v (#3)) AS (#9), min(t.v (#3)) AS (#10), max(t.v (#3)) AS (#11)] + │ └── Aggregate(Partial) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2)] + │ ├── aggregate functions: [sum(t.v (#3)) AS (#8), count(t.v (#3)) AS (#9), min(t.v (#3)) AS (#10), max(t.v (#3)) AS (#11)] + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [t.a (#0)] + │ └── Scan + │ ├── table: default.t (#0) + │ ├── filters: [] + │ ├── order by: [] + │ └── limit: NONE + └── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1 + │ ├── ref_count: 1 + │ ├── channel_size: Some(2) + │ └── Aggregate(Final) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Aggregate(Partial) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [t.a (#0)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1), t.c (#2)] + │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), t.c (#2) -> t.c (#2), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0 + │ ├── ref_count: 1 + │ ├── channel_size: Some(2) + │ └── Aggregate(Final) + │ ├── group items: [t.a (#0) AS (#0)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Aggregate(Partial) + │ ├── group items: [t.a (#0) AS (#0)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [t.a (#0)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1), t.c (#2)] + │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), t.c (#2) -> t.c (#2), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_empty + │ ├── ref_count: 1 + │ ├── channel_size: Some(2) + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [0] + │ └── Aggregate(Final) + │ ├── group items: [] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Exchange(Merge) + │ └── Aggregate(Partial) + │ ├── group items: [] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1), t.c (#2)] + │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), t.c (#2) -> t.c (#2), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── UnionAll + ├── output: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + ├── left: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + ├── right: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + ├── cte_scan_names: [] + ├── logical_recursive_cte_id: None + ├── UnionAll + │ ├── output: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ ├── left: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ ├── right: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ ├── cte_scan_names: [] + │ ├── logical_recursive_cte_id: None + │ ├── UnionAll + │ │ ├── output: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ │ ├── left: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ │ ├── right: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ │ ├── cte_scan_names: [] + │ │ ├── logical_recursive_cte_id: None + │ │ ├── EvalScalar + │ │ │ ├── scalars: [CAST(t.a (#0) AS UInt64 NULL) AS (#0), CAST(t.b (#1) AS UInt64 NULL) AS (#1), CAST(t.c (#2) AS UInt64 NULL) AS (#2), 0 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + │ │ │ └── MaterializedCTERef + │ │ │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ │ │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1), t.c (#2)] + │ │ │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), t.c (#2) -> t.c (#2), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + │ │ └── EvalScalar + │ │ ├── scalars: [CAST(t.a (#0) AS UInt64 NULL) AS (#0), CAST(t.b (#1) AS UInt64 NULL) AS (#1), NULL AS (#2), 4 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + │ │ └── MaterializedCTERef + │ │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1 + │ │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1)] + │ │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + │ └── EvalScalar + │ ├── scalars: [CAST(t.a (#0) AS UInt64 NULL) AS (#0), NULL AS (#1), NULL AS (#2), 6 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0)] + │ └── column mapping: [t.a (#0) -> t.a (#0), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── EvalScalar + ├── scalars: [NULL AS (#0), NULL AS (#1), NULL AS (#2), 7 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + └── MaterializedCTERef + ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_empty + ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11)] + └── column mapping: [sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + + +=== rollup_nearest_parent_cascade === +description: The optional cascade derives every lower ROLLUP level from its closest strict superset. +sql: SELECT a, b, c, sum(v), count(v), min(v), max(v) +FROM t +GROUP BY ROLLUP(a, b, c) +raw_plan: +EvalScalar +├── scalars: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] +└── Aggregate(Initial) + ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2), _grouping_id (#7) AS (#7)] + ├── aggregate functions: [sum(t.v (#3)) AS (#8), count(t.v (#3)) AS (#9), min(t.v (#3)) AS (#10), max(t.v (#3)) AS (#11)] + └── EvalScalar + ├── scalars: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2), t.v (#3) AS (#3)] + └── Scan + ├── table: default.t (#0) + ├── filters: [] + ├── order by: [] + └── limit: NONE + +optimized_plan: +Exchange(Merge) +└── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ ├── ref_count: 2 + │ ├── channel_size: Some(2) + │ └── Aggregate(Final) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2)] + │ ├── aggregate functions: [sum(t.v (#3)) AS (#8), count(t.v (#3)) AS (#9), min(t.v (#3)) AS (#10), max(t.v (#3)) AS (#11)] + │ └── Aggregate(Partial) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1), t.c (#2) AS (#2)] + │ ├── aggregate functions: [sum(t.v (#3)) AS (#8), count(t.v (#3)) AS (#9), min(t.v (#3)) AS (#10), max(t.v (#3)) AS (#11)] + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [t.a (#0)] + │ └── Scan + │ ├── table: default.t (#0) + │ ├── filters: [] + │ ├── order by: [] + │ └── limit: NONE + └── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1 + │ ├── ref_count: 2 + │ ├── channel_size: Some(2) + │ └── Aggregate(Final) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Aggregate(Partial) + │ ├── group items: [t.a (#0) AS (#0), t.b (#1) AS (#1)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [t.a (#0)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1), t.c (#2)] + │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), t.c (#2) -> t.c (#2), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0 + │ ├── ref_count: 2 + │ ├── channel_size: Some(2) + │ └── Aggregate(Final) + │ ├── group items: [t.a (#0) AS (#0)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Aggregate(Partial) + │ ├── group items: [t.a (#0) AS (#0)] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [t.a (#0)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1)] + │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── Sequence(Sequence) + ├── MaterializedCTE + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_empty + │ ├── ref_count: 1 + │ ├── channel_size: Some(2) + │ └── Exchange(Hash) + │ ├── Exchange(Hash): keys: [0] + │ └── Aggregate(Final) + │ ├── group items: [] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── Exchange(Merge) + │ └── Aggregate(Partial) + │ ├── group items: [] + │ ├── aggregate functions: [sum(sum(v) (#8)) AS (#8), sum0(count(v) (#9)) AS (#9), min(min(v) (#10)) AS (#10), max(max(v) (#11)) AS (#11)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0)] + │ └── column mapping: [t.a (#0) -> t.a (#0), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── UnionAll + ├── output: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + ├── left: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + ├── right: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + ├── cte_scan_names: [] + ├── logical_recursive_cte_id: None + ├── UnionAll + │ ├── output: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ ├── left: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ ├── right: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ ├── cte_scan_names: [] + │ ├── logical_recursive_cte_id: None + │ ├── UnionAll + │ │ ├── output: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ │ ├── left: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ │ ├── right: [t.a (#0), t.b (#1), t.c (#2), sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), _grouping_id (#7)] + │ │ ├── cte_scan_names: [] + │ │ ├── logical_recursive_cte_id: None + │ │ ├── EvalScalar + │ │ │ ├── scalars: [CAST(t.a (#0) AS UInt64 NULL) AS (#0), CAST(t.b (#1) AS UInt64 NULL) AS (#1), CAST(t.c (#2) AS UInt64 NULL) AS (#2), 0 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + │ │ │ └── MaterializedCTERef + │ │ │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1_2 + │ │ │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1), t.c (#2)] + │ │ │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), t.c (#2) -> t.c (#2), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + │ │ └── EvalScalar + │ │ ├── scalars: [CAST(t.a (#0) AS UInt64 NULL) AS (#0), CAST(t.b (#1) AS UInt64 NULL) AS (#1), NULL AS (#2), 4 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + │ │ └── MaterializedCTERef + │ │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0_1 + │ │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0), t.b (#1)] + │ │ └── column mapping: [t.a (#0) -> t.a (#0), t.b (#1) -> t.b (#1), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + │ └── EvalScalar + │ ├── scalars: [CAST(t.a (#0) AS UInt64 NULL) AS (#0), NULL AS (#1), NULL AS (#2), 6 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + │ └── MaterializedCTERef + │ ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_cols_0 + │ ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11), t.a (#0)] + │ └── column mapping: [t.a (#0) -> t.a (#0), sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + └── EvalScalar + ├── scalars: [NULL AS (#0), NULL AS (#1), NULL AS (#2), 7 AS (#7), sum(v) (#8) AS (#8), count(v) (#9) AS (#9), min(v) (#10) AS (#10), max(v) (#11) AS (#11)] + └── MaterializedCTERef + ├── cte_name: cte_hierarchical_groupingsets_17760633370082219301_empty + ├── output columns: [sum(v) (#8), count(v) (#9), min(v) (#10), max(v) (#11)] + └── column mapping: [sum(v) (#8) -> sum(v) (#8), count(v) (#9) -> count(v) (#9), min(v) (#10) -> min(v) (#10), max(v) (#11) -> max(v) (#11)] + + diff --git a/src/query/sql/tests/it/optimizer/mod.rs b/src/query/sql/tests/it/optimizer/mod.rs index 2c4c809277b..8f66a927314 100644 --- a/src/query/sql/tests/it/optimizer/mod.rs +++ b/src/query/sql/tests/it/optimizer/mod.rs @@ -15,6 +15,7 @@ mod collect_statistics; mod decorrelate_correlated_aliases; mod eager_aggregation; +mod hierarchical_grouping_sets; mod join_cardinality; mod normalize_scalar; mod outer_join_to_anti; diff --git a/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test b/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test index c815d16d844..929e56ae777 100644 --- a/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test +++ b/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test @@ -1,7 +1,59 @@ statement ok set grouping_sets_to_union = 1; +statement ok +set enable_cascading_grouping_sets = 0; + +include ./group_by_grouping_sets.test + +statement ok +set enable_cascading_grouping_sets = 1; + include ./group_by_grouping_sets.test +statement ok +use default; + +statement ok +create or replace table cascading_grouping_sets (a UInt64, b UInt64, v Nullable(UInt64)); + +statement ok +insert into cascading_grouping_sets values (1, 1, 10), (1, 2, NULL), (2, 1, 5); + +query IIIIIII rowsort +select a, b, sum(v), sum0(v), count(v), min(v), max(v) +from cascading_grouping_sets +group by rollup(a, b) +order by a nulls last, b nulls last; +---- +1 1 10 10 1 10 10 +1 2 NULL 0 0 NULL NULL +1 NULL 10 10 1 10 10 +2 1 5 5 1 5 5 +2 NULL 5 5 1 5 5 +NULL NULL 15 15 2 5 10 + +statement ok +truncate table cascading_grouping_sets; + +# GROUP BY () must still emit one row for empty input in cascade mode. +query IIIIIII +select a, b, sum(v), sum0(v), count(v), min(v), max(v) +from cascading_grouping_sets +group by rollup(a, b); +---- +NULL NULL NULL 0 0 NULL NULL + +statement ok +set enable_cascading_grouping_sets = 0; + +# The shared-base mode uses the same non-Serial empty-grouping producer. +query IIIIIII +select a, b, sum(v), sum0(v), count(v), min(v), max(v) +from cascading_grouping_sets +group by rollup(a, b); +---- +NULL NULL NULL 0 0 NULL NULL + statement ok unset grouping_sets_to_union; From 3b8a26087c26c52f3c888b1ff9fadd62ace7d80d Mon Sep 17 00:00:00 2001 From: coldWater Date: Tue, 4 Aug 2026 21:03:33 +0800 Subject: [PATCH 2/5] fix --- .../service/src/schedulers/fragments/fragmenter.rs | 13 +++++++++++++ .../src/schedulers/fragments/plan_fragment.rs | 7 ++----- 2 files changed, 15 insertions(+), 5 deletions(-) diff --git a/src/query/service/src/schedulers/fragments/fragmenter.rs b/src/query/service/src/schedulers/fragments/fragmenter.rs index 44a463eac6d..d0076e05c31 100644 --- a/src/query/service/src/schedulers/fragments/fragmenter.rs +++ b/src/query/service/src/schedulers/fragments/fragmenter.rs @@ -111,12 +111,23 @@ impl Fragmenter { fragment_id: self.ctx.fragment_id().next_fragment_id(), exchange: None, query_id: self.query_id.clone(), + has_merge_input: false, source_fragments: self.fragments, }); let edges = Self::collect_fragments_edge(fragments.values()); for (source, target) in edges { + let has_merge_input = fragments + .get(&source) + .is_some_and(|fragment| matches!(fragment.exchange, Some(DataExchange::Merge(_)))); + + if has_merge_input { + if let Some(fragment) = fragments.get_mut(&target) { + fragment.has_merge_input = true; + } + } + let Some(fragment) = fragments.get_mut(&source) else { continue; }; @@ -320,6 +331,7 @@ impl DeriveHandle for FragmentDeriveHandle { source_fragments: vec![], fragment_id: source_fragment_id, query_id: self.query_id.clone(), + has_merge_input: false, }; self.fragments.insert(source_fragment_id, source_fragment); @@ -354,6 +366,7 @@ impl DeriveHandle for FragmentDeriveHandle { fragment_id, exchange: None, query_id: self.query_id.clone(), + has_merge_input: false, source_fragments: vec![], }; diff --git a/src/query/service/src/schedulers/fragments/plan_fragment.rs b/src/query/service/src/schedulers/fragments/plan_fragment.rs index b08ced7c014..91e8db2ca5a 100644 --- a/src/query/service/src/schedulers/fragments/plan_fragment.rs +++ b/src/query/service/src/schedulers/fragments/plan_fragment.rs @@ -77,6 +77,7 @@ pub struct PlanFragment { pub fragment_id: usize, pub exchange: Option, pub query_id: String, + pub has_merge_input: bool, // The fragments to ask data from. pub source_fragments: Vec, @@ -103,11 +104,7 @@ impl PlanFragment { fragment_actions.add_action(action); } FragmentType::Intermediate => { - if self - .source_fragments - .iter() - .any(|fragment| matches!(&fragment.exchange, Some(DataExchange::Merge(_)))) - { + if self.has_merge_input { // If this is a intermediate fragment with merge input, // we will only send it to coordinator node. let action = QueryFragmentAction::create( From 7bb3a533c662ab8326872c9b11f8cd162e999d1c Mon Sep 17 00:00:00 2001 From: coldWater Date: Wed, 5 Aug 2026 12:21:05 +0800 Subject: [PATCH 3/5] fix --- .../src/schedulers/fragments/plan_fragment.rs | 46 ++++++++++++++++--- 1 file changed, 40 insertions(+), 6 deletions(-) diff --git a/src/query/service/src/schedulers/fragments/plan_fragment.rs b/src/query/service/src/schedulers/fragments/plan_fragment.rs index 91e8db2ca5a..6e90a2e0398 100644 --- a/src/query/service/src/schedulers/fragments/plan_fragment.rs +++ b/src/query/service/src/schedulers/fragments/plan_fragment.rs @@ -23,6 +23,7 @@ use databend_common_exception::ErrorCode; use databend_common_exception::Result; use databend_common_expression::BlockEntry; use databend_common_expression::Column; +use databend_common_expression::ColumnBuilder; use databend_common_expression::DataBlock; use databend_common_settings::ReplaceIntoShuffleStrategy; use databend_storages_common_table_meta::meta::BlockSlotDescription; @@ -36,6 +37,7 @@ use crate::physical_plans::IPhysicalPlan; use crate::physical_plans::MutationSource; use crate::physical_plans::PhysicalPlan; use crate::physical_plans::PhysicalPlanCast; +use crate::physical_plans::PhysicalPlanMeta; use crate::physical_plans::PhysicalPlanVisitor; use crate::physical_plans::Recluster; use crate::physical_plans::ReplaceDeduplicate; @@ -105,13 +107,45 @@ impl PlanFragment { } FragmentType::Intermediate => { if self.has_merge_input { - // If this is a intermediate fragment with merge input, - // we will only send it to coordinator node. - let action = QueryFragmentAction::create( - Fragmenter::get_local_executor(ctx), - self.plan.clone(), - ); + // Only the coordinator can consume the merge input. Other shuffle + // destinations still need this fragment to receive remote data. + let local_executor = Fragmenter::get_local_executor(ctx); + let action = + QueryFragmentAction::create(local_executor.clone(), self.plan.clone()); fragment_actions.add_action(action); + + if let Some(exchange) = &self.exchange { + let mut empty_plan = self.plan.clone(); + let Some(exchange_sink) = + ExchangeSink::from_mut_physical_plan(&mut empty_plan) + else { + return Err(ErrorCode::Internal( + "Intermediate fragment exchange plan has no ExchangeSink", + )); + }; + exchange_sink.input = PhysicalPlan::new(ConstantTableScan { + meta: PhysicalPlanMeta::new("ConstantTableScan"), + values: exchange_sink + .schema + .fields() + .iter() + .map(|field| { + ColumnBuilder::with_capacity(field.data_type(), 0).build() + }) + .collect(), + num_rows: 0, + output_schema: exchange_sink.schema.clone(), + }); + + for executor in exchange.get_destinations() { + if executor != local_executor { + fragment_actions.add_action(QueryFragmentAction::create( + executor, + empty_plan.clone(), + )); + } + } + } } else { // Otherwise distribute the fragment to all the executors. for executor in Fragmenter::get_executors(ctx) { From af3952449c80902e579d07d053fe182fa3c3415e Mon Sep 17 00:00:00 2001 From: coldWater Date: Wed, 5 Aug 2026 19:48:58 +0800 Subject: [PATCH 4/5] z --- .../rule/agg_rules/grouping_sets_common.rs | 71 +++++++++++++++++++ .../optimizers/rule/agg_rules/mod.rs | 1 + .../agg_rules/rule_grouping_sets_to_union.rs | 22 +++++- .../rule_hierarchical_grouping_sets.rs | 12 ++-- .../group/group_by_grouping_sets.test | 23 ++++++ .../group_by_grouping_sets_union_all.test | 28 ++++++++ 6 files changed, 149 insertions(+), 8 deletions(-) create mode 100644 src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/grouping_sets_common.rs diff --git a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/grouping_sets_common.rs b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/grouping_sets_common.rs new file mode 100644 index 00000000000..db919078e62 --- /dev/null +++ b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/grouping_sets_common.rs @@ -0,0 +1,71 @@ +// Copyright 2021 Datafuse Labs +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use databend_common_exception::Result; + +use crate::Symbol; +use crate::planner::binder::is_grouping_id_item; +use crate::plans::Aggregate; +use crate::plans::BoundColumnRef; +use crate::plans::EvalScalar; +use crate::plans::ScalarItem; + +pub(super) fn ensure_group_items_are_projected( + eval_scalar: &mut EvalScalar, + agg: &Aggregate, + grouping_id_index: Symbol, +) -> Result<()> { + for group_item in agg + .group_items + .iter() + .filter(|item| !is_grouping_id_item(item, grouping_id_index)) + { + if eval_scalar + .items + .iter() + .any(|item| item.index == group_item.index) + { + continue; + } + + eval_scalar.items.push(ScalarItem { + scalar: BoundColumnRef { + span: None, + column: group_item.column_binding(format!("group_item_{}", group_item.index))?, + } + .into(), + index: group_item.index, + }); + } + Ok(()) +} + +pub(super) fn union_output_indexes( + eval_scalar: &EvalScalar, + agg: &Aggregate, + grouping_id_index: Symbol, +) -> Vec { + let mut output_indexes: Vec<_> = eval_scalar.items.iter().map(|item| item.index).collect(); + + for item in agg.aggregate_functions.iter().chain(agg.group_items.iter()) { + if !output_indexes.contains(&item.index) { + output_indexes.push(item.index); + } + } + + if !output_indexes.contains(&grouping_id_index) { + output_indexes.push(grouping_id_index); + } + output_indexes +} diff --git a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/mod.rs b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/mod.rs index 532277a6069..5341a57790b 100644 --- a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/mod.rs +++ b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/mod.rs @@ -13,6 +13,7 @@ // limitations under the License. mod agg_index; +mod grouping_sets_common; mod rule_eager_aggregation; mod rule_fold_count_aggregate; mod rule_grouping_sets_to_union; diff --git a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_grouping_sets_to_union.rs b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_grouping_sets_to_union.rs index 097cdcf7c43..5a2a0377b73 100644 --- a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_grouping_sets_to_union.rs +++ b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_grouping_sets_to_union.rs @@ -21,6 +21,8 @@ use databend_common_exception::Result; use databend_common_expression::Scalar; use databend_common_expression::types::NumberScalar; +use super::grouping_sets_common::ensure_group_items_are_projected; +use super::grouping_sets_common::union_output_indexes; use crate::ScalarExpr; use crate::Symbol; use crate::optimizer::OptimizerContext; @@ -128,6 +130,13 @@ impl Rule for RuleGroupingSetsToUnion { }), }); } + ensure_group_items_are_projected( + &mut eval_scalar, + &agg, + grouping_sets.grouping_id_index, + )?; + let output_indexes = + union_output_indexes(&eval_scalar, &agg, grouping_sets.grouping_id_index); let mut children = Vec::with_capacity(grouping_sets.sets.len()); @@ -209,7 +218,14 @@ impl Rule for RuleGroupingSetsToUnion { }; for scalar in eval_scalar.items.iter_mut() { - visitor.visit(&mut scalar.scalar)?; + if scalar.index == grouping_sets.grouping_id_index { + scalar.scalar = ScalarExpr::ConstantExpr(ConstantExpr { + value: Scalar::Number(NumberScalar::UInt32(grouping_id)), + span: None, + }); + } else { + visitor.visit(&mut scalar.scalar)?; + } } let agg_plan = SExpr::create_unary(agg, cte_consumer.clone()); @@ -221,7 +237,7 @@ impl Rule for RuleGroupingSetsToUnion { let mut result = children.first().unwrap().clone(); for other in children.into_iter().skip(1) { let left_outputs: Vec<(Symbol, Option)> = - eval_scalar.items.iter().map(|x| (x.index, None)).collect(); + output_indexes.iter().map(|index| (*index, None)).collect(); let right_outputs = left_outputs.clone(); let union_plan = UnionAll { @@ -229,7 +245,7 @@ impl Rule for RuleGroupingSetsToUnion { right_outputs, cte_scan_names: vec![], logical_recursive_cte_id: None, - output_indexes: eval_scalar.items.iter().map(|x| x.index).collect(), + output_indexes: output_indexes.clone(), }; result = SExpr::create_binary(Arc::new(union_plan.into()), result, other); } diff --git a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs index 89537506806..326fc1b6ead 100644 --- a/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs +++ b/src/query/sql/src/planner/optimizer/optimizers/rule/agg_rules/rule_hierarchical_grouping_sets.rs @@ -24,6 +24,8 @@ use databend_common_expression::types::DataType; use databend_common_expression::types::NumberDataType; use databend_common_expression::types::NumberScalar; +use super::grouping_sets_common::ensure_group_items_are_projected; +use super::grouping_sets_common::union_output_indexes; use crate::ColumnBindingBuilder; use crate::Symbol; use crate::Visibility; @@ -408,7 +410,7 @@ impl RuleHierarchicalGroupingSetsToUnion { // Step 4: Assemble the complete plan let union_result = - self.create_union_all(&union_branches, eval_scalar, grouping_id_index)?; + self.create_union_all(&union_branches, eval_scalar, agg, grouping_id_index)?; // Step 5: Chain all CTEs in correct dependency order // Sequence semantics: left executes first, right executes after @@ -781,6 +783,8 @@ impl RuleHierarchicalGroupingSetsToUnion { agg: &Aggregate, grouping_id_index: Symbol, ) -> Result<()> { + ensure_group_items_are_projected(eval_scalar, agg, grouping_id_index)?; + let grouping_id = self.calculate_grouping_id(group_columns, &agg.group_items, grouping_id_index); @@ -826,6 +830,7 @@ impl RuleHierarchicalGroupingSetsToUnion { &self, branches: &[SExpr], eval_scalar: &EvalScalar, + agg: &Aggregate, grouping_id_index: Symbol, ) -> Result { if branches.is_empty() { @@ -834,10 +839,7 @@ impl RuleHierarchicalGroupingSetsToUnion { )); } - let mut output_indexes: Vec = eval_scalar.items.iter().map(|x| x.index).collect(); - if !output_indexes.contains(&grouping_id_index) { - output_indexes.push(grouping_id_index); - } + let output_indexes = union_output_indexes(eval_scalar, agg, grouping_id_index); let mut result = branches[0].clone(); for branch in branches.iter().skip(1) { diff --git a/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets.test b/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets.test index b3c87647786..e41d1c927ff 100644 --- a/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets.test +++ b/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets.test @@ -108,6 +108,29 @@ QUALIFY g = 0 a 0 3 b 0 3 +# The grouping-sets rewrite must preserve aggregate and group columns consumed by Window. +query ITTII +SELECT sum(c) AS total, + a, + b, + grouping(a) + grouping(b) AS lochierarchy, + rank() OVER ( + PARTITION BY grouping(a) + grouping(b), + CASE WHEN grouping(b) = 0 THEN a END + ORDER BY sum(c) DESC + ) AS rank_within_parent +FROM t +GROUP BY ROLLUP(a, b) +ORDER BY lochierarchy DESC, a NULLS LAST, b NULLS LAST; +---- +18 NULL NULL 2 1 +7 a NULL 1 2 +11 b NULL 1 1 +3 a A 0 2 +4 a B 0 1 +5 b A 0 2 +6 b B 0 1 + query TTI select a, b, sum(c) as sc from t group by grouping sets ((a,b),(),(b),(a)) order by sc; ---- diff --git a/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test b/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test index 929e56ae777..83f65122cb9 100644 --- a/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test +++ b/tests/sqllogictests/suites/duckdb/sql/aggregate/group/group_by_grouping_sets_union_all.test @@ -20,6 +20,34 @@ create or replace table cascading_grouping_sets (a UInt64, b UInt64, v Nullable( statement ok insert into cascading_grouping_sets values (1, 1, 10), (1, 2, NULL), (2, 1, 5); +# Cover the ordinary Union rewrite separately from the hierarchical alternative. +statement ok +set optimizer_skip_list = 'HierarchicalGroupingSetsToUnion,PushDownPrewhere'; + +query IIIII +select sum(v) as total, + a, + b, + grouping(a) + grouping(b) as lochierarchy, + rank() over ( + partition by grouping(a) + grouping(b), + case when grouping(b) = 0 then a end + order by sum(v) desc + ) as rank_within_parent +from cascading_grouping_sets +group by rollup(a, b) +order by lochierarchy desc, a nulls last, b nulls last; +---- +15 NULL NULL 2 1 +10 1 NULL 1 1 +5 2 NULL 1 2 +10 1 1 0 1 +NULL 1 2 0 2 +5 2 1 0 1 + +statement ok +set optimizer_skip_list = 'PushDownPrewhere'; + query IIIIIII rowsort select a, b, sum(v), sum0(v), count(v), min(v), max(v) from cascading_grouping_sets From cfd3028d1a173af9811519dd23e30b0281740efb Mon Sep 17 00:00:00 2001 From: coldWater Date: Thu, 6 Aug 2026 09:44:12 +0800 Subject: [PATCH 5/5] x --- .../mode/standalone/explain/clustering.test | 2 +- .../explain/explain_grouping_sets.test | 28 +++++++++---------- .../explain/explain_virtual_column.test | 2 +- 3 files changed, 16 insertions(+), 16 deletions(-) diff --git a/tests/sqllogictests/suites/mode/standalone/explain/clustering.test b/tests/sqllogictests/suites/mode/standalone/explain/clustering.test index 0294ceeeba5..d0d984ee630 100644 --- a/tests/sqllogictests/suites/mode/standalone/explain/clustering.test +++ b/tests/sqllogictests/suites/mode/standalone/explain/clustering.test @@ -13,7 +13,7 @@ INSERT INTO test_linear VALUES(2, 1), (2, 2); statement ok ALTER TABLE test_linear RECLUSTER FINAL; -query TT? +query TTT select * exclude(timestamp) from clustering_information('default','test_linear') ---- (a, b) linear {"average_depth":1.0,"average_overlaps":0.0,"block_depth_histogram":{"00001":2},"constant_block_count":0,"p95_depth":1,"p99_depth":1,"total_block_count":2} diff --git a/tests/sqllogictests/suites/mode/standalone/explain/explain_grouping_sets.test b/tests/sqllogictests/suites/mode/standalone/explain/explain_grouping_sets.test index 4f683ab4581..a305bd7f891 100644 --- a/tests/sqllogictests/suites/mode/standalone/explain/explain_grouping_sets.test +++ b/tests/sqllogictests/suites/mode/standalone/explain/explain_grouping_sets.test @@ -165,7 +165,7 @@ Sequence │ │ │ │ │ │ │ └── estimated rows: 100.00 │ │ │ │ │ │ └── EvalScalar │ │ │ │ │ │ ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - │ │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, NULL] + │ │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, NULL, TRY_CAST(group_item_1 (#1) AS UInt8 NULL)] │ │ │ │ │ │ ├── estimated rows: 33.33 │ │ │ │ │ │ └── AggregateFinal │ │ │ │ │ │ ├── output columns: [count(DISTINCT number) (#8), a (#1)] @@ -182,7 +182,7 @@ Sequence │ │ │ │ │ │ └── estimated rows: 100.00 │ │ │ │ │ └── EvalScalar │ │ │ │ │ ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - │ │ │ │ │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), NULL] + │ │ │ │ │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), NULL, TRY_CAST(group_item_2 (#2) AS UInt8 NULL)] │ │ │ │ │ ├── estimated rows: 33.33 │ │ │ │ │ └── AggregateFinal │ │ │ │ │ ├── output columns: [count(DISTINCT number) (#8), b (#2)] @@ -199,7 +199,7 @@ Sequence │ │ │ │ │ └── estimated rows: 100.00 │ │ │ │ └── EvalScalar │ │ │ │ ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - │ │ │ │ ├── expressions: [NULL, NULL, TRY_CAST(group_item (#3) AS UInt8 NULL)] + │ │ │ │ ├── expressions: [NULL, NULL, TRY_CAST(group_item (#3) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ │ │ │ ├── estimated rows: 33.33 │ │ │ │ └── AggregateFinal │ │ │ │ ├── output columns: [count(DISTINCT number) (#8), c (#3)] @@ -216,7 +216,7 @@ Sequence │ │ │ │ └── estimated rows: 100.00 │ │ │ └── EvalScalar │ │ │ ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), NULL] + │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), NULL, TRY_CAST(group_item_1 (#1) AS UInt8 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL)] │ │ │ ├── estimated rows: 33.33 │ │ │ └── AggregateFinal │ │ │ ├── output columns: [count(DISTINCT number) (#8), a (#1), b (#2)] @@ -233,7 +233,7 @@ Sequence │ │ │ └── estimated rows: 100.00 │ │ └── EvalScalar │ │ ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, TRY_CAST(group_item (#3) AS UInt8 NULL)] + │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, TRY_CAST(group_item (#3) AS UInt8 NULL), TRY_CAST(group_item_1 (#1) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ │ ├── estimated rows: 33.33 │ │ └── AggregateFinal │ │ ├── output columns: [count(DISTINCT number) (#8), a (#1), c (#3)] @@ -250,7 +250,7 @@ Sequence │ │ └── estimated rows: 100.00 │ └── EvalScalar │ ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL)] + │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ ├── estimated rows: 33.33 │ └── AggregateFinal │ ├── output columns: [count(DISTINCT number) (#8), b (#2), c (#3)] @@ -267,7 +267,7 @@ Sequence │ └── estimated rows: 100.00 └── EvalScalar ├── output columns: [count(DISTINCT number) (#8), a (#9), b (#10), c (#11)] - ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL)] + ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL), TRY_CAST(group_item_1 (#1) AS UInt8 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] ├── estimated rows: 33.33 └── AggregateFinal ├── output columns: [count(DISTINCT number) (#8), a (#1), b (#2), c (#3)] @@ -440,7 +440,7 @@ Sequence │ │ │ │ │ │ ├── estimated rows: 44.44 │ │ │ │ │ │ ├── EvalScalar │ │ │ │ │ │ │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ │ │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ │ │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_1 (#1) AS UInt8 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ │ │ │ │ │ │ ├── estimated rows: 33.33 │ │ │ │ │ │ │ └── MaterializeCTERef │ │ │ │ │ │ │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_1_2_3 @@ -448,7 +448,7 @@ Sequence │ │ │ │ │ │ │ └── estimated rows: 33.33 │ │ │ │ │ │ └── EvalScalar │ │ │ │ │ │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), NULL, sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), TRY_CAST(group_item (#2) AS UInt8 NULL), NULL, sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_1 (#1) AS UInt8 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL)] │ │ │ │ │ │ ├── estimated rows: 11.11 │ │ │ │ │ │ └── MaterializeCTERef │ │ │ │ │ │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_1_2 @@ -456,7 +456,7 @@ Sequence │ │ │ │ │ │ └── estimated rows: 11.11 │ │ │ │ │ └── EvalScalar │ │ │ │ │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ │ │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_1 (#1) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ │ │ │ │ ├── estimated rows: 11.11 │ │ │ │ │ └── MaterializeCTERef │ │ │ │ │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_1_3 @@ -464,7 +464,7 @@ Sequence │ │ │ │ │ └── estimated rows: 11.11 │ │ │ │ └── EvalScalar │ │ │ │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ │ │ │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ │ │ │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ │ │ │ ├── estimated rows: 11.11 │ │ │ │ └── MaterializeCTERef │ │ │ │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_2_3 @@ -472,7 +472,7 @@ Sequence │ │ │ │ └── estimated rows: 11.11 │ │ │ └── EvalScalar │ │ │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, NULL, sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ │ │ ├── expressions: [TRY_CAST(group_item (#1) AS UInt8 NULL), NULL, NULL, sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_1 (#1) AS UInt8 NULL)] │ │ │ ├── estimated rows: 11.11 │ │ │ └── MaterializeCTERef │ │ │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_1 @@ -480,7 +480,7 @@ Sequence │ │ │ └── estimated rows: 11.11 │ │ └── EvalScalar │ │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), NULL, sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ │ ├── expressions: [NULL, TRY_CAST(group_item (#2) AS UInt8 NULL), NULL, sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_2 (#2) AS UInt8 NULL)] │ │ ├── estimated rows: 11.11 │ │ └── MaterializeCTERef │ │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_2 @@ -488,7 +488,7 @@ Sequence │ │ └── estimated rows: 11.11 │ └── EvalScalar │ ├── output columns: [min(number) (#8), max(number) (#9), sum(number) (#10), count(number) (#11), a (#12), b (#13), c (#14), sum(number) / if(count(number) = 0, 1, count(number)) (#15)] - │ ├── expressions: [NULL, NULL, TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL)] + │ ├── expressions: [NULL, NULL, TRY_CAST(group_item (#3) AS UInt8 NULL), sum(number) (#10) / CAST(if(CAST(count(number) (#11) = 0 AS Boolean NULL), 1, count(number) (#11)) AS UInt64 NULL), TRY_CAST(group_item_3 (#3) AS UInt8 NULL)] │ ├── estimated rows: 11.11 │ └── MaterializeCTERef │ ├── cte_name: cte_hierarchical_groupingsets_16366510952463710337_cols_3 diff --git a/tests/sqllogictests/suites/mode/standalone/explain/explain_virtual_column.test b/tests/sqllogictests/suites/mode/standalone/explain/explain_virtual_column.test index 9e04b409906..5691ce40a21 100644 --- a/tests/sqllogictests/suites/mode/standalone/explain/explain_virtual_column.test +++ b/tests/sqllogictests/suites/mode/standalone/explain/explain_virtual_column.test @@ -787,7 +787,7 @@ CommitSink ├── apply join filters: [#0, #1] └── estimated rows: 1.00 -query TT? +query TTT SELECT * FROM data_main; ---- rec1 cat1 {"metadata":{"category_id":"cat1","timestamp":1625000000000},"timestamp":1625000000000}