Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 20 additions & 34 deletions datafusion/functions-window/src/lead_lag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -311,21 +311,24 @@ impl WindowUDFImpl for WindowShift {
}

fn limit_effect(&self, args: &[Arc<dyn PhysicalExpr>]) -> LimitEffect {
if self.kind == WindowShiftKind::Lag {
return LimitEffect::None;
}
match args {
let amount = match args {
[_, expr, ..] => {
let Some(lit) = expr.downcast_ref::<expressions::Literal>() else {
return LimitEffect::Unknown;
};
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
}
}
}
Expand Down Expand Up @@ -655,39 +658,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
Expand Down
9 changes: 8 additions & 1 deletion datafusion/physical-plan/src/windows/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -295,7 +295,14 @@ impl StandardWindowFunctionExpr for WindowUDFExpr {
}

fn limit_effect(&self) -> LimitEffect {
self.fun.inner().limit_effect(self.args.as_slice())
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::Relative(_) | LimitEffect::Absolute(_) if self.ignore_nulls => {
LimitEffect::Unknown
}
effect => effect,
}
}
}

Expand Down
57 changes: 57 additions & 0 deletions datafusion/sqllogictest/test_files/window.slt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
98 changes: 98 additions & 0 deletions datafusion/sqllogictest/test_files/window_limits.slt
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,104 @@ 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

# 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;
Expand Down
Loading