diff --git a/datafusion/physical-plan/src/coalesce_partitions.rs b/datafusion/physical-plan/src/coalesce_partitions.rs index 9e3811e0ada76..d3de70a3d95f0 100644 --- a/datafusion/physical-plan/src/coalesce_partitions.rs +++ b/datafusion/physical-plan/src/coalesce_partitions.rs @@ -304,7 +304,11 @@ impl ExecutionPlan for CoalescePartitionsExec { } fn cardinality_effect(&self) -> CardinalityEffect { - CardinalityEffect::Equal + if self.fetch.is_none() { + CardinalityEffect::Equal + } else { + CardinalityEffect::LowerEqual + } } /// Tries to swap `projection` with its input, which is known to be a @@ -698,4 +702,21 @@ mod tests { Ok(()) } + + #[test] + fn test_cardinality_effect_with_fetch() { + let input = test::scan_partitioned(4); + + let coalesce = CoalescePartitionsExec::new(Arc::clone(&input)); + assert!(matches!( + coalesce.cardinality_effect(), + CardinalityEffect::Equal + )); + + let coalesce = CoalescePartitionsExec::new(input).with_fetch(Some(10)); + assert!(matches!( + coalesce.cardinality_effect(), + CardinalityEffect::LowerEqual + )); + } } diff --git a/datafusion/physical-plan/src/operator_statistics/mod.rs b/datafusion/physical-plan/src/operator_statistics/mod.rs index e09aa6b86a1a7..49a2baede9242 100644 --- a/datafusion/physical-plan/src/operator_statistics/mod.rs +++ b/datafusion/physical-plan/src/operator_statistics/mod.rs @@ -1913,6 +1913,22 @@ mod tests { Ok(()) } + #[test] + fn test_passthrough_skips_coalesce_partitions_with_fetch() -> Result<()> { + use crate::coalesce_partitions::CoalescePartitionsExec; + + let registry = StatisticsRegistry::with_providers(vec![Arc::new( + PassthroughStatisticsProvider, + )]); + let source = make_source(1000); + let coalesce: Arc = + Arc::new(CoalescePartitionsExec::new(source).with_fetch(Some(10))); + + let stats = compute(®istry, coalesce.as_ref())?; + assert_eq!(stats.base.num_rows, Precision::Exact(10)); + Ok(()) + } + #[test] fn test_chain_priority() -> Result<()> { let mut registry = StatisticsRegistry::new();