From 25a9b33fc1c07984912e41ab0aba1adb321b9ef2 Mon Sep 17 00:00:00 2001 From: osipovartem Date: Thu, 24 Sep 2026 14:22:23 +0300 Subject: [PATCH 1/2] Allow array_reduce accumulator to recover from NULL --- .../functions-nested/src/array_reduce.rs | 5 ++--- .../test_files/array/array_reduce.slt | 19 +++++++++++++++++-- 2 files changed, 19 insertions(+), 5 deletions(-) diff --git a/datafusion/functions-nested/src/array_reduce.rs b/datafusion/functions-nested/src/array_reduce.rs index 9a26d1a6c1c80..45a0ac3576ad6 100644 --- a/datafusion/functions-nested/src/array_reduce.rs +++ b/datafusion/functions-nested/src/array_reduce.rs @@ -271,9 +271,8 @@ impl HigherOrderUDFImpl for ArrayReduce { let mut scatter_indices = Vec::with_capacity(list.len()); for row in 0..list.len() { - let active = accumulator.is_valid(row) - && list.is_valid(row) - && position < offsets[row + 1] - offsets[row]; + let active = + list.is_valid(row) && position < offsets[row + 1] - offsets[row]; if active { source_indices.push(u64::try_from(offsets[row] + position).map_err( |error| internal_datafusion_err!("invalid list index: {error}"), diff --git a/datafusion/sqllogictest/test_files/array/array_reduce.slt b/datafusion/sqllogictest/test_files/array/array_reduce.slt index 23f8072eeb190..e9da518a17b9d 100644 --- a/datafusion/sqllogictest/test_files/array/array_reduce.slt +++ b/datafusion/sqllogictest/test_files/array/array_reduce.slt @@ -62,7 +62,7 @@ SELECT array_reduce([1, NULL, 2], 0, (acc, value) -> acc + coalesce(value, 0)); ---- 3 -# A null merge result is terminal and cannot be recovered by a later element. +# A later element may recover a null accumulator. query I SELECT array_reduce( [1, 2], @@ -70,7 +70,12 @@ SELECT array_reduce( (acc, value) -> CASE WHEN value = 1 THEN NULL ELSE coalesce(acc, 0) + value END ); ---- -NULL +2 + +query I +SELECT array_reduce([1, 2], NULL::BIGINT, (acc, value) -> coalesce(acc, 0) + value); +---- +3 query R SELECT array_reduce([1.2, 2.3], 0, (acc, value) -> acc + value); @@ -94,5 +99,15 @@ FROM reduce_t; 7 NULL +query I rowsort +SELECT array_reduce(values, initial, + (acc, value) -> CASE WHEN value = 1 THEN NULL ELSE coalesce(acc, 0) + value END) +FROM reduce_t; +---- +5 +7 +9 +NULL + statement ok set datafusion.sql_parser.dialect = generic; From aabecd0f2894f85eef09e5aa1370806b6efec52d Mon Sep 17 00:00:00 2001 From: osipovartem Date: Thu, 24 Sep 2026 14:25:41 +0300 Subject: [PATCH 2/2] Cover null element and captured-column recovery --- .../test_files/array/array_reduce.slt | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/datafusion/sqllogictest/test_files/array/array_reduce.slt b/datafusion/sqllogictest/test_files/array/array_reduce.slt index e9da518a17b9d..dba41ca4c6797 100644 --- a/datafusion/sqllogictest/test_files/array/array_reduce.slt +++ b/datafusion/sqllogictest/test_files/array/array_reduce.slt @@ -77,6 +77,11 @@ SELECT array_reduce([1, 2], NULL::BIGINT, (acc, value) -> coalesce(acc, 0) + val ---- 3 +query I +SELECT array_reduce([1, NULL, 2], 0, (acc, value) -> coalesce(acc, 0) + value); +---- +2 + query R SELECT array_reduce([1.2, 2.3], 0, (acc, value) -> acc + value); ---- @@ -109,5 +114,15 @@ FROM reduce_t; 9 NULL +query I rowsort +SELECT array_reduce(values, initial, + (acc, value) -> CASE WHEN value = 1 THEN NULL ELSE coalesce(acc, 0) + value + extra END) +FROM reduce_t; +---- +25 +29 +7 +NULL + statement ok set datafusion.sql_parser.dialect = generic;