From 922c1cf45adb7ff33253194a6253dccfc5ca1618 Mon Sep 17 00:00:00 2001 From: linfeng Date: Wed, 16 Sep 2026 20:26:20 +0800 Subject: [PATCH 1/2] fix: correct LEAD/LAG IGNORE NULLS evaluation and limit pushdown --- datafusion/functions-window/src/lead_lag.rs | 37 ++++-------- datafusion/physical-plan/src/windows/mod.rs | 8 ++- datafusion/sqllogictest/test_files/window.slt | 57 +++++++++++++++++++ .../sqllogictest/test_files/window_limits.slt | 40 +++++++++++++ 4 files changed, 114 insertions(+), 28 deletions(-) diff --git a/datafusion/functions-window/src/lead_lag.rs b/datafusion/functions-window/src/lead_lag.rs index fea4a1a4aadda..37a7727a66dcc 100644 --- a/datafusion/functions-window/src/lead_lag.rs +++ b/datafusion/functions-window/src/lead_lag.rs @@ -655,39 +655,22 @@ impl PartitionEvaluator for WindowShiftEvaluator { // Stores the necessary non-null entry number further than the current row. let non_null_row_count = offset_magnitude(self.shift_offset); - if self.non_null_offsets.is_empty() { - // When empty, fill non_null offsets with the data further than the current row. - let mut offset_val = 1; - for idx in range.start + 1..range.end { - if array.is_valid(idx) { - self.non_null_offsets.push_back(offset_val); - offset_val = 1; - } else { - offset_val += 1; - } - // It is enough to keep track of `non_null_row_count + 1` non-null offset. - // further data is unnecessary for the result. - if self.non_null_offsets.len() == non_null_row_count.saturating_add(1) - { - break; - } + // Resume after the last cached non-null row, even when the cache is + // non-empty but does not yet contain enough rows for this offset. + let mut total_offset: usize = self.non_null_offsets.iter().sum(); + for next_idx in range.start + total_offset + 1..range.end { + if self.non_null_offsets.len() == non_null_row_count { + break; } - } else if range.end < len && array.is_valid(range.end) { - // Update `non_null_offsets` with the new end data. - if array.is_valid(range.end) { - // When non-null, append a new offset. - self.non_null_offsets.push_back(1); - } else { - // When null, increment offset count of the last entry - let last_idx = self.non_null_offsets.len() - 1; - self.non_null_offsets[last_idx] += 1; + if array.is_valid(next_idx) { + let next_offset = next_idx - range.start; + self.non_null_offsets.push_back(next_offset - total_offset); + total_offset = next_offset; } } // Find the nonNULL row index that shifted by offset comparing to current row index idx = if self.non_null_offsets.len() >= non_null_row_count { - let total_offset: usize = - self.non_null_offsets.iter().take(non_null_row_count).sum(); Some(range.start + total_offset) } else { None diff --git a/datafusion/physical-plan/src/windows/mod.rs b/datafusion/physical-plan/src/windows/mod.rs index 7c5f55f661f38..7a242d774f78c 100644 --- a/datafusion/physical-plan/src/windows/mod.rs +++ b/datafusion/physical-plan/src/windows/mod.rs @@ -295,7 +295,13 @@ impl StandardWindowFunctionExpr for WindowUDFExpr { } fn limit_effect(&self) -> LimitEffect { - self.fun.inner().limit_effect(self.args.as_slice()) + if self.ignore_nulls { + // The function's offset counts non-null values, so it cannot bound + // the number of input rows needed when NULLs are skipped. + LimitEffect::Unknown + } else { + self.fun.inner().limit_effect(self.args.as_slice()) + } } } diff --git a/datafusion/sqllogictest/test_files/window.slt b/datafusion/sqllogictest/test_files/window.slt index 6ffe1b4cc087c..e2adb33bc1517 100644 --- a/datafusion/sqllogictest/test_files/window.slt +++ b/datafusion/sqllogictest/test_files/window.slt @@ -4482,6 +4482,63 @@ NULL 19 statement ok set datafusion.execution.batch_size = 100; +statement ok +CREATE TABLE shifts AS VALUES + (1, 10), (2, 20), (3, 30), (4, 40), + (5, NULL), (6, NULL), (7, 70), (8, 80), (9, NULL); + +# A partially populated lookahead cache must scan past NULLs to find +# the requested non-null row. Negative offsets reverse the direction. +query IIIIII +SELECT column1, column2, + LEAD(column2, 2) IGNORE NULLS OVER w, + LAG(column2, -2, -1) IGNORE NULLS OVER w, + LAG(column2, 2) IGNORE NULLS OVER w, + LEAD(column2, -2) IGNORE NULLS OVER w +FROM shifts +WINDOW w AS (ORDER BY column1) +ORDER BY column1; +---- +1 10 30 30 NULL NULL +2 20 40 40 NULL NULL +3 30 70 70 10 10 +4 40 80 80 20 20 +5 NULL 80 80 30 30 +6 NULL 80 80 30 30 +7 70 NULL -1 30 30 +8 80 NULL -1 40 40 +9 NULL NULL -1 70 70 + +# The result must be the same when input arrives in separate batches. +statement ok +SET datafusion.execution.batch_size = 1; + +query IIIIII +SELECT column1, column2, + LEAD(column2, 2) IGNORE NULLS OVER w, + LAG(column2, -2, -1) IGNORE NULLS OVER w, + LAG(column2, 2) IGNORE NULLS OVER w, + LEAD(column2, -2) IGNORE NULLS OVER w +FROM shifts +WINDOW w AS (ORDER BY column1) +ORDER BY column1; +---- +1 10 30 30 NULL NULL +2 20 40 40 NULL NULL +3 30 70 70 10 10 +4 40 80 80 20 20 +5 NULL 80 80 30 30 +6 NULL 80 80 30 30 +7 70 NULL -1 30 30 +8 80 NULL -1 40 40 +9 NULL NULL -1 70 70 + +statement ok +SET datafusion.execution.batch_size = 100; + +statement ok +DROP TABLE shifts; + # Tests schema and data are in sync for mixed nulls and not nulls values for builtin window function query T select lag(a) over (order by a ASC NULLS FIRST) as x1 diff --git a/datafusion/sqllogictest/test_files/window_limits.slt b/datafusion/sqllogictest/test_files/window_limits.slt index 1498aa68aa251..2bb3a1c62331c 100644 --- a/datafusion/sqllogictest/test_files/window_limits.slt +++ b/datafusion/sqllogictest/test_files/window_limits.slt @@ -112,6 +112,46 @@ physical_plan 04)------SortExec: TopK(fetch=5), expr=[empno@0 ASC NULLS LAST], preserve_partitioning=[false] 05)--------DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/testing/data/csv/aggregate_test_100_with_dates.csv]]}, projection=[empno], file_type=csv, has_header=true +# Non-null offsets cannot be used as physical row counts for LIMIT pushdown. +statement ok +SET datafusion.optimizer.enable_window_limits = false; + +query III +SELECT id, + LEAD(v) IGNORE NULLS OVER w, + LAG(v, -1, -1) IGNORE NULLS OVER w +FROM (VALUES (1, 10), (2, NULL), (3, NULL), (4, 40)) AS t(id, v) +WINDOW w AS (ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) +ORDER BY id +LIMIT 1; +---- +1 40 40 + +statement ok +SET datafusion.optimizer.enable_window_limits = true; + +query III +SELECT id, + LEAD(v) IGNORE NULLS OVER w, + LAG(v, -1, -1) IGNORE NULLS OVER w +FROM (VALUES (1, 10), (2, NULL), (3, NULL), (4, 40)) AS t(id, v) +WINDOW w AS (ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) +ORDER BY id +LIMIT 1; +---- +1 40 40 + +# Negative-offset LAG must also block pushdown when no LEAD is present. +query II +SELECT id, LAG(v, -1) IGNORE NULLS OVER ( + ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW +) +FROM (VALUES (1, 10), (2, NULL), (3, NULL), (4, 40)) AS t(id, v) +ORDER BY id +LIMIT 1; +---- +1 40 + # Should use the max of leads statement ok set datafusion.optimizer.enable_window_limits = false; From b5b535a5c4f55fcb73e13272b3112dc0adaadfe0 Mon Sep 17 00:00:00 2001 From: linfeng Date: Mon, 21 Sep 2026 22:58:31 +0800 Subject: [PATCH 2/2] preserve safe limit pushdown for window functions --- datafusion/functions-window/src/lead_lag.rs | 17 +++--- datafusion/physical-plan/src/windows/mod.rs | 9 +-- .../sqllogictest/test_files/window_limits.slt | 58 +++++++++++++++++++ 3 files changed, 73 insertions(+), 11 deletions(-) diff --git a/datafusion/functions-window/src/lead_lag.rs b/datafusion/functions-window/src/lead_lag.rs index 37a7727a66dcc..3f13170c03de1 100644 --- a/datafusion/functions-window/src/lead_lag.rs +++ b/datafusion/functions-window/src/lead_lag.rs @@ -311,10 +311,7 @@ impl WindowUDFImpl for WindowShift { } fn limit_effect(&self, args: &[Arc]) -> LimitEffect { - if self.kind == WindowShiftKind::Lag { - return LimitEffect::None; - } - match args { + let amount = match args { [_, expr, ..] => { let Some(lit) = expr.downcast_ref::() else { return LimitEffect::Unknown; @@ -322,10 +319,16 @@ impl WindowUDFImpl for WindowShift { let ScalarValue::Int64(Some(amount)) = lit.value() else { return LimitEffect::Unknown; // we should only get int64 from the parser }; - LimitEffect::Relative((*amount).max(0) as usize) + *amount } - [_] => LimitEffect::Relative(1), // default value - _ => LimitEffect::Unknown, // invalid arguments + [_] => 1, // default value + _ => return LimitEffect::Unknown, // invalid arguments + }; + let shift_offset = self.kind.shift_offset(Some(amount)); + if shift_offset < 0 { + LimitEffect::Relative(offset_magnitude(shift_offset)) + } else { + LimitEffect::None } } } diff --git a/datafusion/physical-plan/src/windows/mod.rs b/datafusion/physical-plan/src/windows/mod.rs index 7a242d774f78c..aeb49010d281d 100644 --- a/datafusion/physical-plan/src/windows/mod.rs +++ b/datafusion/physical-plan/src/windows/mod.rs @@ -295,12 +295,13 @@ impl StandardWindowFunctionExpr for WindowUDFExpr { } fn limit_effect(&self) -> LimitEffect { - if self.ignore_nulls { + match self.fun.inner().limit_effect(self.args.as_slice()) { // The function's offset counts non-null values, so it cannot bound // the number of input rows needed when NULLs are skipped. - LimitEffect::Unknown - } else { - self.fun.inner().limit_effect(self.args.as_slice()) + LimitEffect::Relative(_) | LimitEffect::Absolute(_) if self.ignore_nulls => { + LimitEffect::Unknown + } + effect => effect, } } } diff --git a/datafusion/sqllogictest/test_files/window_limits.slt b/datafusion/sqllogictest/test_files/window_limits.slt index 2bb3a1c62331c..11c96904f4ea2 100644 --- a/datafusion/sqllogictest/test_files/window_limits.slt +++ b/datafusion/sqllogictest/test_files/window_limits.slt @@ -152,6 +152,64 @@ LIMIT 1; ---- 1 40 +# Negative-offset LAG needs lookahead even without IGNORE NULLS. +query II +SELECT id, LAG(v, -1) OVER ( + ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW +) +FROM (VALUES (1, 10), (2, 20), (3, 30), (4, 40)) AS t(id, v) +ORDER BY id +LIMIT 1; +---- +1 20 + +# Causal IGNORE NULLS functions should retain TopK below the window. +query IIIIIII +SELECT id, + LAG(v) IGNORE NULLS OVER w, + LEAD(v, -1) IGNORE NULLS OVER w, + LEAD(v, 0) IGNORE NULLS OVER w, + FIRST_VALUE(v) IGNORE NULLS OVER w, + LAST_VALUE(v) IGNORE NULLS OVER w, + NTH_VALUE(v, 2) IGNORE NULLS OVER w +FROM (VALUES (1, 10), (2, NULL), (3, 30), (4, 40)) AS t(id, v) +WINDOW w AS (ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) +ORDER BY id +LIMIT 3; +---- +1 NULL NULL 10 10 10 NULL +2 10 10 NULL 10 10 NULL +3 10 10 30 10 30 30 + +query TT +EXPLAIN +SELECT id, + LAG(v) IGNORE NULLS OVER w, + LEAD(v, -1) IGNORE NULLS OVER w, + LEAD(v, 0) IGNORE NULLS OVER w, + FIRST_VALUE(v) IGNORE NULLS OVER w, + LAST_VALUE(v) IGNORE NULLS OVER w, + NTH_VALUE(v, 2) IGNORE NULLS OVER w +FROM (VALUES (1, 10), (2, NULL), (3, 30), (4, 40)) AS t(id, v) +WINDOW w AS (ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) +ORDER BY id +LIMIT 3; +---- +logical_plan +01)Sort: t.id ASC NULLS LAST, fetch=3 +02)--Projection: t.id, lag(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v,Int64(-1)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v,Int64(0)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, first_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, last_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, nth_value(t.v,Int64(2)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW +03)----WindowAggr: windowExpr=[[lag(t.v)IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v, Int64(-1))IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v, Int64(0))IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, first_value(t.v)IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, last_value(t.v)IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, nth_value(t.v, Int64(2))IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] +04)------SubqueryAlias: t +05)--------Projection: column1 AS id, column2 AS v +06)----------Values: (Int64(1), Int64(10)), (Int64(2), Int64(NULL)), (Int64(3), Int64(30)), (Int64(4), Int64(40)) +physical_plan +01)ProjectionExec: expr=[id@0 as id, lag(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@2 as lag(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v,Int64(-1)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as lead(t.v,Int64(-1)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v,Int64(0)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@4 as lead(t.v,Int64(0)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, first_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@5 as first_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, last_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@6 as last_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, nth_value(t.v,Int64(2)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@7 as nth_value(t.v,Int64(2)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW] +02)--GlobalLimitExec: skip=0, fetch=3 +03)----BoundedWindowAggExec: wdw=[lag(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Int64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v,Int64(-1)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lead(t.v,Int64(-1)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Int64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, lead(t.v,Int64(0)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lead(t.v,Int64(0)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Int64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, first_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "first_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Int64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, last_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "last_value(t.v) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Int64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, nth_value(t.v,Int64(2)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "nth_value(t.v,Int64(2)) IGNORE NULLS ORDER BY [t.id ASC NULLS LAST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Int64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] +04)------ProjectionExec: expr=[column1@0 as id, column2@1 as v] +05)--------SortExec: TopK(fetch=3), expr=[column1@0 ASC NULLS LAST], preserve_partitioning=[false] +06)----------DataSourceExec: partitions=1, partition_sizes=[1] + # Should use the max of leads statement ok set datafusion.optimizer.enable_window_limits = false;