Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -382,46 +382,53 @@ fn sliced_bytes(array: &dyn Array) -> Result<usize> {
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::<i32>(0)[data.offset()..=data.offset() + data.len()];
for array in arrays {
let (offsets, values) = match array.as_list_opt::<i32>() {
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::<i64>(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<ArrayRef> = arrays
.iter()
.map(|array| {
let list = array.as_list::<i64>();
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<ArrayData> = arrays
DataType::FixedSizeList(_, _) => {
let children: Vec<ArrayRef> = 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<ArrayData> = arrays
let children: Vec<ArrayRef> = 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)
}),
Expand Down Expand Up @@ -645,11 +652,9 @@ impl Gather {
}

fn merge(&mut self, arrays: &[ArrayRef], bytes: usize) -> Option<ArrayRef> {
let data: Vec<ArrayData> = 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() {
Expand Down Expand Up @@ -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(())
}

Expand Down Expand Up @@ -1753,4 +1758,176 @@ mod tests {
}
Ok(())
}

fn list_of(values: ArrayRef, lengths: Vec<usize>) -> 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::<Int32Type, _, _>(
(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<RecordBatch> {
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<i32> = (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<dyn ExecutionPlan>,
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(())
}
}
Loading