From f1008f3f76f53b09ad73b30dfaf126ae74f2a4a2 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Tue, 15 Sep 2026 21:11:49 +0100 Subject: [PATCH 1/7] perf(array): reuse probe state in RLE, RunEnd and PCO MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Port of #9844 onto the reduced probe API. Encodings implement one `probe_scalar(state, index, ctx)` and read children through `state.slot(i)`; the same body serves one-off and repeated reads. - FastLanes RLE: `RleProbeState` keeps the slice's base value offset. `state.split()` holds it while the indices, offsets and values children are read through `ProbeChildren::slot`. - RunEnd: the ends child probe is held across the binary search, so a repeated read keeps the child's preparation within and between searches; the values child is probed for the selected run. - PCO: `PcoProbeState` keeps the validity mask, prefix ranks for non-null positions, page boundaries and the most recently decoded page. A one-off read decompresses the single row as before. - Tests cover one-off and repeated access over sliced and nullable inputs, nested RunEnd(PCO, RunEnd(PCO, PCO)) state reuse and drop counts, lazy validity, and page-cache eviction. The PCO crate gains a scalar-probe benchmark and an example. Short local run of the new bench, medians, 1024 clustered lookups: PCO 5.54 ms via `execute_scalar` vs 24 µs via a repeated probe; RunEnd(PCO, PCO) 24.5 ms vs 4.1 ms; RLE nullable 176 µs vs 133 µs. Scattered PCO lookups are unchanged, each landing on a new page. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_016CqrLKgPqFYGZK5sjk1qe7 --- Cargo.lock | 3 + encodings/fastlanes/src/rle/mod.rs | 1 + encodings/fastlanes/src/rle/probe.rs | 86 +++++ encodings/fastlanes/src/rle/probe/tests.rs | 94 ++++++ .../fastlanes/src/rle/vtable/operations.rs | 38 +-- encodings/pco/Cargo.toml | 6 + encodings/pco/benches/probe.rs | 198 +++++++++++ encodings/pco/examples/probe.rs | 42 +++ encodings/pco/src/array.rs | 19 +- encodings/pco/src/lib.rs | 1 + encodings/pco/src/probe.rs | 159 +++++++++ encodings/pco/src/probe/tests.rs | 311 ++++++++++++++++++ encodings/runend/src/lib.rs | 1 + encodings/runend/src/ops.rs | 14 +- encodings/runend/src/probe.rs | 50 +++ encodings/runend/src/probe/tests.rs | 102 ++++++ 16 files changed, 1090 insertions(+), 35 deletions(-) create mode 100644 encodings/fastlanes/src/rle/probe.rs create mode 100644 encodings/fastlanes/src/rle/probe/tests.rs create mode 100644 encodings/pco/benches/probe.rs create mode 100644 encodings/pco/examples/probe.rs create mode 100644 encodings/pco/src/probe.rs create mode 100644 encodings/pco/src/probe/tests.rs create mode 100644 encodings/runend/src/probe.rs create mode 100644 encodings/runend/src/probe/tests.rs diff --git a/Cargo.lock b/Cargo.lock index a2dd783356b..88f8586fa61 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11411,6 +11411,7 @@ dependencies = [ name = "vortex-pco" version = "0.1.0" dependencies = [ + "codspeed-divan-compat", "pco", "prost 0.14.4", "rstest", @@ -11418,7 +11419,9 @@ dependencies = [ "vortex-arrow", "vortex-buffer", "vortex-error", + "vortex-fastlanes", "vortex-mask", + "vortex-runend", "vortex-session", ] diff --git a/encodings/fastlanes/src/rle/mod.rs b/encodings/fastlanes/src/rle/mod.rs index 25610742c19..57b37c7201c 100644 --- a/encodings/fastlanes/src/rle/mod.rs +++ b/encodings/fastlanes/src/rle/mod.rs @@ -9,6 +9,7 @@ pub use array::RLESlots; mod compute; mod kernel; +mod probe; mod vtable; pub use vtable::RLE; diff --git a/encodings/fastlanes/src/rle/probe.rs b/encodings/fastlanes/src/rle/probe.rs new file mode 100644 index 00000000000..8f80aa09228 --- /dev/null +++ b/encodings/fastlanes/src/rle/probe.rs @@ -0,0 +1,86 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! RLE scalar access through the indices, offsets, and values slots. + +use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; +use vortex_array::scalar::Scalar; +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::FL_CHUNK_SIZE; +use crate::RLE; +use crate::rle::RLEArrayExt; +use crate::rle::RLESlots; + +/// The slice's base value offset, retained across repeated lookups. +#[derive(Default)] +pub struct RleProbeState { + base: Option, +} + +pub(crate) fn scalar_at( + state: &mut ProbeState<'_, RLE>, + index: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let array = state.array(); + let logical_index = array.offset() + index; + // Hold the retained state while reading children: the base offset is cached across rows. + let (mut retained, mut children) = state.split(); + let code = children + .slot(RLESlots::INDICES)? + .ok_or_else(|| vortex_err!("RLE indices slot is missing"))? + .execute_scalar(logical_index, ctx)?; + let Some(code) = code.as_primitive().as_::() else { + return Ok(Scalar::null(array.dtype().clone())); + }; + + let chunk = logical_index / FL_CHUNK_SIZE; + let offset = if chunk == 0 { + 0 + } else { + let base = match retained.as_deref().and_then(|state| state.base) { + Some(base) => base, + None => { + let value = read_offset( + children + .slot(RLESlots::VALUES_IDX_OFFSETS)? + .ok_or_else(|| vortex_err!("RLE offsets slot is missing"))? + .execute_scalar(0, ctx)?, + )?; + if let Some(state) = retained.as_mut() { + state.base = Some(value); + } + value + } + }; + read_offset( + children + .slot(RLESlots::VALUES_IDX_OFFSETS)? + .ok_or_else(|| vortex_err!("RLE offsets slot is missing"))? + .execute_scalar(chunk, ctx)?, + )? + .checked_sub(base) + .ok_or_else(|| vortex_err!("RLE offsets precede the slice base"))? + }; + let value_index = offset + .checked_add(code) + .ok_or_else(|| vortex_err!("RLE value index overflow"))?; + let scalar = children + .slot(RLESlots::VALUES)? + .ok_or_else(|| vortex_err!("RLE values slot is missing"))? + .execute_scalar(value_index, ctx)?; + Scalar::try_new(array.dtype().clone(), scalar.into_value()) +} + +fn read_offset(scalar: Scalar) -> VortexResult { + scalar + .as_primitive() + .as_::() + .ok_or_else(|| vortex_err!("RLE offset must be a non-null usize")) +} + +#[cfg(test)] +mod tests; diff --git a/encodings/fastlanes/src/rle/probe/tests.rs b/encodings/fastlanes/src/rle/probe/tests.rs new file mode 100644 index 00000000000..fec027edf58 --- /dev/null +++ b/encodings/fastlanes/src/rle/probe/tests.rs @@ -0,0 +1,94 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use rstest::rstest; +use vortex_array::ArrayProbe; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::RepeatedArrayProbe; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::assert_arrays_eq; +use vortex_array::builders::builder_with_capacity_in; +use vortex_array::builtins::ArrayBuiltins; +use vortex_array::scalar_fn::fns::operators::Operator; +use vortex_array::validity::Validity; +use vortex_error::VortexResult; + +use crate::RLE; +use crate::RLEData; + +/// A one-off or retained probe over `array`, the retained one living in `retained`. +fn probe_for<'a>( + array: &'a ArrayRef, + retained: &'a mut Option, + repeated: bool, +) -> ArrayProbe<'a> { + if repeated { + retained.insert(array.repeated_probe()).as_probe() + } else { + array.probe() + } +} + +#[rstest] +fn random_access_across_chunks_and_nulls( + #[values(false, true)] repeated: bool, + #[values(false, true)] sliced: bool, +) -> VortexResult<()> { + let mut ctx = crate::test::SESSION.create_execution_ctx(); + let input = + PrimitiveArray::from_option_iter((0..8192u32).map(|i| (i % 11 != 0).then_some(i / 16))); + let encoded = RLEData::encode(input.as_view(), &mut ctx)?.into_array(); + let range = if sliced { 1777..7333 } else { 0..8192 }; + let source = if sliced { + encoded + .slice(range.clone())? + .execute::(&mut ctx)? + } else { + encoded + }; + assert!(source.is::()); + let input = input.slice(range)?; + let indices = [0u32, 1, 1023, 1024, 2048, 2047, 4097, 11, 33, 17, 0]; + let mut actual = builder_with_capacity_in(source.dtype(), indices.len(), ctx.allocator()); + let mut retained = None; + let mut probe = probe_for(&source, &mut retained, repeated); + for index in indices { + actual.append_scalar(&probe.execute_scalar(index as usize, &mut ctx)?)?; + } + assert_arrays_eq!( + actual.finish(), + input.take(PrimitiveArray::from_iter(indices).into_array())?, + &mut ctx + ); + assert!(probe.execute_scalar(source.len(), &mut ctx).is_err()); + Ok(()) +} + +#[rstest] +fn lazy_validity_does_not_evaluate_unrequested_rows( + #[values(false, true)] repeated: bool, +) -> VortexResult<()> { + let mut ctx = crate::test::SESSION.create_execution_ctx(); + let numerators = PrimitiveArray::from_iter(vec![1u32; 1024]).into_array(); + let denominators = + PrimitiveArray::from_iter((0..1024).map(|i| u32::from(i != 1023))).into_array(); + let validity = numerators + .binary(denominators, Operator::Div)? + .binary(numerators, Operator::Eq)?; + let array = RLE::try_new( + PrimitiveArray::from_iter([42u32]).into_array(), + PrimitiveArray::new(vec![0u16; 1024], Validity::Array(validity)).into_array(), + PrimitiveArray::from_iter([0u64]).into_array(), + 0, + 1024, + )? + .into_array(); + let expected = array.execute_scalar(0, &mut ctx)?; + let mut retained = None; + let mut probe = probe_for(&array, &mut retained, repeated); + assert_eq!(probe.execute_scalar(0, &mut ctx)?, expected); + assert!(probe.execute_scalar(1023, &mut ctx).is_err()); + Ok(()) +} diff --git a/encodings/fastlanes/src/rle/vtable/operations.rs b/encodings/fastlanes/src/rle/vtable/operations.rs index ba6a624f1f5..92060060d23 100644 --- a/encodings/fastlanes/src/rle/vtable/operations.rs +++ b/encodings/fastlanes/src/rle/vtable/operations.rs @@ -3,42 +3,32 @@ use vortex_array::ArrayView; use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; use vortex_array::scalar::Scalar; use vortex_array::vtable::OperationsVTable; -use vortex_error::VortexExpect; use vortex_error::VortexResult; use super::RLE; -use crate::FL_CHUNK_SIZE; -use crate::rle::RLEArrayExt; -use crate::rle::RLEArraySlotsExt; +use crate::rle::probe; +use crate::rle::probe::RleProbeState; impl OperationsVTable for RLE { - type ProbeState = (); + type ProbeState = RleProbeState; + + fn probe_scalar( + state: &mut ProbeState<'_, RLE>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + probe::scalar_at(state, index, ctx) + } fn scalar_at( array: ArrayView<'_, RLE>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let offset_in_chunk = array.offset(); - let chunk_relative_idx = array - .indices() - .execute_scalar(offset_in_chunk + index, ctx)?; - - let chunk_relative_idx = chunk_relative_idx - .as_primitive() - .as_::() - .vortex_expect("Index must not be null"); - - let chunk_id = (offset_in_chunk + index) / FL_CHUNK_SIZE; - let value_idx_offset = array.values_idx_offset(chunk_id, ctx); - - let scalar = array - .values() - .execute_scalar(value_idx_offset + chunk_relative_idx, ctx)?; - - Scalar::try_new(array.dtype().clone(), scalar.into_value()) + probe::scalar_at(&mut ProbeState::once(array), index, ctx) } } @@ -54,9 +44,9 @@ mod tests { use vortex_array::validity::Validity; use vortex_buffer::Buffer; use vortex_buffer::buffer; + use vortex_error::VortexExpect; use vortex_session::VortexSession; - use super::*; use crate::RLE; use crate::RLEArray; use crate::RLEData; diff --git a/encodings/pco/Cargo.toml b/encodings/pco/Cargo.toml index 1683d4875e9..c3a1e0019ef 100644 --- a/encodings/pco/Cargo.toml +++ b/encodings/pco/Cargo.toml @@ -26,7 +26,13 @@ vortex-mask = { workspace = true } vortex-session = { workspace = true } [dev-dependencies] +divan = { workspace = true } +vortex-fastlanes = { workspace = true } +vortex-runend = { workspace = true } rstest = { workspace = true } vortex-array = { workspace = true, features = ["_test-harness"] } vortex-arrow = { workspace = true } +[[bench]] +name = "probe" +harness = false diff --git a/encodings/pco/benches/probe.rs b/encodings/pco/benches/probe.rs new file mode 100644 index 00000000000..4ecc47074cb --- /dev/null +++ b/encodings/pco/benches/probe.rs @@ -0,0 +1,198 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Scalar access comparisons, including preparation and teardown for every group of lookups. + +use std::hint::black_box; +use std::sync::LazyLock; + +use divan::Bencher; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::validity::Validity; +use vortex_error::VortexExpect; +use vortex_fastlanes::RLEData; +use vortex_pco::Pco; +use vortex_runend::RunEnd; +use vortex_session::VortexSession; + +fn main() { + divan::main(); +} + +const LEN: usize = 16_384; +// Number of accesses, nullable, scattered (otherwise clustered within 256 rows). +const CASES: &[(usize, bool, bool)] = &[ + (1, false, false), + (1, true, false), + (64, false, false), + (64, true, false), + (64, false, true), + (64, true, true), + (1024, false, false), + (1024, true, false), + (1024, false, true), + (1024, true, true), +]; + +static SESSION: LazyLock = LazyLock::new(|| { + let session = vortex_array::array_session(); + vortex_fastlanes::initialize(&session); + vortex_runend::initialize(&session); + session +}); + +fn input(nullable: bool) -> PrimitiveArray { + let validity = if nullable { + Validity::from_iter((0..LEN).map(|i| i % 11 != 0)) + } else { + Validity::NonNullable + }; + PrimitiveArray::new( + (0..LEN) + .map(|i| u32::try_from(i / 16).vortex_expect("fixture values fit u32")) + .collect::>(), + validity, + ) +} + +fn indices(count: usize, scattered: bool) -> Vec { + let span = if scattered { LEN } else { 256 }; + let base = if scattered { 0 } else { 4096 }; + let mut seed = 42u64; + (0..count) + .map(|_| { + seed = seed.wrapping_mul(6364136223846793005).wrapping_add(1); + base + ((seed >> 32) as usize % span) + }) + .collect() +} + +fn rle(nullable: bool) -> ArrayRef { + let mut ctx = SESSION.create_execution_ctx(); + RLEData::encode(input(nullable).as_view(), &mut ctx) + .vortex_expect("RLE compression") + .into_array() +} + +fn pco(nullable: bool) -> ArrayRef { + let mut ctx = SESSION.create_execution_ctx(); + Pco::from_primitive(input(nullable).as_view(), 8, 1024, &mut ctx) + .vortex_expect("PCO compression") + .into_array() +} + +fn runend_pco(nullable: bool) -> ArrayRef { + let mut ctx = SESSION.create_execution_ctx(); + let runs = u32::try_from(LEN / 4).vortex_expect("run count fits u32"); + let ends = PrimitiveArray::from_iter((1..=runs).map(|run| run * 4)); + let values = PrimitiveArray::new( + (0..runs).collect::>(), + if nullable { + Validity::from_iter((0..runs).map(|run| run % 11 != 0)) + } else { + Validity::NonNullable + }, + ); + RunEnd::try_new( + Pco::from_primitive(ends.as_view(), 8, 1024, &mut ctx) + .vortex_expect("PCO ends compression") + .into_array(), + Pco::from_primitive(values.as_view(), 8, 1024, &mut ctx) + .vortex_expect("PCO values compression") + .into_array(), + &mut ctx, + ) + .vortex_expect("RunEnd construction") + .into_array() +} + +fn execute_scalar(bencher: Bencher, array: ArrayRef, indices: &[usize]) { + bencher + .with_inputs(|| SESSION.create_execution_ctx()) + .bench_refs(|ctx| { + for &index in indices { + black_box( + array + .execute_scalar(black_box(index), ctx) + .vortex_expect("scalar access"), + ); + } + }); +} + +fn probe(bencher: Bencher, array: ArrayRef, indices: &[usize], repeated: bool) { + bencher + .with_inputs(|| SESSION.create_execution_ctx()) + .bench_refs(|ctx| { + let mut retained = repeated.then(|| array.repeated_probe()); + let mut probe = match &mut retained { + Some(retained) => retained.as_probe(), + None => array.probe(), + }; + for &index in indices { + black_box( + probe + .execute_scalar(black_box(index), ctx) + .vortex_expect("probe access"), + ); + } + }); +} + +#[divan::bench(args = CASES)] +fn rle_probe(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { + probe( + bencher, + rle(nullable), + &indices(count, scattered), + count != 1, + ); +} + +#[divan::bench(args = CASES)] +fn pco_probe(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { + probe( + bencher, + pco(nullable), + &indices(count, scattered), + count != 1, + ); +} + +#[divan::bench(args = [false, true])] +fn rle_repeated_first(bencher: Bencher, nullable: bool) { + probe(bencher, rle(nullable), &indices(1, false), true); +} + +#[divan::bench(args = [false, true])] +fn pco_repeated_first(bencher: Bencher, nullable: bool) { + probe(bencher, pco(nullable), &indices(1, false), true); +} + +#[divan::bench(args = CASES)] +fn rle_execute_scalar(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { + execute_scalar(bencher, rle(nullable), &indices(count, scattered)); +} + +#[divan::bench(args = CASES)] +fn pco_execute_scalar(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { + execute_scalar(bencher, pco(nullable), &indices(count, scattered)); +} + +#[divan::bench(args = CASES)] +fn runend_pco_execute_scalar(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { + execute_scalar(bencher, runend_pco(nullable), &indices(count, scattered)); +} + +#[divan::bench(args = CASES)] +fn runend_pco_probe(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { + probe( + bencher, + runend_pco(nullable), + &indices(count, scattered), + count != 1, + ); +} diff --git a/encodings/pco/examples/probe.rs b/encodings/pco/examples/probe.rs new file mode 100644 index 00000000000..a8183043216 --- /dev/null +++ b/encodings/pco/examples/probe.rs @@ -0,0 +1,42 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Run with `cargo run -p vortex-pco --example probe`. + +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_error::VortexResult; +use vortex_fastlanes::RLEData; +use vortex_pco::Pco; +use vortex_runend::RunEnd; + +fn main() -> VortexResult<()> { + let session = vortex_array::array_session(); + vortex_fastlanes::initialize(&session); + vortex_runend::initialize(&session); + let mut ctx = session.create_execution_ctx(); + let values = + PrimitiveArray::from_option_iter((0..4096u32).map(|i| (i % 11 != 0).then_some(i / 16))); + let rle = RLEData::encode(values.as_view(), &mut ctx)?.into_array(); + let pco = Pco::from_primitive(values.as_view(), 3, 1024, &mut ctx)?.into_array(); + let ends = PrimitiveArray::from_iter((1..=4096u32).map(|run| run * 4)); + let ends = Pco::from_primitive(ends.as_view(), 3, 1024, &mut ctx)?.into_array(); + let runend_pco = RunEnd::try_new(ends, pco.clone(), &mut ctx)?.into_array(); + + for (name, array) in [("RLE", rle), ("PCO", pco), ("RunEnd(PCO, PCO)", runend_pco)] { + // A one-off probe borrows the array and retains nothing. + let scalar = array.probe().execute_scalar(17, &mut ctx)?; + println!("{name} single lookup: {scalar}"); + + // A repeated probe keeps child probes; each PCO child keeps its own decoded page. + let mut probe = array.repeated_probe(); + for index in [17, 33, 22, 1025, 17] { + println!( + "{name}[{index}] = {}", + probe.execute_scalar(index, &mut ctx)? + ); + } + } + Ok(()) +} diff --git a/encodings/pco/src/array.rs b/encodings/pco/src/array.rs index 97a17fe6a3c..2d9d563f473 100644 --- a/encodings/pco/src/array.rs +++ b/encodings/pco/src/array.rs @@ -29,6 +29,7 @@ use vortex_array::EqMode; use vortex_array::ExecutionCtx; use vortex_array::ExecutionResult; use vortex_array::IntoArray; +use vortex_array::ProbeState; use vortex_array::TypedArrayRef; use vortex_array::array_slots; use vortex_array::arrays::Primitive; @@ -778,19 +779,23 @@ impl ValidityVTable for Pco { } impl OperationsVTable for Pco { - type ProbeState = (); + type ProbeState = crate::probe::PcoProbeState; + + fn probe_scalar( + state: &mut ProbeState<'_, Pco>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + let array = state.array(); + crate::probe::scalar_at(array, index, state.retained(), ctx) + } fn scalar_at( array: ArrayView<'_, Pco>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let unsliced_validity = array.unsliced_validity(); - array - ._slice(index, index + 1) - .decompress(&unsliced_validity, ctx)? - .into_array() - .execute_scalar(0, ctx) + crate::probe::scalar_at(array, index, None, ctx) } } diff --git a/encodings/pco/src/lib.rs b/encodings/pco/src/lib.rs index 601bd0826bd..70d71abf7d2 100644 --- a/encodings/pco/src/lib.rs +++ b/encodings/pco/src/lib.rs @@ -20,6 +20,7 @@ mod array; mod compute; +mod probe; mod rules; mod slice; diff --git a/encodings/pco/src/probe.rs b/encodings/pco/src/probe.rs new file mode 100644 index 00000000000..771172200a3 --- /dev/null +++ b/encodings/pco/src/probe.rs @@ -0,0 +1,159 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! PCO probes retain one decoded page and index the compacted, non-null value positions. + +use std::ops::Range; + +use pco::data_types::Number; +use pco::data_types::NumberType; +use pco::match_number_enum; +use pco::wrapped::FileDecompressor; +use vortex_array::ArrayView; +use vortex_array::ExecutionCtx; +use vortex_array::IntoArray; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::dtype::NativePType; +use vortex_array::dtype::half; +use vortex_array::match_each_native_ptype; +use vortex_array::scalar::Scalar; +use vortex_array::validity::Validity; +use vortex_buffer::BufferMut; +use vortex_error::VortexResult; +use vortex_error::vortex_err; +use vortex_mask::Mask; + +use crate::Pco; +use crate::PcoArrayExt; +use crate::array::number_type_from_ptype; +use crate::array::vortex_err_from_pco; + +const RANK_STRIDE: usize = 512; + +/// State retained by repeated PCO probes. Construction performs no allocation. +/// +/// The mask and rank index use unsliced logical rows; page boundaries use compacted value +/// positions. Only the most recently accessed page is retained, bounding decoded storage. +#[derive(Default)] +pub struct PcoProbeState { + validity: Option, + rank: Vec, + pages: Vec, + decoded: Option<(Range, PrimitiveArray)>, + #[cfg(test)] + decoded_pages: usize, + #[cfg(test)] + tracking: tests::TrackedProbeState, +} + +struct Page { + values: Range, + chunk: usize, +} + +pub(crate) fn scalar_at( + array: ArrayView<'_, Pco>, + index: usize, + state: Option<&mut PcoProbeState>, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let Some(state) = state else { + let validity = array.unsliced_validity(); + return array + ._slice(index, index + 1) + .decompress(&validity, ctx)? + .into_array() + .execute_scalar(0, ctx); + }; + let mask = match &mut state.validity { + Some(mask) => mask, + slot @ None => slot.insert( + array + .unsliced_validity() + .execute_mask(array.unsliced_n_rows(), ctx)?, + ), + }; + let logical_index = array.slice_start() + index; + + let value_index = match mask { + Mask::AllTrue(_) => logical_index, + Mask::AllFalse(_) => unreachable!("probe dispatch checks validity"), + Mask::Values(values) => { + let bits = values.bit_buffer(); + if state.rank.is_empty() { + state.rank.reserve(bits.len().div_ceil(RANK_STRIDE)); + let mut count = 0; + for start in (0..bits.len()).step_by(RANK_STRIDE) { + state.rank.push(count); + count += bits.count_range(start, (start + RANK_STRIDE).min(bits.len())); + } + } + let block = logical_index / RANK_STRIDE; + state.rank[block] + bits.count_range(block * RANK_STRIDE, logical_index) + } + }; + + let (range, values) = match &state.decoded { + Some(decoded) if decoded.0.contains(&value_index) => decoded, + _ => { + if state.pages.is_empty() { + let mut start = 0; + for (chunk, metadata) in array.metadata.chunks.iter().enumerate() { + for page in &metadata.pages { + let end = start + page.n_values as usize; + state.pages.push(Page { + values: start..end, + chunk, + }); + start = end; + } + } + } + let page_index = state + .pages + .partition_point(|page| page.values.end <= value_index); + let page = state + .pages + .get(page_index) + .ok_or_else(|| vortex_err!("Missing PCO page for value {value_index}"))?; + let decoded = match_number_enum!( + number_type_from_ptype(array.dtype().as_ptype()), + NumberType => { decode_page::(array, page, array.pages[page_index].as_slice(), ctx)? } + ); + #[cfg(test)] + { + state.decoded_pages += 1; + state.tracking.record_decode(); + } + state.decoded.insert((page.values.clone(), decoded)) + } + }; + Ok(match_each_native_ptype!(values.ptype(), |T| { + Scalar::primitive( + values.as_slice::()[value_index - range.start], + array.dtype().nullability(), + ) + })) +} + +fn decode_page( + array: ArrayView<'_, Pco>, + page: &Page, + buffer: &[u8], + ctx: &mut ExecutionCtx, +) -> VortexResult { + let (file, _) = + FileDecompressor::new(array.metadata.header.as_slice()).map_err(vortex_err_from_pco)?; + let (mut chunk, _) = file + .chunk_decompressor::(array.chunk_metas[page.chunk].as_ref()) + .map_err(vortex_err_from_pco)?; + let mut decoder = chunk + .page_decompressor(buffer, page.values.len()) + .map_err(vortex_err_from_pco)?; + let mut values = BufferMut::::zeroed_in(page.values.len(), ctx.allocator().clone()); + decoder.read(&mut values).map_err(vortex_err_from_pco)?; + Ok(PrimitiveArray::new(values.freeze(), Validity::NonNullable)) +} + +#[cfg(test)] +mod tests; diff --git a/encodings/pco/src/probe/tests.rs b/encodings/pco/src/probe/tests.rs new file mode 100644 index 00000000000..ee42d2db596 --- /dev/null +++ b/encodings/pco/src/probe/tests.rs @@ -0,0 +1,311 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::cell::Cell; + +use rstest::rstest; +use vortex_array::ArrayProbe; +use vortex_array::ArrayRef; +use vortex_array::ExecutionCtx; +use vortex_array::IntoArray; +use vortex_array::RepeatedArrayProbe; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::assert_arrays_eq; +use vortex_array::builders::builder_with_capacity_in; +use vortex_array::validity::Validity; +use vortex_error::VortexResult; +use vortex_runend::RunEnd; + +use super::PcoProbeState; +use crate::Pco; + +#[derive(Clone, Copy, Debug, Default, PartialEq)] +struct ProbeCounts { + initialized: usize, + decoded: usize, + dropped: usize, +} + +thread_local! { + static COUNTS: Cell = Cell::default(); +} + +pub(super) struct TrackedProbeState; + +impl Default for TrackedProbeState { + fn default() -> Self { + let counts = COUNTS.get(); + COUNTS.set(ProbeCounts { + initialized: counts.initialized + 1, + ..counts + }); + Self + } +} + +impl TrackedProbeState { + pub(super) fn record_decode(&self) { + let counts = COUNTS.get(); + COUNTS.set(ProbeCounts { + decoded: counts.decoded + 1, + ..counts + }); + } +} + +impl Drop for TrackedProbeState { + fn drop(&mut self) { + let counts = COUNTS.get(); + COUNTS.set(ProbeCounts { + dropped: counts.dropped + 1, + ..counts + }); + } +} + +/// A one-off or retained probe over `array`, the retained one living in `retained`. +fn probe_for<'a>( + array: &'a ArrayRef, + retained: &'a mut Option, + repeated: bool, +) -> ArrayProbe<'a> { + if repeated { + retained.insert(array.repeated_probe()).as_probe() + } else { + array.probe() + } +} + +fn stacked_runend( + runs: u32, + page_size: usize, + nested: bool, + nullable: bool, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let ends = PrimitiveArray::from_iter((1..=runs).map(|run| run * 4)); + let ends = Pco::from_primitive(ends.as_view(), 3, page_size, ctx)?.into_array(); + let values = if nested { + stacked_runend(runs / 4, page_size, false, nullable, ctx)? + } else { + let values = PrimitiveArray::new( + (0..runs).collect::>(), + if nullable { + Validity::from_iter((0..runs).map(|run| run % 11 != 0)) + } else { + Validity::NonNullable + }, + ); + Pco::from_primitive(values.as_view(), 3, page_size, ctx)?.into_array() + }; + Ok(RunEnd::try_new(ends, values, ctx)?.into_array()) +} + +#[test] +fn stacked_runend_reuses_each_child_state_and_drops_it_once() -> VortexResult<()> { + let session = vortex_array::array_session(); + vortex_runend::initialize(&session); + let mut ctx = session.create_execution_ctx(); + // RunEnd(PCO, RunEnd(PCO, PCO)): three leaves, each fitting in one page. + let array = stacked_runend(256, 512, true, false, &mut ctx)?; + COUNTS.set(ProbeCounts::default()); + for lifetime in 1..=2 { + let mut probe = array.repeated_probe(); + assert_eq!(COUNTS.get().initialized, 3 * (lifetime - 1)); + assert!(probe.execute_scalar(array.len(), &mut ctx).is_err()); + assert_eq!(COUNTS.get().initialized, 3 * (lifetime - 1)); + for index in [1, 5, 127, 255, 511, 1023, 0, 1] { + assert_eq!( + probe.execute_scalar(index, &mut ctx)?, + u32::try_from(index / 16)?.into() + ); + assert_eq!( + COUNTS.get(), + ProbeCounts { + initialized: 3 * lifetime, + decoded: 3 * lifetime, + dropped: 3 * (lifetime - 1), + } + ); + } + drop(probe); + assert_eq!(COUNTS.get().dropped, 3 * lifetime); + } + let before = COUNTS.get(); + let mut once = array.probe(); + for index in [1, 17, 511] { + assert_eq!( + once.execute_scalar(index, &mut ctx)?, + u32::try_from(index / 16)?.into() + ); + } + assert_eq!(COUNTS.get(), before); + Ok(()) +} + +#[rstest] +fn stacked_runend_random_access( + #[values(false, true)] repeated: bool, + #[values(false, true)] nested: bool, + #[values(false, true)] nullable: bool, + #[values(false, true)] sliced: bool, +) -> VortexResult<()> { + let session = vortex_array::array_session(); + vortex_runend::initialize(&session); + let mut ctx = session.create_execution_ctx(); + let encoded = stacked_runend(4096, 128, nested, nullable, &mut ctx)?; + let range = if sliced { 777..15333 } else { 0..16384 }; + let source = if sliced { + encoded + .slice(range.clone())? + .execute::(&mut ctx)? + } else { + encoded + }; + assert!(source.is::()); + let indices = [ + 0, + 1, + 3, + 4, + 15, + 16, + 44, + 176, + 511, + 512, + 1023, + 1024, + 4097, + source.len() - 1, + 0, + ]; + let mut actual = builder_with_capacity_in(source.dtype(), indices.len(), ctx.allocator()); + let mut retained = None; + let mut probe = probe_for(&source, &mut retained, repeated); + for &index in &indices { + actual.append_scalar(&probe.execute_scalar(index, &mut ctx)?)?; + } + let expected = indices + .map(|index| u32::try_from((range.start + index) / if nested { 16 } else { 4 })) + .into_iter() + .collect::, _>>()?; + let validity = if nullable { + Validity::from_iter(expected.iter().map(|value| value % 11 != 0)) + } else { + Validity::NonNullable + }; + let expected = PrimitiveArray::new(expected, validity); + assert_arrays_eq!(actual.finish(), expected, &mut ctx); + Ok(()) +} + +#[rstest] +fn sliced_nullable_random_access( + #[values(false, true)] repeated: bool, + #[values(false, true)] sliced: bool, +) -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = + PrimitiveArray::from_option_iter((0..4096i32).map(|i| (i % 7 != 0).then_some(i * 19))); + let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?.into_array(); + let range = if sliced { 777..3333 } else { 0..4096 }; + let source = encoded.slice(range.clone())?; + assert!(source.is::()); + let input = input.slice(range)?; + let indices = [0u32, 1, 127, 128, 512, 2048, 129, 2, 1024, 0]; + let mut actual = builder_with_capacity_in(source.dtype(), indices.len(), ctx.allocator()); + let mut retained = None; + let mut probe = probe_for(&source, &mut retained, repeated); + for index in indices { + actual.append_scalar(&probe.execute_scalar(index as usize, &mut ctx)?)?; + } + assert_arrays_eq!( + actual.finish(), + input.take(PrimitiveArray::from_iter(indices).into_array())?, + &mut ctx + ); + assert!(probe.execute_scalar(source.len(), &mut ctx).is_err()); + Ok(()) +} + +#[rstest] +#[case(PrimitiveArray::from_iter([1u16, 9, 32768, 65535]))] +#[case(PrimitiveArray::from_iter([i64::MIN, -1, 0, i64::MAX]))] +#[case(PrimitiveArray::from_iter([1.25f64, -2.5, 0.0, f64::INFINITY]))] +fn preserves_physical_type(#[case] input: PrimitiveArray) -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let encoded = Pco::from_primitive(input.as_view(), 3, 2, &mut ctx)?.into_array(); + let mut probe = encoded.repeated_probe(); + let mut actual = builder_with_capacity_in(input.dtype(), input.len(), ctx.allocator()); + for i in 0..input.len() { + actual.append_scalar(&probe.execute_scalar(i, &mut ctx)?)?; + } + assert_arrays_eq!(actual.finish(), input, &mut ctx); + Ok(()) +} + +#[test] +fn decodes_once_per_cached_page_and_evicts() -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = PrimitiveArray::from_iter(0..4096i32); + let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?; + let mut state = PcoProbeState::default(); + for index in [1, 5, 2, 100, 7] { + assert_eq!( + super::scalar_at(encoded.as_view(), index, Some(&mut state), &mut ctx)?, + i32::try_from(index)?.into() + ); + } + assert_eq!(state.decoded_pages, 1); + super::scalar_at(encoded.as_view(), 2048, Some(&mut state), &mut ctx)?; + assert_eq!(state.decoded_pages, 2); + super::scalar_at(encoded.as_view(), 1, Some(&mut state), &mut ctx)?; + assert_eq!(state.decoded_pages, 3); + Ok(()) +} + +#[test] +fn all_null_access_does_not_decode() -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = PrimitiveArray::new(vec![0i32; 128], Validity::AllInvalid); + let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?; + let encoded = encoded.into_array(); + COUNTS.set(ProbeCounts::default()); + let mut probe = encoded.repeated_probe(); + assert!(probe.execute_scalar(42, &mut ctx)?.is_null()); + // The row is null, so the encoding is never entered and no state is built. + assert_eq!(COUNTS.get(), ProbeCounts::default()); + Ok(()) +} + +#[test] +fn probe_retains_pages_after_source_is_dropped() -> VortexResult<()> { + let session = vortex_array::array_session(); + vortex_runend::initialize(&session); + let mut ctx = session.create_execution_ctx(); + let array = stacked_runend(256, 512, true, false, &mut ctx)?; + COUNTS.set(ProbeCounts::default()); + let probe = RepeatedArrayProbe::new(array.clone()); + drop(array); + let mut moved = (probe, ()); + for index in [1, 5, 127, 255, 511, 1023, 0, 1] { + assert_eq!( + moved.0.execute_scalar(index, &mut ctx)?, + u32::try_from(index / 16)?.into() + ); + assert_eq!( + COUNTS.get(), + ProbeCounts { + initialized: 3, + decoded: 3, + dropped: 0 + } + ); + } + assert!(moved.0.execute_scalar(4096, &mut ctx).is_err()); + drop(moved); + assert_eq!(COUNTS.get().dropped, 3); + Ok(()) +} diff --git a/encodings/runend/src/lib.rs b/encodings/runend/src/lib.rs index b9afd20151c..98c661e6e36 100644 --- a/encodings/runend/src/lib.rs +++ b/encodings/runend/src/lib.rs @@ -15,6 +15,7 @@ pub mod decompress_bool; mod iter; mod kernel; pub mod ops; +mod probe; mod rules; #[cfg(test)] #[cfg(not(codspeed))] diff --git a/encodings/runend/src/ops.rs b/encodings/runend/src/ops.rs index 46949d84e1b..637fa8ecd25 100644 --- a/encodings/runend/src/ops.rs +++ b/encodings/runend/src/ops.rs @@ -4,6 +4,7 @@ use vortex_array::ArrayRef; use vortex_array::ArrayView; use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; use vortex_array::match_each_unsigned_integer_ptype; use vortex_array::scalar::Scalar; use vortex_array::search_sorted::SearchResult; @@ -14,19 +15,24 @@ use vortex_array::vtable::OperationsVTable; use vortex_error::VortexResult; use crate::RunEnd; -use crate::array::RunEndArrayExt; -use crate::array::RunEndArraySlotsExt; impl OperationsVTable for RunEnd { type ProbeState = (); + fn probe_scalar( + state: &mut ProbeState<'_, RunEnd>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + crate::probe::scalar_at(state, index, ctx) + } + fn scalar_at( array: ArrayView<'_, RunEnd>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let physical_index = array.find_physical_index(index, ctx)?; - array.values().execute_scalar(physical_index, ctx) + crate::probe::scalar_at(&mut ProbeState::once(array), index, ctx) } } diff --git a/encodings/runend/src/probe.rs b/encodings/runend/src/probe.rs new file mode 100644 index 00000000000..abcb3d22e0d --- /dev/null +++ b/encodings/runend/src/probe.rs @@ -0,0 +1,50 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! RunEnd probes retain their ends and values child probes across lookups. + +use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; +use vortex_array::scalar::Scalar; +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::RunEnd; +use crate::RunEndArrayExt; +use crate::RunEndSlots; + +pub(crate) fn scalar_at( + state: &mut ProbeState<'_, RunEnd>, + index: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let array = state.array(); + let logical_index = array + .offset() + .checked_add(index) + .ok_or_else(|| vortex_err!("RunEnd logical index overflow"))?; + // Search for the first end strictly greater than the logical index. Every comparison + // uses the same ends probe, so a retained read keeps the child's preparation within and + // between searches. + let mut ends = state + .slot(RunEndSlots::ENDS)? + .ok_or_else(|| vortex_err!("RunEnd ends slot is missing"))?; + let mut left = 0; + let mut right = ends.array().len(); + while left < right { + let mid = left + (right - left) / 2; + let end = usize::try_from(&ends.execute_scalar(mid, ctx)?)?; + if end <= logical_index { + left = mid + 1; + } else { + right = mid; + } + } + state + .slot(RunEndSlots::VALUES)? + .ok_or_else(|| vortex_err!("RunEnd values slot is missing"))? + .execute_scalar(left, ctx) +} + +#[cfg(test)] +mod tests; diff --git a/encodings/runend/src/probe/tests.rs b/encodings/runend/src/probe/tests.rs new file mode 100644 index 00000000000..c57b930b61e --- /dev/null +++ b/encodings/runend/src/probe/tests.rs @@ -0,0 +1,102 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use rstest::rstest; +use vortex_array::ArrayProbe; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::RepeatedArrayProbe; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::arrays::VarBinViewArray; +use vortex_array::assert_arrays_eq; +use vortex_array::builders::builder_with_capacity_in; +use vortex_array::builtins::ArrayBuiltins; +use vortex_array::dtype::DType; +use vortex_array::dtype::Nullability; +use vortex_array::scalar_fn::fns::operators::Operator; +use vortex_array::validity::Validity; +use vortex_buffer::Buffer; +use vortex_buffer::buffer; +use vortex_error::VortexResult; + +use crate::RunEnd; +use crate::tests::SESSION; + +/// A one-off or retained probe over `array`, the retained one living in `retained`. +fn probe_for<'a>( + array: &'a ArrayRef, + retained: &'a mut Option, + repeated: bool, +) -> ArrayProbe<'a> { + if repeated { + retained.insert(array.repeated_probe()).as_probe() + } else { + array.probe() + } +} + +#[rstest] +fn sliced_strings_and_null_runs(#[values(false, true)] repeated: bool) -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let array = RunEnd::try_new_offset_length( + buffer![2u8, 5, 9].into_array(), + VarBinViewArray::from_iter( + [Some("first"), None, Some("last")], + DType::Utf8(Nullability::Nullable), + ) + .into_array(), + 1, + 7, + &mut ctx, + )? + .into_array(); + let mut retained = None; + let mut probe = probe_for(&array, &mut retained, repeated); + let mut actual = builder_with_capacity_in(array.dtype(), 5, ctx.allocator()); + for index in [6, 0, 1, 3, 4] { + actual.append_scalar(&probe.execute_scalar(index, &mut ctx)?)?; + } + assert_arrays_eq!( + actual.finish(), + VarBinViewArray::from_iter( + [Some("last"), Some("first"), None, None, Some("last")], + DType::Utf8(Nullability::Nullable) + ), + &mut ctx + ); + assert!(probe.execute_scalar(array.len(), &mut ctx).is_err()); + Ok(()) +} + +#[test] +fn empty_probe_checks_bounds() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let array = RunEnd::try_new( + Buffer::::empty().into_array(), + Buffer::::empty().into_array(), + &mut ctx, + )? + .into_array(); + assert!(array.probe().execute_scalar(0, &mut ctx).is_err()); + assert!(array.repeated_probe().execute_scalar(0, &mut ctx).is_err()); + Ok(()) +} + +#[test] +fn lazy_value_validity_only_evaluates_requested_runs() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let numerators = buffer![1u32, 1].into_array(); + let validity = numerators + .binary(buffer![1u32, 0].into_array(), Operator::Div)? + .binary(numerators, Operator::Eq)?; + let values = PrimitiveArray::new(buffer![42u32, 99], Validity::Array(validity)); + let array = + RunEnd::try_new(buffer![4u32, 8].into_array(), values.into_array(), &mut ctx)?.into_array(); + let expected = array.execute_scalar(0, &mut ctx)?; + let mut probe = array.repeated_probe(); + assert_eq!(probe.execute_scalar(0, &mut ctx)?, expected); + assert_eq!(probe.execute_scalar(3, &mut ctx)?, expected); + assert!(probe.execute_scalar(4, &mut ctx).is_err()); + Ok(()) +} From d4ad1023acc475899d43164985399191e3d6f3f0 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Tue, 15 Sep 2026 22:48:27 +0100 Subject: [PATCH 2/7] bench(pco): read the scalar bench through a repeated probe Drops the probe bench added by this PR; the scalar bench from #9896 is the baseline, and reading it through one `RepeatedArrayProbe` shows the retained page decode on the clustered cases. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_016CqrLKgPqFYGZK5sjk1qe7 Signed-off-by: Joe Isaacs --- encodings/pco/Cargo.toml | 4 - encodings/pco/benches/probe.rs | 198 -------------------------------- encodings/pco/benches/scalar.rs | 17 ++- 3 files changed, 12 insertions(+), 207 deletions(-) delete mode 100644 encodings/pco/benches/probe.rs diff --git a/encodings/pco/Cargo.toml b/encodings/pco/Cargo.toml index b80a81547d2..eb9ab5ac771 100644 --- a/encodings/pco/Cargo.toml +++ b/encodings/pco/Cargo.toml @@ -33,10 +33,6 @@ rstest = { workspace = true } vortex-array = { workspace = true, features = ["_test-harness"] } vortex-arrow = { workspace = true } -[[bench]] -name = "probe" -harness = false - [[bench]] name = "scalar" harness = false diff --git a/encodings/pco/benches/probe.rs b/encodings/pco/benches/probe.rs deleted file mode 100644 index 4ecc47074cb..00000000000 --- a/encodings/pco/benches/probe.rs +++ /dev/null @@ -1,198 +0,0 @@ -// SPDX-License-Identifier: Apache-2.0 -// SPDX-FileCopyrightText: Copyright the Vortex contributors - -//! Scalar access comparisons, including preparation and teardown for every group of lookups. - -use std::hint::black_box; -use std::sync::LazyLock; - -use divan::Bencher; -use vortex_array::ArrayRef; -use vortex_array::IntoArray; -use vortex_array::VortexSessionExecute; -use vortex_array::arrays::PrimitiveArray; -use vortex_array::validity::Validity; -use vortex_error::VortexExpect; -use vortex_fastlanes::RLEData; -use vortex_pco::Pco; -use vortex_runend::RunEnd; -use vortex_session::VortexSession; - -fn main() { - divan::main(); -} - -const LEN: usize = 16_384; -// Number of accesses, nullable, scattered (otherwise clustered within 256 rows). -const CASES: &[(usize, bool, bool)] = &[ - (1, false, false), - (1, true, false), - (64, false, false), - (64, true, false), - (64, false, true), - (64, true, true), - (1024, false, false), - (1024, true, false), - (1024, false, true), - (1024, true, true), -]; - -static SESSION: LazyLock = LazyLock::new(|| { - let session = vortex_array::array_session(); - vortex_fastlanes::initialize(&session); - vortex_runend::initialize(&session); - session -}); - -fn input(nullable: bool) -> PrimitiveArray { - let validity = if nullable { - Validity::from_iter((0..LEN).map(|i| i % 11 != 0)) - } else { - Validity::NonNullable - }; - PrimitiveArray::new( - (0..LEN) - .map(|i| u32::try_from(i / 16).vortex_expect("fixture values fit u32")) - .collect::>(), - validity, - ) -} - -fn indices(count: usize, scattered: bool) -> Vec { - let span = if scattered { LEN } else { 256 }; - let base = if scattered { 0 } else { 4096 }; - let mut seed = 42u64; - (0..count) - .map(|_| { - seed = seed.wrapping_mul(6364136223846793005).wrapping_add(1); - base + ((seed >> 32) as usize % span) - }) - .collect() -} - -fn rle(nullable: bool) -> ArrayRef { - let mut ctx = SESSION.create_execution_ctx(); - RLEData::encode(input(nullable).as_view(), &mut ctx) - .vortex_expect("RLE compression") - .into_array() -} - -fn pco(nullable: bool) -> ArrayRef { - let mut ctx = SESSION.create_execution_ctx(); - Pco::from_primitive(input(nullable).as_view(), 8, 1024, &mut ctx) - .vortex_expect("PCO compression") - .into_array() -} - -fn runend_pco(nullable: bool) -> ArrayRef { - let mut ctx = SESSION.create_execution_ctx(); - let runs = u32::try_from(LEN / 4).vortex_expect("run count fits u32"); - let ends = PrimitiveArray::from_iter((1..=runs).map(|run| run * 4)); - let values = PrimitiveArray::new( - (0..runs).collect::>(), - if nullable { - Validity::from_iter((0..runs).map(|run| run % 11 != 0)) - } else { - Validity::NonNullable - }, - ); - RunEnd::try_new( - Pco::from_primitive(ends.as_view(), 8, 1024, &mut ctx) - .vortex_expect("PCO ends compression") - .into_array(), - Pco::from_primitive(values.as_view(), 8, 1024, &mut ctx) - .vortex_expect("PCO values compression") - .into_array(), - &mut ctx, - ) - .vortex_expect("RunEnd construction") - .into_array() -} - -fn execute_scalar(bencher: Bencher, array: ArrayRef, indices: &[usize]) { - bencher - .with_inputs(|| SESSION.create_execution_ctx()) - .bench_refs(|ctx| { - for &index in indices { - black_box( - array - .execute_scalar(black_box(index), ctx) - .vortex_expect("scalar access"), - ); - } - }); -} - -fn probe(bencher: Bencher, array: ArrayRef, indices: &[usize], repeated: bool) { - bencher - .with_inputs(|| SESSION.create_execution_ctx()) - .bench_refs(|ctx| { - let mut retained = repeated.then(|| array.repeated_probe()); - let mut probe = match &mut retained { - Some(retained) => retained.as_probe(), - None => array.probe(), - }; - for &index in indices { - black_box( - probe - .execute_scalar(black_box(index), ctx) - .vortex_expect("probe access"), - ); - } - }); -} - -#[divan::bench(args = CASES)] -fn rle_probe(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { - probe( - bencher, - rle(nullable), - &indices(count, scattered), - count != 1, - ); -} - -#[divan::bench(args = CASES)] -fn pco_probe(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { - probe( - bencher, - pco(nullable), - &indices(count, scattered), - count != 1, - ); -} - -#[divan::bench(args = [false, true])] -fn rle_repeated_first(bencher: Bencher, nullable: bool) { - probe(bencher, rle(nullable), &indices(1, false), true); -} - -#[divan::bench(args = [false, true])] -fn pco_repeated_first(bencher: Bencher, nullable: bool) { - probe(bencher, pco(nullable), &indices(1, false), true); -} - -#[divan::bench(args = CASES)] -fn rle_execute_scalar(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { - execute_scalar(bencher, rle(nullable), &indices(count, scattered)); -} - -#[divan::bench(args = CASES)] -fn pco_execute_scalar(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { - execute_scalar(bencher, pco(nullable), &indices(count, scattered)); -} - -#[divan::bench(args = CASES)] -fn runend_pco_execute_scalar(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { - execute_scalar(bencher, runend_pco(nullable), &indices(count, scattered)); -} - -#[divan::bench(args = CASES)] -fn runend_pco_probe(bencher: Bencher, (count, nullable, scattered): (usize, bool, bool)) { - probe( - bencher, - runend_pco(nullable), - &indices(count, scattered), - count != 1, - ); -} diff --git a/encodings/pco/benches/scalar.rs b/encodings/pco/benches/scalar.rs index eff55fceb8c..6a928c4b4a0 100644 --- a/encodings/pco/benches/scalar.rs +++ b/encodings/pco/benches/scalar.rs @@ -1,8 +1,9 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors -//! Scalar reads out of a PCO array. Cases are `(access_count, nullable, scattered)`; clustered -//! indices stay inside one page so a retained decode can be reused, scattered ones cross pages. +//! Scalar reads out of a PCO array through one repeated probe. Cases are +//! `(access_count, nullable, scattered)`; clustered indices stay inside one page so the retained +//! decode is reused, scattered ones cross pages. use std::sync::LazyLock; @@ -65,11 +66,17 @@ fn scalar_access(bencher: Bencher, (count, nullable, scattered): (usize, bool, b let array = pco(nullable); let indices = indices(count, scattered); bencher - .with_inputs(|| (SESSION.create_execution_ctx(), Vec::with_capacity(count))) - .bench_refs(|(ctx, scalars)| { + .with_inputs(|| { + ( + SESSION.create_execution_ctx(), + array.repeated_probe(), + Vec::with_capacity(count), + ) + }) + .bench_refs(|(ctx, probe, scalars)| { for &index in &indices { scalars.push( - array + probe .execute_scalar(index, ctx) .vortex_expect("scalar access"), ); From 0d5143bc759c1e569d3d68e4cf9bf65f55aa90a1 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Tue, 15 Sep 2026 23:06:09 +0100 Subject: [PATCH 3/7] perf(pco): retain every decoded page in the probe state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Each page is decoded at most once per probe, so scattered reads stop paying a page decode per row. Storage grows with the pages touched and never beyond the decompressed array. The last page is checked before the binary search so clustered reads keep their cost. Local medians, 1024 lookups, before vs after: scattered 5.4 ms vs 122 µs; clustered 25 µs vs 24 µs. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_016CqrLKgPqFYGZK5sjk1qe7 Signed-off-by: Joe Isaacs --- encodings/pco/src/probe.rs | 75 ++++++++++++++++++-------------- encodings/pco/src/probe/tests.rs | 11 +++-- 2 files changed, 51 insertions(+), 35 deletions(-) diff --git a/encodings/pco/src/probe.rs b/encodings/pco/src/probe.rs index 771172200a3..f429081d042 100644 --- a/encodings/pco/src/probe.rs +++ b/encodings/pco/src/probe.rs @@ -1,7 +1,7 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors -//! PCO probes retain one decoded page and index the compacted, non-null value positions. +//! PCO probes retain every decoded page and index the compacted, non-null value positions. use std::ops::Range; @@ -33,13 +33,14 @@ const RANK_STRIDE: usize = 512; /// State retained by repeated PCO probes. Construction performs no allocation. /// /// The mask and rank index use unsliced logical rows; page boundaries use compacted value -/// positions. Only the most recently accessed page is retained, bounding decoded storage. +/// positions. Every page decoded so far is retained, so the state grows with the pages touched +/// and never beyond the decompressed array. #[derive(Default)] pub struct PcoProbeState { validity: Option, rank: Vec, pages: Vec, - decoded: Option<(Range, PrimitiveArray)>, + last_page: usize, #[cfg(test)] decoded_pages: usize, #[cfg(test)] @@ -49,6 +50,7 @@ pub struct PcoProbeState { struct Page { values: Range, chunk: usize, + decoded: Option, } pub(crate) fn scalar_at( @@ -93,44 +95,52 @@ pub(crate) fn scalar_at( } }; - let (range, values) = match &state.decoded { - Some(decoded) if decoded.0.contains(&value_index) => decoded, - _ => { - if state.pages.is_empty() { - let mut start = 0; - for (chunk, metadata) in array.metadata.chunks.iter().enumerate() { - for page in &metadata.pages { - let end = start + page.n_values as usize; - state.pages.push(Page { - values: start..end, - chunk, - }); - start = end; - } - } + if state.pages.is_empty() { + let mut start = 0; + for (chunk, metadata) in array.metadata.chunks.iter().enumerate() { + for page in &metadata.pages { + let end = start + page.n_values as usize; + state.pages.push(Page { + values: start..end, + chunk, + decoded: None, + }); + start = end; } - let page_index = state - .pages - .partition_point(|page| page.values.end <= value_index); - let page = state - .pages - .get(page_index) - .ok_or_else(|| vortex_err!("Missing PCO page for value {value_index}"))?; + } + } + let page_index = if state.pages[state.last_page].values.contains(&value_index) { + state.last_page + } else { + state + .pages + .partition_point(|page| page.values.end <= value_index) + }; + state.last_page = page_index; + let page = state + .pages + .get_mut(page_index) + .ok_or_else(|| vortex_err!("Missing PCO page for value {value_index}"))?; + let values = match &mut page.decoded { + Some(decoded) => decoded, + slot @ None => { let decoded = match_number_enum!( number_type_from_ptype(array.dtype().as_ptype()), - NumberType => { decode_page::(array, page, array.pages[page_index].as_slice(), ctx)? } + NumberType => { + decode_page::(array, page.chunk, page.values.len(), array.pages[page_index].as_slice(), ctx)? + } ); #[cfg(test)] { state.decoded_pages += 1; state.tracking.record_decode(); } - state.decoded.insert((page.values.clone(), decoded)) + slot.insert(decoded) } }; Ok(match_each_native_ptype!(values.ptype(), |T| { Scalar::primitive( - values.as_slice::()[value_index - range.start], + values.as_slice::()[value_index - page.values.start], array.dtype().nullability(), ) })) @@ -138,19 +148,20 @@ pub(crate) fn scalar_at( fn decode_page( array: ArrayView<'_, Pco>, - page: &Page, + chunk: usize, + n_values: usize, buffer: &[u8], ctx: &mut ExecutionCtx, ) -> VortexResult { let (file, _) = FileDecompressor::new(array.metadata.header.as_slice()).map_err(vortex_err_from_pco)?; let (mut chunk, _) = file - .chunk_decompressor::(array.chunk_metas[page.chunk].as_ref()) + .chunk_decompressor::(array.chunk_metas[chunk].as_ref()) .map_err(vortex_err_from_pco)?; let mut decoder = chunk - .page_decompressor(buffer, page.values.len()) + .page_decompressor(buffer, n_values) .map_err(vortex_err_from_pco)?; - let mut values = BufferMut::::zeroed_in(page.values.len(), ctx.allocator().clone()); + let mut values = BufferMut::::zeroed_in(n_values, ctx.allocator().clone()); decoder.read(&mut values).map_err(vortex_err_from_pco)?; Ok(PrimitiveArray::new(values.freeze(), Validity::NonNullable)) } diff --git a/encodings/pco/src/probe/tests.rs b/encodings/pco/src/probe/tests.rs index ee42d2db596..8b1af6e3edb 100644 --- a/encodings/pco/src/probe/tests.rs +++ b/encodings/pco/src/probe/tests.rs @@ -247,7 +247,7 @@ fn preserves_physical_type(#[case] input: PrimitiveArray) -> VortexResult<()> { } #[test] -fn decodes_once_per_cached_page_and_evicts() -> VortexResult<()> { +fn decodes_each_page_once() -> VortexResult<()> { let mut ctx = vortex_array::array_session().create_execution_ctx(); let input = PrimitiveArray::from_iter(0..4096i32); let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?; @@ -261,8 +261,13 @@ fn decodes_once_per_cached_page_and_evicts() -> VortexResult<()> { assert_eq!(state.decoded_pages, 1); super::scalar_at(encoded.as_view(), 2048, Some(&mut state), &mut ctx)?; assert_eq!(state.decoded_pages, 2); - super::scalar_at(encoded.as_view(), 1, Some(&mut state), &mut ctx)?; - assert_eq!(state.decoded_pages, 3); + for index in [1, 2049, 100, 2100] { + assert_eq!( + super::scalar_at(encoded.as_view(), index, Some(&mut state), &mut ctx)?, + i32::try_from(index)?.into() + ); + } + assert_eq!(state.decoded_pages, 2); Ok(()) } From a0f1a8d00b8ea73da624c7052c58f9a677dd0e5a Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Wed, 16 Sep 2026 09:30:07 +0100 Subject: [PATCH 4/7] refactor(pco): drop the probe example and test-only probe state Removes the example and the unused fastlanes dev-dependency, and the `cfg(test)` counters on `PcoProbeState`; the tests that read them are replaced by behavioural checks. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_016CqrLKgPqFYGZK5sjk1qe7 Signed-off-by: Joe Isaacs --- Cargo.lock | 1 - encodings/pco/Cargo.toml | 1 - encodings/pco/examples/probe.rs | 42 ---------- encodings/pco/src/probe.rs | 9 --- encodings/pco/src/probe/tests.rs | 132 +------------------------------ 5 files changed, 2 insertions(+), 183 deletions(-) delete mode 100644 encodings/pco/examples/probe.rs diff --git a/Cargo.lock b/Cargo.lock index 88f8586fa61..be52fc7f3d2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11419,7 +11419,6 @@ dependencies = [ "vortex-arrow", "vortex-buffer", "vortex-error", - "vortex-fastlanes", "vortex-mask", "vortex-runend", "vortex-session", diff --git a/encodings/pco/Cargo.toml b/encodings/pco/Cargo.toml index eb9ab5ac771..363c11bce6f 100644 --- a/encodings/pco/Cargo.toml +++ b/encodings/pco/Cargo.toml @@ -27,7 +27,6 @@ vortex-session = { workspace = true } [dev-dependencies] divan = { workspace = true } -vortex-fastlanes = { workspace = true } vortex-runend = { workspace = true } rstest = { workspace = true } vortex-array = { workspace = true, features = ["_test-harness"] } diff --git a/encodings/pco/examples/probe.rs b/encodings/pco/examples/probe.rs deleted file mode 100644 index a8183043216..00000000000 --- a/encodings/pco/examples/probe.rs +++ /dev/null @@ -1,42 +0,0 @@ -// SPDX-License-Identifier: Apache-2.0 -// SPDX-FileCopyrightText: Copyright the Vortex contributors - -//! Run with `cargo run -p vortex-pco --example probe`. - -use vortex_array::IntoArray; -use vortex_array::VortexSessionExecute; -use vortex_array::arrays::PrimitiveArray; -use vortex_error::VortexResult; -use vortex_fastlanes::RLEData; -use vortex_pco::Pco; -use vortex_runend::RunEnd; - -fn main() -> VortexResult<()> { - let session = vortex_array::array_session(); - vortex_fastlanes::initialize(&session); - vortex_runend::initialize(&session); - let mut ctx = session.create_execution_ctx(); - let values = - PrimitiveArray::from_option_iter((0..4096u32).map(|i| (i % 11 != 0).then_some(i / 16))); - let rle = RLEData::encode(values.as_view(), &mut ctx)?.into_array(); - let pco = Pco::from_primitive(values.as_view(), 3, 1024, &mut ctx)?.into_array(); - let ends = PrimitiveArray::from_iter((1..=4096u32).map(|run| run * 4)); - let ends = Pco::from_primitive(ends.as_view(), 3, 1024, &mut ctx)?.into_array(); - let runend_pco = RunEnd::try_new(ends, pco.clone(), &mut ctx)?.into_array(); - - for (name, array) in [("RLE", rle), ("PCO", pco), ("RunEnd(PCO, PCO)", runend_pco)] { - // A one-off probe borrows the array and retains nothing. - let scalar = array.probe().execute_scalar(17, &mut ctx)?; - println!("{name} single lookup: {scalar}"); - - // A repeated probe keeps child probes; each PCO child keeps its own decoded page. - let mut probe = array.repeated_probe(); - for index in [17, 33, 22, 1025, 17] { - println!( - "{name}[{index}] = {}", - probe.execute_scalar(index, &mut ctx)? - ); - } - } - Ok(()) -} diff --git a/encodings/pco/src/probe.rs b/encodings/pco/src/probe.rs index f429081d042..5fb31893d0d 100644 --- a/encodings/pco/src/probe.rs +++ b/encodings/pco/src/probe.rs @@ -41,10 +41,6 @@ pub struct PcoProbeState { rank: Vec, pages: Vec, last_page: usize, - #[cfg(test)] - decoded_pages: usize, - #[cfg(test)] - tracking: tests::TrackedProbeState, } struct Page { @@ -130,11 +126,6 @@ pub(crate) fn scalar_at( decode_page::(array, page.chunk, page.values.len(), array.pages[page_index].as_slice(), ctx)? } ); - #[cfg(test)] - { - state.decoded_pages += 1; - state.tracking.record_decode(); - } slot.insert(decoded) } }; diff --git a/encodings/pco/src/probe/tests.rs b/encodings/pco/src/probe/tests.rs index 8b1af6e3edb..57ff92f8ccd 100644 --- a/encodings/pco/src/probe/tests.rs +++ b/encodings/pco/src/probe/tests.rs @@ -1,8 +1,6 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors -use std::cell::Cell; - use rstest::rstest; use vortex_array::ArrayProbe; use vortex_array::ArrayRef; @@ -17,53 +15,8 @@ use vortex_array::validity::Validity; use vortex_error::VortexResult; use vortex_runend::RunEnd; -use super::PcoProbeState; use crate::Pco; -#[derive(Clone, Copy, Debug, Default, PartialEq)] -struct ProbeCounts { - initialized: usize, - decoded: usize, - dropped: usize, -} - -thread_local! { - static COUNTS: Cell = Cell::default(); -} - -pub(super) struct TrackedProbeState; - -impl Default for TrackedProbeState { - fn default() -> Self { - let counts = COUNTS.get(); - COUNTS.set(ProbeCounts { - initialized: counts.initialized + 1, - ..counts - }); - Self - } -} - -impl TrackedProbeState { - pub(super) fn record_decode(&self) { - let counts = COUNTS.get(); - COUNTS.set(ProbeCounts { - decoded: counts.decoded + 1, - ..counts - }); - } -} - -impl Drop for TrackedProbeState { - fn drop(&mut self) { - let counts = COUNTS.get(); - COUNTS.set(ProbeCounts { - dropped: counts.dropped + 1, - ..counts - }); - } -} - /// A one-off or retained probe over `array`, the retained one living in `retained`. fn probe_for<'a>( array: &'a ArrayRef, @@ -102,48 +55,6 @@ fn stacked_runend( Ok(RunEnd::try_new(ends, values, ctx)?.into_array()) } -#[test] -fn stacked_runend_reuses_each_child_state_and_drops_it_once() -> VortexResult<()> { - let session = vortex_array::array_session(); - vortex_runend::initialize(&session); - let mut ctx = session.create_execution_ctx(); - // RunEnd(PCO, RunEnd(PCO, PCO)): three leaves, each fitting in one page. - let array = stacked_runend(256, 512, true, false, &mut ctx)?; - COUNTS.set(ProbeCounts::default()); - for lifetime in 1..=2 { - let mut probe = array.repeated_probe(); - assert_eq!(COUNTS.get().initialized, 3 * (lifetime - 1)); - assert!(probe.execute_scalar(array.len(), &mut ctx).is_err()); - assert_eq!(COUNTS.get().initialized, 3 * (lifetime - 1)); - for index in [1, 5, 127, 255, 511, 1023, 0, 1] { - assert_eq!( - probe.execute_scalar(index, &mut ctx)?, - u32::try_from(index / 16)?.into() - ); - assert_eq!( - COUNTS.get(), - ProbeCounts { - initialized: 3 * lifetime, - decoded: 3 * lifetime, - dropped: 3 * (lifetime - 1), - } - ); - } - drop(probe); - assert_eq!(COUNTS.get().dropped, 3 * lifetime); - } - let before = COUNTS.get(); - let mut once = array.probe(); - for index in [1, 17, 511] { - assert_eq!( - once.execute_scalar(index, &mut ctx)?, - u32::try_from(index / 16)?.into() - ); - } - assert_eq!(COUNTS.get(), before); - Ok(()) -} - #[rstest] fn stacked_runend_random_access( #[values(false, true)] repeated: bool, @@ -247,51 +158,22 @@ fn preserves_physical_type(#[case] input: PrimitiveArray) -> VortexResult<()> { } #[test] -fn decodes_each_page_once() -> VortexResult<()> { - let mut ctx = vortex_array::array_session().create_execution_ctx(); - let input = PrimitiveArray::from_iter(0..4096i32); - let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?; - let mut state = PcoProbeState::default(); - for index in [1, 5, 2, 100, 7] { - assert_eq!( - super::scalar_at(encoded.as_view(), index, Some(&mut state), &mut ctx)?, - i32::try_from(index)?.into() - ); - } - assert_eq!(state.decoded_pages, 1); - super::scalar_at(encoded.as_view(), 2048, Some(&mut state), &mut ctx)?; - assert_eq!(state.decoded_pages, 2); - for index in [1, 2049, 100, 2100] { - assert_eq!( - super::scalar_at(encoded.as_view(), index, Some(&mut state), &mut ctx)?, - i32::try_from(index)?.into() - ); - } - assert_eq!(state.decoded_pages, 2); - Ok(()) -} - -#[test] -fn all_null_access_does_not_decode() -> VortexResult<()> { +fn all_null_access_returns_null() -> VortexResult<()> { let mut ctx = vortex_array::array_session().create_execution_ctx(); let input = PrimitiveArray::new(vec![0i32; 128], Validity::AllInvalid); let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?; let encoded = encoded.into_array(); - COUNTS.set(ProbeCounts::default()); let mut probe = encoded.repeated_probe(); assert!(probe.execute_scalar(42, &mut ctx)?.is_null()); - // The row is null, so the encoding is never entered and no state is built. - assert_eq!(COUNTS.get(), ProbeCounts::default()); Ok(()) } #[test] -fn probe_retains_pages_after_source_is_dropped() -> VortexResult<()> { +fn probe_outlives_source_and_moves() -> VortexResult<()> { let session = vortex_array::array_session(); vortex_runend::initialize(&session); let mut ctx = session.create_execution_ctx(); let array = stacked_runend(256, 512, true, false, &mut ctx)?; - COUNTS.set(ProbeCounts::default()); let probe = RepeatedArrayProbe::new(array.clone()); drop(array); let mut moved = (probe, ()); @@ -300,17 +182,7 @@ fn probe_retains_pages_after_source_is_dropped() -> VortexResult<()> { moved.0.execute_scalar(index, &mut ctx)?, u32::try_from(index / 16)?.into() ); - assert_eq!( - COUNTS.get(), - ProbeCounts { - initialized: 3, - decoded: 3, - dropped: 0 - } - ); } assert!(moved.0.execute_scalar(4096, &mut ctx).is_err()); - drop(moved); - assert_eq!(COUNTS.get().dropped, 3); Ok(()) } From 07dbdae90cf0b39993f608b3472848823604af1e Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Wed, 16 Sep 2026 09:31:52 +0100 Subject: [PATCH 5/7] test(pco): probe tests use only PCO arrays Drops the runend dev-dependency; random access is covered over nullable, sliced and repeated PCO arrays directly. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_016CqrLKgPqFYGZK5sjk1qe7 Signed-off-by: Joe Isaacs --- Cargo.lock | 1 - encodings/pco/Cargo.toml | 1 - encodings/pco/src/probe/tests.rs | 101 +++---------------------------- 3 files changed, 10 insertions(+), 93 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index be52fc7f3d2..275db7ef8ee 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11420,7 +11420,6 @@ dependencies = [ "vortex-buffer", "vortex-error", "vortex-mask", - "vortex-runend", "vortex-session", ] diff --git a/encodings/pco/Cargo.toml b/encodings/pco/Cargo.toml index 363c11bce6f..749464c865c 100644 --- a/encodings/pco/Cargo.toml +++ b/encodings/pco/Cargo.toml @@ -27,7 +27,6 @@ vortex-session = { workspace = true } [dev-dependencies] divan = { workspace = true } -vortex-runend = { workspace = true } rstest = { workspace = true } vortex-array = { workspace = true, features = ["_test-harness"] } vortex-arrow = { workspace = true } diff --git a/encodings/pco/src/probe/tests.rs b/encodings/pco/src/probe/tests.rs index 57ff92f8ccd..f4c9740cf2c 100644 --- a/encodings/pco/src/probe/tests.rs +++ b/encodings/pco/src/probe/tests.rs @@ -4,7 +4,6 @@ use rstest::rstest; use vortex_array::ArrayProbe; use vortex_array::ArrayRef; -use vortex_array::ExecutionCtx; use vortex_array::IntoArray; use vortex_array::RepeatedArrayProbe; use vortex_array::VortexSessionExecute; @@ -13,7 +12,6 @@ use vortex_array::assert_arrays_eq; use vortex_array::builders::builder_with_capacity_in; use vortex_array::validity::Validity; use vortex_error::VortexResult; -use vortex_runend::RunEnd; use crate::Pco; @@ -30,96 +28,18 @@ fn probe_for<'a>( } } -fn stacked_runend( - runs: u32, - page_size: usize, - nested: bool, - nullable: bool, - ctx: &mut ExecutionCtx, -) -> VortexResult { - let ends = PrimitiveArray::from_iter((1..=runs).map(|run| run * 4)); - let ends = Pco::from_primitive(ends.as_view(), 3, page_size, ctx)?.into_array(); - let values = if nested { - stacked_runend(runs / 4, page_size, false, nullable, ctx)? - } else { - let values = PrimitiveArray::new( - (0..runs).collect::>(), - if nullable { - Validity::from_iter((0..runs).map(|run| run % 11 != 0)) - } else { - Validity::NonNullable - }, - ); - Pco::from_primitive(values.as_view(), 3, page_size, ctx)?.into_array() - }; - Ok(RunEnd::try_new(ends, values, ctx)?.into_array()) -} - #[rstest] -fn stacked_runend_random_access( +fn random_access( #[values(false, true)] repeated: bool, - #[values(false, true)] nested: bool, #[values(false, true)] nullable: bool, #[values(false, true)] sliced: bool, ) -> VortexResult<()> { - let session = vortex_array::array_session(); - vortex_runend::initialize(&session); - let mut ctx = session.create_execution_ctx(); - let encoded = stacked_runend(4096, 128, nested, nullable, &mut ctx)?; - let range = if sliced { 777..15333 } else { 0..16384 }; - let source = if sliced { - encoded - .slice(range.clone())? - .execute::(&mut ctx)? - } else { - encoded - }; - assert!(source.is::()); - let indices = [ - 0, - 1, - 3, - 4, - 15, - 16, - 44, - 176, - 511, - 512, - 1023, - 1024, - 4097, - source.len() - 1, - 0, - ]; - let mut actual = builder_with_capacity_in(source.dtype(), indices.len(), ctx.allocator()); - let mut retained = None; - let mut probe = probe_for(&source, &mut retained, repeated); - for &index in &indices { - actual.append_scalar(&probe.execute_scalar(index, &mut ctx)?)?; - } - let expected = indices - .map(|index| u32::try_from((range.start + index) / if nested { 16 } else { 4 })) - .into_iter() - .collect::, _>>()?; - let validity = if nullable { - Validity::from_iter(expected.iter().map(|value| value % 11 != 0)) + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = if nullable { + PrimitiveArray::from_option_iter((0..4096i32).map(|i| (i % 7 != 0).then_some(i * 19))) } else { - Validity::NonNullable + PrimitiveArray::from_iter((0..4096i32).map(|i| i * 19)) }; - let expected = PrimitiveArray::new(expected, validity); - assert_arrays_eq!(actual.finish(), expected, &mut ctx); - Ok(()) -} - -#[rstest] -fn sliced_nullable_random_access( - #[values(false, true)] repeated: bool, - #[values(false, true)] sliced: bool, -) -> VortexResult<()> { - let mut ctx = vortex_array::array_session().create_execution_ctx(); - let input = - PrimitiveArray::from_option_iter((0..4096i32).map(|i| (i % 7 != 0).then_some(i * 19))); let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?.into_array(); let range = if sliced { 777..3333 } else { 0..4096 }; let source = encoded.slice(range.clone())?; @@ -170,17 +90,16 @@ fn all_null_access_returns_null() -> VortexResult<()> { #[test] fn probe_outlives_source_and_moves() -> VortexResult<()> { - let session = vortex_array::array_session(); - vortex_runend::initialize(&session); - let mut ctx = session.create_execution_ctx(); - let array = stacked_runend(256, 512, true, false, &mut ctx)?; + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = PrimitiveArray::from_iter(0..4096u32); + let array = Pco::from_primitive(input.as_view(), 3, 512, &mut ctx)?.into_array(); let probe = RepeatedArrayProbe::new(array.clone()); drop(array); let mut moved = (probe, ()); - for index in [1, 5, 127, 255, 511, 1023, 0, 1] { + for index in [1, 5, 127, 255, 511, 1023, 0, 1, 4095] { assert_eq!( moved.0.execute_scalar(index, &mut ctx)?, - u32::try_from(index / 16)?.into() + u32::try_from(index)?.into() ); } assert!(moved.0.execute_scalar(4096, &mut ctx).is_err()); From 41fe85301097032d4b6925d98ce36ae32b94b7ce Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Tue, 22 Sep 2026 10:30:57 +0000 Subject: [PATCH 6/7] refactor(pco): make the probe's last page an Option `last_page` defaulted to 0, which claimed a previous probe had landed in page 0 before any probe had run and indexed `pages` before it was known to be non-empty. Making it `Option` states that there is no last page until one is resolved. The page lookup that follows can only miss on an out-of-range value index, which the mask and page table rule out, so it becomes a `vortex_expect` with a static message rather than an error. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01THeokTcUrxdbZcPrDy8x6D --- encodings/pco/src/probe.rs | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/encodings/pco/src/probe.rs b/encodings/pco/src/probe.rs index 5fb31893d0d..7dee35000a4 100644 --- a/encodings/pco/src/probe.rs +++ b/encodings/pco/src/probe.rs @@ -19,8 +19,8 @@ use vortex_array::match_each_native_ptype; use vortex_array::scalar::Scalar; use vortex_array::validity::Validity; use vortex_buffer::BufferMut; +use vortex_error::VortexExpect; use vortex_error::VortexResult; -use vortex_error::vortex_err; use vortex_mask::Mask; use crate::Pco; @@ -40,7 +40,8 @@ pub struct PcoProbeState { validity: Option, rank: Vec, pages: Vec, - last_page: usize, + /// The page the previous probe landed in, `None` until the first probe resolves one. + last_page: Option, } struct Page { @@ -105,18 +106,17 @@ pub(crate) fn scalar_at( } } } - let page_index = if state.pages[state.last_page].values.contains(&value_index) { - state.last_page - } else { - state + let page_index = match state.last_page { + Some(last) if state.pages[last].values.contains(&value_index) => last, + _ => state .pages - .partition_point(|page| page.values.end <= value_index) + .partition_point(|page| page.values.end <= value_index), }; - state.last_page = page_index; + state.last_page = Some(page_index); let page = state .pages .get_mut(page_index) - .ok_or_else(|| vortex_err!("Missing PCO page for value {value_index}"))?; + .vortex_expect("PCO pages cover every valid value index"); let values = match &mut page.decoded { Some(decoded) => decoded, slot @ None => { From 96bb45002ff9d6ba724c60d5cef681d07dada933 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Tue, 22 Sep 2026 10:31:06 +0000 Subject: [PATCH 7/7] perf(pco): decode probe pages into an uninitialized buffer The page decompressor writes every element of the page, so zeroing the buffer first is a wasted pass over it. Reserve the capacity and set the length instead, as the range decompressor already does. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01THeokTcUrxdbZcPrDy8x6D --- encodings/pco/src/probe.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/encodings/pco/src/probe.rs b/encodings/pco/src/probe.rs index 7dee35000a4..2012d8abf51 100644 --- a/encodings/pco/src/probe.rs +++ b/encodings/pco/src/probe.rs @@ -152,7 +152,10 @@ fn decode_page( let mut decoder = chunk .page_decompressor(buffer, n_values) .map_err(vortex_err_from_pco)?; - let mut values = BufferMut::::zeroed_in(n_values, ctx.allocator().clone()); + let mut values = BufferMut::::with_capacity_in(n_values, ctx.allocator().clone()); + // SAFETY: the buffer reserves `n_values` elements, and the page decompressor was built for + // exactly that count, so `read` writes every element before anything observes it. + unsafe { values.set_len(n_values) }; decoder.read(&mut values).map_err(vortex_err_from_pco)?; Ok(PrimitiveArray::new(values.freeze(), Validity::NonNullable)) }