fix(query): validate DistAnalyzeExec child count (#8510)

* fix(query): validate DistAnalyzeExec child count

Signed-off-by: discord9 <discord9@163.com>

* test(query): avoid implicit clone lint

Signed-off-by: discord9 <discord9@163.com>

---------

Signed-off-by: discord9 <discord9@163.com>
This commit is contained in:
discord9
2026-07-15 09:59:24 +08:00
committed by GitHub
parent 2123108db0
commit 623145e635
+75 -2
View File
@@ -33,7 +33,7 @@ use datafusion::physical_plan::{
DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties, accept,
};
use datafusion_common::tree_node::{TreeNode, TreeNodeRecursion};
use datafusion_common::{DataFusionError, internal_err};
use datafusion_common::{DataFusionError, assert_eq_or_internal_err, internal_err};
use datafusion_physical_expr::{Distribution, EquivalenceProperties, Partitioning};
use futures::StreamExt;
use serde::Serialize;
@@ -172,8 +172,13 @@ impl ExecutionPlan for DistAnalyzeExec {
self: Arc<Self>,
mut children: Vec<Arc<dyn ExecutionPlan>>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
assert_eq_or_internal_err!(
children.len(),
1,
"DistAnalyzeExec requires exactly one child"
);
Ok(Arc::new(Self::new(
children.pop().unwrap(),
children.swap_remove(0),
self.verbose,
self.format,
)))
@@ -381,3 +386,71 @@ impl Display for JsonMetrics {
write!(f, "{}", serde_json::to_string(self).unwrap())
}
}
#[cfg(test)]
mod tests {
use datafusion::physical_plan::empty::EmptyExec;
use super::*;
fn empty_plan(name: &str) -> Arc<dyn ExecutionPlan> {
Arc::new(EmptyExec::new(Arc::new(Schema::new(vec![Field::new(
name,
DataType::Utf8,
true,
)]))))
}
#[test]
fn qbs_dist_analyze_rejects_zero_children() {
let analyze = Arc::new(DistAnalyzeExec::new(
empty_plan("original"),
false,
AnalyzeFormat::TEXT,
));
assert!(ExecutionPlan::with_new_children(analyze, vec![]).is_err());
}
#[test]
fn qbs_dist_analyze_rejects_multiple_children() {
let analyze = Arc::new(DistAnalyzeExec::new(
empty_plan("original"),
false,
AnalyzeFormat::TEXT,
));
let result = ExecutionPlan::with_new_children(
analyze,
vec![empty_plan("first"), empty_plan("second")],
);
if let Ok(plan) = result {
let retained = plan
.as_any()
.downcast_ref::<DistAnalyzeExec>()
.unwrap()
.input()
.schema()
.field(0)
.name()
.clone();
panic!("expected an arity error for multiple children, but retained `{retained}`");
}
}
#[test]
fn qbs_dist_analyze_accepts_exactly_one_child() {
let analyze = Arc::new(DistAnalyzeExec::new(
empty_plan("original"),
false,
AnalyzeFormat::TEXT,
));
let replacement = empty_plan("replacement");
let rebuilt = ExecutionPlan::with_new_children(analyze, vec![replacement]).unwrap();
let rebuilt = rebuilt.as_any().downcast_ref::<DistAnalyzeExec>().unwrap();
assert_eq!(rebuilt.input().schema().field(0).name(), "replacement");
}
}