From 85115d26b087cdfbe4f768b4adb70065d40f45c4 Mon Sep 17 00:00:00 2001 From: Christian Del Monte Date: Fri, 28 Aug 2026 12:38:19 +0200 Subject: [PATCH 1/2] fix: prevent CSE from extracting window functions --- .../optimizer/src/common_subexpr_eliminate.rs | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/datafusion/optimizer/src/common_subexpr_eliminate.rs b/datafusion/optimizer/src/common_subexpr_eliminate.rs index 94d61b74879c6..cac06befb8241 100644 --- a/datafusion/optimizer/src/common_subexpr_eliminate.rs +++ b/datafusion/optimizer/src/common_subexpr_eliminate.rs @@ -732,6 +732,7 @@ impl CSEController for ExprCSEController<'_> { | Expr::Wildcard { .. } | Expr::Lambda(_) | Expr::LambdaVariable(_) + | Expr::WindowFunction(..) ); let is_aggr = matches!(node, Expr::AggregateFunction(..)); @@ -856,6 +857,7 @@ mod test { use crate::test::udfs::leaf_udf_expr; use crate::test::*; use datafusion_expr::test::function_stub::{avg, sum}; + use datafusion_functions_window::row_number::row_number_udwf; macro_rules! assert_optimized_plan_equal { ( @@ -1932,4 +1934,21 @@ mod test { " ) } + + #[test] + fn test_window_function_is_not_extracted() -> Result<()> { + let wexpr = row_number_udwf().call(vec![]); + + let plan = LogicalPlanBuilder::empty(true) + .window(vec![wexpr.clone(), wexpr.alias("aliased")])? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + WindowAggr: windowExpr=[[row_number() ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING, row_number() ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS aliased]] + EmptyRelation: rows=1 + " + ) + } } From 99327920a6c2a870bb752fb40f1cff07cf64f7e8 Mon Sep 17 00:00:00 2001 From: Christian Del Monte Date: Wed, 2 Sep 2026 10:29:13 +0200 Subject: [PATCH 2/2] test: add execution regression for window CSE --- datafusion/core/tests/dataframe/mod.rs | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/datafusion/core/tests/dataframe/mod.rs b/datafusion/core/tests/dataframe/mod.rs index 44ab14e6cc137..83adf1b019c06 100644 --- a/datafusion/core/tests/dataframe/mod.rs +++ b/datafusion/core/tests/dataframe/mod.rs @@ -206,6 +206,23 @@ async fn with_column_window_functions() -> DataFusionResult<()> { Ok(()) } +#[tokio::test] +async fn duplicated_window_functions_can_be_executed() -> Result<()> { + let wexpr = datafusion::functions_window::row_number::row_number_udwf().call(vec![]); + + let plan = LogicalPlanBuilder::empty(true) + .window(vec![wexpr.clone(), wexpr.alias("aliased")])? + .build()?; + + let ctx = SessionContext::new(); + + let collected = DataFrame::new(ctx.state(), plan).collect().await?; + + assert_eq!(collected.iter().map(|b| b.num_rows()).sum::(), 1); + + Ok(()) +} + #[tokio::test] async fn test_coalesce_schema() -> Result<()> { let ctx = SessionContext::new();