diff --git a/native/vendor/datafusion-physical-plan/src/sorts/sort/late_materialize.rs b/native/vendor/datafusion-physical-plan/src/sorts/sort/late_materialize.rs index c73d914f4db..a0933de1ac6 100644 --- a/native/vendor/datafusion-physical-plan/src/sorts/sort/late_materialize.rs +++ b/native/vendor/datafusion-physical-plan/src/sorts/sort/late_materialize.rs @@ -382,46 +382,53 @@ fn sliced_bytes(array: &dyn Array) -> Result { Ok(bytes) } -fn offsets_fit(arrays: &[ArrayData]) -> bool { +fn offsets_fit(arrays: &[ArrayRef]) -> bool { match arrays[0].data_type() { DataType::List(_) | DataType::Map(_, _) => { let mut total = 0usize; let mut children = Vec::with_capacity(arrays.len()); - for data in arrays { - let offsets = - &data.buffer::(0)[data.offset()..=data.offset() + data.len()]; + for array in arrays { + let (offsets, values) = match array.as_list_opt::() { + Some(list) => (list.value_offsets(), Arc::clone(list.values())), + None => { + let map = array.as_map(); + ( + map.value_offsets(), + Arc::new(map.entries().clone()) as ArrayRef, + ) + } + }; let start = offsets[0] as usize; - let end = offsets[data.len()] as usize; + let end = offsets[array.len()] as usize; total += end - start; - children.push(data.child_data()[0].slice(start, end - start)); + children.push(values.slice(start, end - start)); } total <= i32::MAX as usize && offsets_fit(&children) } DataType::LargeList(_) => { - let mut children = Vec::with_capacity(arrays.len()); - for data in arrays { - let offsets = - &data.buffer::(0)[data.offset()..=data.offset() + data.len()]; - let start = offsets[0] as usize; - let end = offsets[data.len()] as usize; - children.push(data.child_data()[0].slice(start, end - start)); - } + let children: Vec = arrays + .iter() + .map(|array| { + let list = array.as_list::(); + let offsets = list.value_offsets(); + let start = offsets[0] as usize; + let end = offsets[list.len()] as usize; + list.values().slice(start, end - start) + }) + .collect(); offsets_fit(&children) } - DataType::FixedSizeList(_, size) => { - let size = *size as usize; - let children: Vec = arrays + DataType::FixedSizeList(_, _) => { + let children: Vec = arrays .iter() - .map(|data| { - data.child_data()[0].slice(data.offset() * size, data.len() * size) - }) + .map(|array| Arc::clone(array.as_fixed_size_list().values())) .collect(); offsets_fit(&children) } DataType::Struct(fields) => (0..fields.len()).all(|field| { - let children: Vec = arrays + let children: Vec = arrays .iter() - .map(|data| data.child_data()[field].slice(data.offset(), data.len())) + .map(|array| Arc::clone(array.as_struct().column(field))) .collect(); offsets_fit(&children) }), @@ -645,11 +652,9 @@ impl Gather { } fn merge(&mut self, arrays: &[ArrayRef], bytes: usize) -> Option { - let data: Vec = arrays.iter().map(|array| array.to_data()).collect(); - if !offsets_fit(&data) { + if !offsets_fit(arrays) { return None; } - drop(data); let needed = self.needed() + bytes; let size = self.reservation.size(); if size < needed && self.reservation.try_grow(needed - size).is_err() { @@ -1680,12 +1685,12 @@ mod tests { .add_child_data(ArrayData::new_null(&DataType::Null, len as usize)) .build() }; - let small = list(1)?; - let big = list(i32::MAX - 1)?; - assert!(offsets_fit(&[small.clone(), small.clone()])); - assert!(offsets_fit(&[big.clone(), small.clone()])); - assert!(!offsets_fit(&[big.clone(), small.clone(), small])); - assert!(!offsets_fit(&[big.clone(), big])); + let small = arrow::array::make_array(list(1)?); + let big = arrow::array::make_array(list(i32::MAX - 1)?); + assert!(offsets_fit(&[Arc::clone(&small), Arc::clone(&small)])); + assert!(offsets_fit(&[Arc::clone(&big), Arc::clone(&small)])); + assert!(!offsets_fit(&[Arc::clone(&big), Arc::clone(&small), small])); + assert!(!offsets_fit(&[Arc::clone(&big), big])); Ok(()) } @@ -1753,4 +1758,176 @@ mod tests { } Ok(()) } + + fn list_of(values: ArrayRef, lengths: Vec) -> ArrayRef { + Arc::new(ListArray::new( + Arc::new(Field::new_list_field(values.data_type().clone(), true)), + arrow::buffer::OffsetBuffer::from_lengths(lengths), + values, + None, + )) + } + + #[test] + fn offsets_fit_reads_a_list_nested_in_a_sliced_list() { + let inner = ListArray::from_iter_primitive::( + (0..17).map(|i| Some(vec![Some(i)])), + ); + let outer = list_of(Arc::new(inner), vec![17, 0]); + let tail = outer.slice(1, 1); + assert!(offsets_fit(&[Arc::clone(&tail), tail])); + } + + #[test] + fn offsets_fit_reads_a_struct_nested_in_a_sliced_list() { + let item = arrow::array::StructArray::from(vec![( + Arc::new(Field::new("x", DataType::Int32, true)), + Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef, + )]); + let list = list_of(Arc::new(item), vec![2, 1]); + let tail = list.slice(1, 1); + assert!(offsets_fit(&[Arc::clone(&tail), tail, list])); + } + + #[test] + fn nested_list_offsets_that_would_overflow_are_not_concatenated() -> Result<()> { + let list = + |len: usize| list_of(Arc::new(arrow::array::NullArray::new(len)), vec![len]); + let nested = |len: usize| -> ArrayRef { + let wrapped = arrow::array::StructArray::from(vec![( + Arc::new(Field::new("l", list(len).data_type().clone(), true)), + list(len), + )]); + list_of(Arc::new(wrapped), vec![1]) + }; + let small = nested(1); + let big = nested(i32::MAX as usize - 1); + assert!(offsets_fit(&[Arc::clone(&big), Arc::clone(&small)])); + assert!(!offsets_fit(&[Arc::clone(&big), Arc::clone(&small), small])); + assert!(!offsets_fit(&[Arc::clone(&big), big])); + Ok(()) + } + + fn sliced_nested_batches(count: usize, rows: usize, seed: u64) -> Vec { + use arrow::array::{MapArray, NullBufferBuilder, StructArray}; + let total = count * rows; + let mut state = seed | 1; + let mut next = move || { + state ^= state << 13; + state ^= state >> 7; + state ^= state << 17; + state + }; + let mut keys: Vec = (0..total as i32).collect(); + for i in (1..total).rev() { + keys.swap(i, (next() % (i as u64 + 1)) as usize); + } + let (mut item_lengths, mut item_nulls) = (vec![], NullBufferBuilder::new(total)); + let (mut xs, mut tag_lengths, mut tags) = (vec![], vec![], vec![]); + let (mut attr_lengths, mut attr_keys, mut attr_values) = (vec![], vec![], vec![]); + let mut pad = vec![]; + for row in 0..total { + let v = next(); + if v % 7 == 0 { + item_lengths.push(0); + item_nulls.append_null(); + } else { + let items = (v % 4) as usize; + item_lengths.push(items); + item_nulls.append_non_null(); + for item in 0..items { + xs.push((v >> (item * 8)) as i32); + let n = ((v >> (item * 4)) % 3) as usize; + tag_lengths.push(n); + for t in 0..n { + tags.push(format!("tag-{row}-{item}-{t}")); + } + } + } + let attrs = ((v >> 32) % 3) as usize; + attr_lengths.push(attrs); + for a in 0..attrs { + attr_keys.push(format!("key{a}")); + attr_values.push((v % 5 != 0).then(|| format!("value-{row}-{a}"))); + } + pad.push(format!("{row:0>80}")); + } + let tags = list_of(Arc::new(StringArray::from(tags)), tag_lengths); + let item = StructArray::from(vec![ + ( + Arc::new(Field::new("x", DataType::Int32, true)), + Arc::new(Int32Array::from(xs)) as ArrayRef, + ), + ( + Arc::new(Field::new("tags", tags.data_type().clone(), true)), + tags, + ), + ]); + let items: ArrayRef = Arc::new(ListArray::new( + Arc::new(Field::new_list_field(item.data_type().clone(), true)), + arrow::buffer::OffsetBuffer::from_lengths(item_lengths), + Arc::new(item), + item_nulls.finish(), + )); + let entries = StructArray::from(vec![ + ( + Arc::new(Field::new("key", DataType::Utf8, false)), + Arc::new(StringArray::from(attr_keys)) as ArrayRef, + ), + ( + Arc::new(Field::new("value", DataType::Utf8, true)), + Arc::new(StringArray::from(attr_values)) as ArrayRef, + ), + ]); + let attrs: ArrayRef = Arc::new(MapArray::new( + Arc::new(Field::new("entries", entries.data_type().clone(), false)), + arrow::buffer::OffsetBuffer::from_lengths(attr_lengths), + entries, + None, + false, + )); + let batch = RecordBatch::try_from_iter(vec![ + ("k", Arc::new(Int32Array::from(keys)) as ArrayRef), + ("items", items), + ("attrs", attrs), + ("pad", Arc::new(StringArray::from(pad)) as ArrayRef), + ]) + .unwrap(); + (0..count).map(|i| batch.slice(i * rows, rows)).collect() + } + + #[tokio::test] + async fn sliced_nested_rows_match_the_reference() -> Result<()> { + let input = sliced_nested_batches(40, 50, 17); + let schema = input[0].schema(); + let ordering = LexOrdering::new(vec![PhysicalSortExpr::new( + col("k", &schema)?, + SortOptions::default(), + )]) + .unwrap(); + assert!(LateMaterialization::select(&input[0], &ordering)?.is_some()); + let all = concat_batches(&schema, &input)?; + let indices = arrow::compute::sort_to_indices(all.column(0), None, None)?; + let expected = take_record_batch(&all, &indices)?; + let bytes: usize = input.iter().map(|b| b.get_sliced_size().unwrap()).sum(); + for pool in [None, Some(PeakPool::new(bytes / 8))] { + let source = TestMemoryExec::try_new_exec( + std::slice::from_ref(&input), + Arc::clone(&schema), + None, + )?; + let sort = Arc::new(SortExec::new(ordering.clone(), source)); + let output = collect( + Arc::clone(&sort) as Arc, + context(pool.clone().map(|p| p as _), 32, 64 << 10), + ) + .await?; + assert_eq!(concat_batches(&schema, &output)?, expected); + if let Some(pool) = pool { + assert!(sort.metrics().unwrap().spill_count().unwrap() > 0); + assert_eq!(pool.reserved(), 0); + } + } + Ok(()) + } }