Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
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 @@ -302,7 +302,6 @@ mod tests {
Arc::clone(&element_dtype),
Nullability::NonNullable,
0,
0,
ctx.allocator(),
);
list.clone()
Expand All @@ -314,7 +313,6 @@ mod tests {
Arc::clone(&element_dtype),
Nullability::NonNullable,
0,
0,
ctx.allocator(),
);
list.clone()
Expand All @@ -333,7 +331,6 @@ mod tests {
element_dtype,
Nullability::NonNullable,
0,
0,
vortex_buffer::BufferAllocatorRef::static_ref(),
);
listview
Expand Down
12 changes: 2 additions & 10 deletions encodings/sparse/src/canonical.rs
Original file line number Diff line number Diff line change
Expand Up @@ -402,34 +402,26 @@ fn execute_sparse_lists(
fill_value,
values_dtype,
len,
total_canonical_values,
nullability,
ctx,
)
})
}))
}

#[expect(clippy::too_many_arguments)]
fn execute_sparse_lists_inner<I: IntegerPType, O: OffsetBuilderPType>(
patch_indices: &[I],
patch_values: ListViewArray,
fill_value: &Scalar,
values_dtype: Arc<DType>,
len: usize,
total_canonical_values: usize,
nullability: Nullability,
ctx: &mut ExecutionCtx,
) -> ArrayRef {
// Create the builder with appropriate types. It is easy to just use the same type for both
// `offsets` and `sizes` since we have no other constraints.
let mut builder = ListViewBuilder::<O, O>::with_capacity_in(
values_dtype,
nullability,
total_canonical_values,
len,
ctx.allocator(),
);
let mut builder =
ListViewBuilder::<O, O>::with_capacity_in(values_dtype, nullability, len, ctx.allocator());
// The fill's elements become an array once, up front. Every gap then appends that same array,
// so the fill's elements are stored once for the whole result however many gaps reference them.
let fill_elements = list_scalar_elements_array(fill_value.as_list(), ctx.allocator());
Expand Down
4 changes: 4 additions & 0 deletions vortex-array/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,10 @@ harness = false
name = "chunked_fsl_canonicalize"
harness = false

[[bench]]
name = "chunked_nested_canonicalize"
harness = false

[[bench]]
name = "scalar_at_struct"
harness = false
Expand Down
76 changes: 76 additions & 0 deletions vortex-array/benches/chunked_nested_canonicalize.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

#![expect(clippy::unwrap_used)]

use std::sync::LazyLock;

use divan::Bencher;
use vortex_array::ArrayRef;
use vortex_array::Canonical;
use vortex_array::IntoArray;
use vortex_array::VortexSessionExecute;
use vortex_array::array_session;
use vortex_array::arrays::ChunkedArray;
use vortex_array::arrays::ListArray;
use vortex_array::arrays::ListViewArray;
use vortex_array::arrays::PrimitiveArray;
use vortex_array::arrays::StructArray;
use vortex_array::validity::Validity;
use vortex_buffer::Buffer;
use vortex_session::VortexSession;

static SESSION: LazyLock<VortexSession> = LazyLock::new(array_session);
const ROWS: usize = 1_024;
const CHUNKS: &[usize] = &[2, 32];

fn main() {
LazyLock::force(&SESSION);
divan::main();
}

fn canonicalize(bencher: Bencher, chunk: ArrayRef, nchunks: usize) {
let array = ChunkedArray::try_new(
std::iter::repeat_n(chunk.clone(), nchunks),
chunk.dtype().clone(),
)
.unwrap()
.into_array();
bencher
.with_inputs(|| (&array, SESSION.create_execution_ctx()))
.bench_refs(|(array, ctx)| array.clone().execute::<Canonical>(ctx).unwrap());
}

#[divan::bench(args = CHUNKS, consts = [1, 8])]
fn structs<const FIELDS: usize>(bencher: Bencher, nchunks: usize) {
let field = PrimitiveArray::from_iter(0..ROWS as u64).into_array();
let chunk =
StructArray::try_from_iter((0..FIELDS).map(|idx| (format!("field_{idx}"), field.clone())))
.unwrap();
canonicalize(bencher, chunk.into_array(), nchunks);
}

#[divan::bench(args = CHUNKS)]
fn lists(bencher: Bencher, nchunks: usize) {
let elements = PrimitiveArray::from_iter(0..4 * ROWS as u64).into_array();
let offsets: Buffer<u64> = (0..=ROWS).map(|idx| 4 * idx as u64).collect();
let chunk = ListArray::try_new(elements, offsets.into_array(), Validity::NonNullable).unwrap();
canonicalize(bencher, chunk.into_array(), nchunks);
}

#[divan::bench(args = CHUNKS, consts = [false, true])]
fn listviews<const OVERLAPPING: bool>(bencher: Bencher, nchunks: usize) {
let step = if OVERLAPPING { 1 } else { 4 };
let elements = PrimitiveArray::from_iter(0..(step * ROWS + 4) as u64).into_array();
let offsets: Buffer<u64> = (0..ROWS).map(|idx| (step * idx) as u64).collect();
let sizes: Buffer<u64> = std::iter::repeat_n(4, ROWS).collect();
let chunk = ListViewArray::new(
elements,
offsets.into_array(),
sizes.into_array(),
Validity::NonNullable,
);
// SAFETY: with step 4, every view ends exactly where the next begins.
let chunk = unsafe { chunk.with_zero_copy_to_list(!OVERLAPPING) };
canonicalize(bencher, chunk.into_array(), nchunks);
}
2 changes: 0 additions & 2 deletions vortex-array/benches/listview_builder_extend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,6 @@ fn extend_from_array_zctl(bencher: Bencher, (num_lists, list_size): (usize, usiz
let mut builder = ListViewBuilder::<u64, u64>::with_capacity_in(
Arc::new(DType::Primitive(I32, NonNullable)),
NonNullable,
num_lists * list_size,
num_lists,
ctx.allocator(),
);
Expand All @@ -99,7 +98,6 @@ fn extend_from_array_non_zctl_overlapping(
let mut builder = ListViewBuilder::<u64, u64>::with_capacity_in(
Arc::new(DType::Primitive(I32, NonNullable)),
Nullable,
num_lists * list_size,
num_lists,
ctx.allocator(),
);
Expand Down
7 changes: 0 additions & 7 deletions vortex-array/src/array/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -121,9 +121,6 @@ pub(crate) trait DynArrayData: 'static + private::Sealed + Send + Sync + Debug {
/// Returns the array as a reference to a generic [`Any`] trait object.
fn as_any(&self) -> &dyn Any;

/// Returns the array as a mutable reference to a generic [`Any`] trait object.
fn as_any_mut(&mut self) -> &mut dyn Any;

/// Returns the [`Validity`] of the array.
fn validity(&self, this: &ArrayRef) -> VortexResult<Validity>;

Expand Down Expand Up @@ -271,10 +268,6 @@ impl<V: VTable> DynArrayData for ArrayData<V> {
self
}

fn as_any_mut(&mut self) -> &mut dyn Any {
self
}

fn validity(&self, this: &ArrayRef) -> VortexResult<Validity> {
if this.dtype().is_nullable() {
let view = unsafe { ArrayView::new_unchecked(this, &self.data) };
Expand Down
25 changes: 23 additions & 2 deletions vortex-array/src/array/typed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ use crate::validity::Validity;
/// `ArrayRef` stores `Arc<ArrayInner<dyn DynArrayData>>` — a single 16-byte fat pointer.
/// Metadata is accessed via `self.0.*` (a normal struct field read through the Arc),
/// while encoding-specific methods go through `self.0.data` (vtable dispatch).
#[derive(Clone)]
pub(crate) struct ArrayInner<D: ?Sized> {
pub(crate) len: usize,
pub(crate) encoding_id: ArrayId,
Expand Down Expand Up @@ -306,8 +307,28 @@ impl<V: VTable> Array<V> {
/// Returns `None` when this handle is not the unique owner of the backing allocation.
pub fn data_mut(&mut self) -> Option<&mut V::TypedArrayData> {
let store = self.inner.inner_mut()?;
let array_inner = store.data.as_any_mut().downcast_mut::<ArrayData<V>>();
Some(&mut array_inner?.data)
// NOTE(ngates): use downcast_mut_unchecked when it becomes stable
debug_assert!(store.data.as_any().is::<ArrayData<V>>());
// SAFETY: `Array<V>` guarantees the inner is `ArrayData<V>`, the same invariant the shared
// path in `downcast_inner` relies on. Reaching it through `Any` instead spends a virtual
// `as_any_mut` and a `TypeId` comparison on a type the compiler already knows, which the
// executor pays once per chunk through `with_next_child_slot`.
let array_inner =
unsafe { &mut *std::ptr::from_mut(&mut store.data).cast::<ArrayData<V>>() };
Some(&mut array_inner.data)
}

/// Mutates encoding data, copying shared array metadata and slots while retaining statistics.
///
/// The update must preserve the logical values and array invariants.
pub(crate) fn with_data_mut(self, update: impl FnOnce(&mut V::TypedArrayData)) -> Self {
// SAFETY: Array<V> guarantees the inner is ArrayData<V>.
let mut inner = unsafe { self.inner.downcast_inner_unchecked::<V>() };
update(&mut Arc::make_mut(&mut inner).data.data);
Self {
inner: ArrayRef::from_inner(inner),
_phantom: PhantomData,
}
}

/// Returns the full typed array construction parts if this handle owns the allocation.
Expand Down
1 change: 0 additions & 1 deletion vortex-array/src/arrays/arbitrary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -315,7 +315,6 @@ fn random_list_with_offset_type<O: OffsetBuilderPType>(
let mut builder = ListViewBuilder::<O, O>::with_capacity_in(
Arc::clone(elem_dtype),
null,
array_length,
10,
vortex_buffer::BufferAllocatorRef::static_ref(),
);
Expand Down
2 changes: 1 addition & 1 deletion vortex-array/src/arrays/bool/vtable/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@ impl VTable for Bool {
let Some(builder) = builder.as_any_mut().downcast_mut::<BoolBuilder>() else {
vortex_bail!("append_to_builder for Bool requires a BoolBuilder");
};
builder.append_bool_array(&array.into_owned(), ctx)
builder.append_bool_array(array, ctx)
}

fn execute(array: Array<Self>, _ctx: &mut ExecutionCtx) -> VortexResult<ExecutionResult> {
Expand Down
57 changes: 20 additions & 37 deletions vortex-array/src/arrays/chunked/array.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use std::fmt::Display;
use std::fmt::Formatter;

use futures::stream;
use vortex_buffer::BufferMut;
use vortex_buffer::Buffer;
use vortex_error::VortexExpect;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
Expand Down Expand Up @@ -49,8 +49,8 @@ pub struct ChunkedSlots {
#[derive(Clone, Debug)]
pub struct ChunkedData {
pub(super) chunk_offsets: Vec<usize>,
/// This is used to find the next child to execute when in executing into a builder.
pub(super) next_builder_slot: usize,
/// The next child to visit while executing slots or appending into a builder.
pub(super) next_child_slot: usize,
}

impl Display for ChunkedData {
Expand Down Expand Up @@ -79,20 +79,18 @@ pub trait ChunkedArrayExt: TypedArrayRef<Chunked> {
.vortex_expect("validated chunk slot")
}

fn iter_chunks<'a>(&'a self) -> Box<dyn Iterator<Item = &'a ArrayRef> + 'a> {
Box::new(
self.as_ref().slots()[ChunkedSlots::CHUNKS_OFFSET..]
.iter()
.map(|slot| slot.as_ref().vortex_expect("validated chunk slot")),
)
fn iter_chunks(&self) -> impl Iterator<Item = &ArrayRef> {
self.as_ref().slots()[ChunkedSlots::CHUNKS_OFFSET..]
.iter()
.map(|slot| slot.as_ref().vortex_expect("validated chunk slot"))
}

fn chunks(&self) -> Vec<ArrayRef> {
self.iter_chunks().cloned().collect()
}

fn non_empty_chunks<'a>(&'a self) -> Box<dyn Iterator<Item = &'a ArrayRef> + 'a> {
Box::new(self.iter_chunks().filter(|chunk| !chunk.is_empty()))
fn non_empty_chunks(&self) -> impl Iterator<Item = &ArrayRef> {
self.iter_chunks().filter(|chunk| !chunk.is_empty())
}

/// Returns the cached chunk boundary offsets.
Expand Down Expand Up @@ -135,18 +133,19 @@ impl ChunkedData {
pub(super) fn new(chunk_offsets: Vec<usize>) -> Self {
Self {
chunk_offsets,
next_builder_slot: ChunkedSlots::CHUNKS_OFFSET,
next_child_slot: ChunkedSlots::CHUNKS_OFFSET,
}
}

pub(super) fn make_chunk_offsets_array(chunk_offsets: &[usize]) -> ArrayRef {
let mut chunk_offsets_buf = BufferMut::<u64>::with_capacity(chunk_offsets.len());
for &offset in chunk_offsets {
let offset = u64::try_from(offset)
.vortex_expect("chunk offset must fit in u64 for serialization");
unsafe { chunk_offsets_buf.push_unchecked(offset) }
}
PrimitiveArray::new(chunk_offsets_buf.freeze(), Validity::NonNullable).into_array()
let chunk_offsets_buf =
Buffer::from_trusted_len_iter(chunk_offsets.iter().copied().map(|offset| {
u64::try_from(offset)
.vortex_expect("chunk offset must fit in u64 for serialization")
}));

unsafe { PrimitiveArray::new_unchecked(chunk_offsets_buf, Validity::NonNullable) }
.into_array()
}

/// Validates the components that would be used to create a `ChunkedArray`.
Expand Down Expand Up @@ -190,24 +189,8 @@ impl Array<Chunked> {
Ok(ArrayParts::new(Chunked, dtype, len, ChunkedData::new(chunk_offsets)).with_slots(slots))
}

pub(super) fn with_next_builder_slot(mut self, next_builder_slot: usize) -> Self {
if let Some(data) = self.data_mut() {
data.next_builder_slot = next_builder_slot;
return self;
}
// This is the slow path that will be hit at most once per execution since the second one
// *MUST* have execlusive access due to this copy.
let stats = self.statistics().to_owned();
let mut data = self.data().clone();
data.next_builder_slot = next_builder_slot;
// SAFETY: we only modified next_builder_slot which doesn't affect array invariants.
unsafe {
Array::from_parts_unchecked(
ArrayParts::new(Chunked, self.dtype().clone(), self.len(), data)
.with_slots(self.slots().iter().cloned().collect::<ArraySlots>()),
)
}
.with_stats_set(stats)
pub(super) fn with_next_child_slot(self, next_child_slot: usize) -> Self {
self.with_data_mut(|data| data.next_child_slot = next_child_slot)
}

/// Constructs a new `ChunkedArray`.
Expand Down
Loading
Loading