From 1f133b0f7cb03afc82bd749b3909488663c56fa5 Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Sat, 22 Aug 2026 16:34:04 -0700 Subject: [PATCH] fix: avoid casting unreachable nested list values --- datafusion/common/src/nested_struct.rs | 370 +++++++++++++++++- datafusion/sqllogictest/test_files/struct.slt | 20 + 2 files changed, 382 insertions(+), 8 deletions(-) diff --git a/datafusion/common/src/nested_struct.rs b/datafusion/common/src/nested_struct.rs index 3f5e4ab4b78e2..16cebeccaa552 100644 --- a/datafusion/common/src/nested_struct.rs +++ b/datafusion/common/src/nested_struct.rs @@ -22,9 +22,10 @@ use arrow::{ GenericListViewArray, MapArray, StructArray, UInt64Array, downcast_integer, make_array, new_null_array, }, - buffer::NullBuffer, + buffer::{NullBuffer, ScalarBuffer}, compute::{CastOptions, can_cast_types, cast_with_options, take}, datatypes::{DataType, DataType::Struct, Field, FieldRef}, + error::ArrowError, }; use std::{ collections::{HashMap, HashSet}, @@ -241,6 +242,18 @@ fn cast_list_column( cast_options: &CastOptions, ) -> Result { let source_list = source_col.as_list::(); + let offsets = source_list.value_offsets(); + let needs_compaction = offsets[0] != O::usize_as(0) + || offsets[offsets.len() - 1].as_usize() != source_list.values().len() + || source_list + .offsets() + .has_non_empty_nulls(source_list.nulls()); + let compacted_list = if needs_compaction { + Some(compact_list_values(source_list)?) + } else { + None + }; + let source_list = compacted_list.as_ref().unwrap_or(source_list); let cast_values = cast_column( source_list.values(), @@ -263,23 +276,127 @@ fn cast_list_view_column( cast_options: &CastOptions, ) -> Result { let source_list = source_col.as_list_view::(); + let compacted_values = compact_list_view_values(source_list)?; + let (offsets, sizes, values) = match compacted_values.as_ref() { + Some((offsets, sizes, values)) => (offsets, sizes, values), + None => ( + source_list.offsets(), + source_list.sizes(), + source_list.values(), + ), + }; - let cast_values = cast_column( - source_list.values(), - target_inner_field.data_type(), - cast_options, - )?; + let cast_values = cast_column(values, target_inner_field.data_type(), cast_options)?; let result = GenericListViewArray::::try_new( Arc::clone(target_inner_field), - source_list.offsets().clone(), - source_list.sizes().clone(), + offsets.clone(), + sizes.clone(), cast_values, source_list.nulls().cloned(), )?; Ok(Arc::new(result)) } +fn compact_list_values( + list: &GenericListArray, +) -> Result> { + let indices = UInt64Array::from_iter_values(0..list.len() as u64); + Ok(take(list, &indices, None)?.as_list::().clone()) +} + +type CompactedListView = (ScalarBuffer, ScalarBuffer, ArrayRef); + +/// Selects the union of child ranges reachable from valid ListView rows. +/// +/// ListView ranges may overlap or appear in any order, so selecting each row +/// independently could duplicate a large number of child values. Merging the +/// ranges first preserves sharing while excluding unreachable values. +fn compact_list_view_values( + list: &GenericListViewArray, +) -> Result>> { + let mut dense_end = 0; + let mut is_dense = true; + for row in 0..list.len() { + if list.is_null(row) || list.value_sizes()[row] == O::usize_as(0) { + continue; + } + let start = list.value_offsets()[row].as_usize(); + if start != dense_end { + is_dense = false; + break; + } + dense_end = start + list.value_sizes()[row].as_usize(); + } + if is_dense && dense_end == list.values().len() { + return Ok(None); + } + + let mut ranges = list + .value_offsets() + .iter() + .zip(list.value_sizes()) + .enumerate() + .filter(|(row, (_, size))| list.is_valid(*row) && **size != O::usize_as(0)) + .map(|(_, (offset, size))| { + let start = offset.as_usize(); + (start, start + size.as_usize()) + }) + .collect::>(); + ranges.sort_unstable_by_key(|range| range.0); + + let mut merged_ranges: Vec<(usize, usize)> = Vec::with_capacity(ranges.len()); + for (start, end) in ranges { + if let Some((_, previous_end)) = merged_ranges.last_mut() + && start <= *previous_end + { + *previous_end = (*previous_end).max(end); + } else { + merged_ranges.push((start, end)); + } + } + + if merged_ranges.as_slice() == [(0, list.values().len())] { + return Ok(None); + } + + let mut range_bases = Vec::with_capacity(merged_ranges.len()); + let mut selected_len = 0; + for &(start, end) in &merged_ranges { + range_bases.push(selected_len); + selected_len += end - start; + } + + let mut offsets = Vec::with_capacity(list.len()); + let mut sizes = Vec::with_capacity(list.len()); + for row in 0..list.len() { + let size = list.value_sizes()[row]; + if list.is_null(row) || size == O::usize_as(0) { + offsets.push(O::usize_as(0)); + sizes.push(O::usize_as(0)); + continue; + } + + let source_offset = list.value_offsets()[row].as_usize(); + let range_index = + merged_ranges.partition_point(|range| range.0 <= source_offset) - 1; + let new_offset = + range_bases[range_index] + source_offset - merged_ranges[range_index].0; + offsets.push(O::from_usize(new_offset).ok_or_else(|| { + ArrowError::ComputeError("ListView offset overflow during compaction".into()) + })?); + sizes.push(size); + } + + let child_indices = UInt64Array::from_iter_values( + merged_ranges + .iter() + .flat_map(|&(start, end)| (start..end).map(|index| index as u64)), + ); + let values = take(list.values(), &child_indices, None)?; + Ok(Some((offsets.into(), sizes.into(), values))) +} + fn cast_fixed_size_list_column( source_col: &ArrayRef, target_inner_field: &FieldRef, @@ -2482,6 +2599,243 @@ mod tests { assert!(b_col.iter().all(|v| v.is_none())); } + fn list_struct_values(values: Vec<&str>) -> ArrayRef { + Arc::new(StructArray::from(vec![( + arc_field("value", DataType::Utf8), + Arc::new(StringArray::from(values)) as ArrayRef, + )])) + } + + fn list_struct_fields() -> (FieldRef, FieldRef) { + ( + arc_field("item", struct_type(vec![field("value", DataType::Utf8)])), + arc_field("item", struct_type(vec![field("value", DataType::Int32)])), + ) + } + + fn assert_visible_list_value(list: &ArrayRef, row: usize, expected: i32) { + let values = list.as_list::().value(row); + let values = values.as_struct(); + assert_eq!( + get_column_as!(values, "value", Int32Array).value(0), + expected + ); + } + + #[test] + fn test_sliced_list_ignores_unreachable_invalid_nested_value() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new( + ListArray::new( + source_field, + OffsetBuffer::new(vec![0, 1, 2].into()), + list_struct_values(vec!["bad", "2"]), + None, + ) + .slice(1, 1), + ); + let target_type = DataType::List(target_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + + assert_eq!(result.data_type(), &target_type); + assert_visible_list_value(&result, 0, 2); + } + + #[test] + fn test_null_list_parent_hides_invalid_nested_value() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new(ListArray::new( + source_field, + OffsetBuffer::new(vec![0, 1, 2].into()), + list_struct_values(vec!["bad", "2"]), + Some(NullBuffer::from(vec![false, true])), + )); + let target_type = DataType::List(target_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + + assert_eq!(result.data_type(), &target_type); + assert!(result.is_null(0)); + assert_visible_list_value(&result, 1, 2); + } + + #[test] + fn test_sliced_large_list_ignores_unreachable_invalid_nested_value() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new( + GenericListArray::::new( + source_field, + OffsetBuffer::new(vec![0, 1, 2].into()), + list_struct_values(vec!["bad", "2"]), + None, + ) + .slice(1, 1), + ); + let target_type = DataType::LargeList(target_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + let values = result.as_list::().value(0); + + assert_eq!(result.data_type(), &target_type); + assert_eq!( + get_column_as!(values.as_struct(), "value", Int32Array).value(0), + 2 + ); + } + + #[test] + fn test_list_view_ignores_unreachable_invalid_nested_values() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new(ListViewArray::new( + source_field, + ScalarBuffer::from(vec![0i32, 1, 2]), + ScalarBuffer::from(vec![1i32, 1, 1]), + list_struct_values(vec!["bad", "2", "bad"]), + Some(NullBuffer::from(vec![false, true, true])), + )); + let source_col = source_col.slice(0, 2); + let target_type = DataType::ListView(target_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + let result = result.as_list_view::(); + let values = result.value(1); + + assert_eq!(result.data_type(), &target_type); + assert!(result.is_null(0)); + assert_eq!(result.value_sizes(), &[0, 1]); + assert_eq!( + get_column_as!(values.as_struct(), "value", Int32Array).value(0), + 2 + ); + } + + #[test] + fn test_list_view_compacts_overlapping_out_of_order_ranges() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new(ListViewArray::new( + source_field, + ScalarBuffer::from(vec![2i32, 0, 2]), + ScalarBuffer::from(vec![2i32, 1, 1]), + list_struct_values(vec!["1", "bad", "2", "3", "bad"]), + None, + )); + let target_type = DataType::ListView(target_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + let result = result.as_list_view::(); + let first = result.value(0); + let second = result.value(1); + + assert_eq!(result.value_offsets(), &[1, 0, 1]); + assert_eq!(result.value_sizes(), &[2, 1, 1]); + assert_eq!( + get_column_as!(first.as_struct(), "value", Int32Array).values(), + &[2, 3] + ); + assert_eq!( + get_column_as!(second.as_struct(), "value", Int32Array).values(), + &[1] + ); + } + + #[test] + fn test_nested_sparse_list_view_compacts_backing_once() { + let (source_inner_field, target_inner_field) = list_struct_fields(); + let inner = ListViewArray::new( + source_inner_field, + ScalarBuffer::from(vec![0i32, 1, 2, 3, 5]), + ScalarBuffer::from(vec![1i32, 1, 1, 2, 1]), + list_struct_values(vec!["bad", "1", "bad", "2", "3", "bad"]), + None, + ); + let source_outer_field = arc_field("item", inner.data_type().clone()); + let target_outer_field = + arc_field("item", DataType::ListView(Arc::clone(&target_inner_field))); + let source_col: ArrayRef = Arc::new(ListViewArray::new( + source_outer_field, + ScalarBuffer::from(vec![1i32, 3]), + ScalarBuffer::from(vec![1i32, 1]), + Arc::new(inner), + None, + )); + let target_type = DataType::ListView(target_outer_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + let outer = result.as_list_view::(); + let inner = outer.values().as_list_view::(); + + assert_eq!(outer.value_offsets(), &[0, 1]); + assert_eq!(outer.value_sizes(), &[1, 1]); + assert_eq!(inner.len(), 2); + assert_eq!(inner.value_offsets(), &[0, 1]); + assert_eq!(inner.value_sizes(), &[1, 2]); + assert_eq!(inner.values().len(), 3); + + let first = outer.value(0); + let first = first.as_list_view::().value(0); + assert_eq!( + get_column_as!(first.as_struct(), "value", Int32Array).values(), + &[1] + ); + let second = outer.value(1); + let second = second.as_list_view::().value(0); + assert_eq!( + get_column_as!(second.as_struct(), "value", Int32Array).values(), + &[2, 3] + ); + } + + #[test] + fn test_sliced_large_list_view_ignores_unreachable_invalid_nested_value() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new( + GenericListViewArray::::new( + source_field, + ScalarBuffer::from(vec![0i64, 1]), + ScalarBuffer::from(vec![1i64, 1]), + list_struct_values(vec!["bad", "2"]), + None, + ) + .slice(1, 1), + ); + let target_type = DataType::LargeListView(target_field); + + let result = + cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).unwrap(); + let result = result.as_list_view::(); + let values = result.value(0); + + assert_eq!(result.data_type(), &target_type); + assert_eq!(result.value_offsets(), &[0]); + assert_eq!(result.value_sizes(), &[1]); + assert_eq!( + get_column_as!(values.as_struct(), "value", Int32Array).value(0), + 2 + ); + } + + #[test] + fn test_large_list_view_visible_invalid_nested_value_returns_error() { + let (source_field, target_field) = list_struct_fields(); + let source_col: ArrayRef = Arc::new(GenericListViewArray::::new( + source_field, + ScalarBuffer::from(vec![0i64]), + ScalarBuffer::from(vec![1i64]), + list_struct_values(vec!["bad"]), + None, + )); + let target_type = DataType::LargeListView(target_field); + + assert!(cast_column(&source_col, &target_type, &DEFAULT_CAST_OPTIONS).is_err()); + } + fn fixed_size_list_struct_field(fields: Vec<(&str, DataType)>) -> FieldRef { arc_field( "item", diff --git a/datafusion/sqllogictest/test_files/struct.slt b/datafusion/sqllogictest/test_files/struct.slt index 941bab7c2d866..a2fc523501ed2 100644 --- a/datafusion/sqllogictest/test_files/struct.slt +++ b/datafusion/sqllogictest/test_files/struct.slt @@ -1732,3 +1732,23 @@ drop view leaf_view; statement ok drop table leaf_base; + +# LIMIT can slice a List batch while retaining unreachable backing children +statement ok +create table list_cast_limit as +select * from (values + (1, [struct('1')]), + (2, [struct('2')]), + (3, [struct('bad')]) +) as t(i, l); + +query ? +select arrow_cast(l, 'List(Struct("c0": Int32))') +from list_cast_limit +limit 2; +---- +[{c0: 1}] +[{c0: 2}] + +statement ok +drop table list_cast_limit;