Skip to content
Merged
Show file tree
Hide file tree
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
25 changes: 21 additions & 4 deletions encodings/alp/src/alp/ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,32 +3,41 @@

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 vortex_error::vortex_err;

use crate::ALP;
use crate::ALPArrayExt;
use crate::ALPArraySlotsExt;
use crate::ALPFloat;
use crate::ALPSlots;
use crate::match_each_alp_float_ptype;

impl OperationsVTable<ALP> for ALP {
type ProbeState = ();

fn scalar_at(
array: ArrayView<'_, ALP>,
fn probe_scalar(
state: &mut ProbeState<'_, ALP>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
if !state.is_valid(index, ctx)? {
return Ok(Scalar::null(state.array().dtype().clone()));
}
let array = state.array();
if let Some(patches) = array.patches()
&& let Some(patch) = patches.get_patched(index)?
{
return patch.cast(array.dtype());
}

let encoded_val = array.encoded().execute_scalar(index, ctx)?;
let encoded_val = state
.slot(ALPSlots::ENCODED)?
.ok_or_else(|| vortex_err!("ALP encoded slot is missing"))?
.execute_scalar(index, ctx)?;

Ok(match_each_alp_float_ptype!(array.dtype().as_ptype(), |T| {
let encoded_val: <T as ALPFloat>::ALPInt =
Expand All @@ -39,4 +48,12 @@ impl OperationsVTable<ALP> for ALP {
)
}))
}

fn scalar_at(
array: ArrayView<'_, ALP>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
Self::probe_scalar(&mut ProbeState::once(array), index, ctx)
}
}
64 changes: 42 additions & 22 deletions encodings/datetime-parts/src/ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,27 +3,30 @@

use vortex_array::ArrayView;
use vortex_array::ExecutionCtx;
use vortex_array::ProbeState;
use vortex_array::dtype::DType;
use vortex_array::extension::datetime::Timestamp;
use vortex_array::scalar::Scalar;
use vortex_array::vtable::OperationsVTable;
use vortex_error::VortexExpect;
use vortex_error::VortexResult;
use vortex_error::vortex_err;
use vortex_error::vortex_panic;

use crate::DateTimeParts;
use crate::array::DateTimePartsArraySlotsExt;
use crate::DateTimePartsSlots;
use crate::timestamp;
use crate::timestamp::TimestampParts;

impl OperationsVTable<DateTimeParts> for DateTimeParts {
type ProbeState = ();

fn scalar_at(
array: ArrayView<'_, DateTimeParts>,
fn probe_scalar(
state: &mut ProbeState<'_, DateTimeParts>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
let array = state.array();
let DType::Extension(ext) = array.dtype().clone() else {
vortex_panic!(
"DateTimePartsArray must have extension dtype, found {}",
Expand All @@ -35,28 +38,19 @@ impl OperationsVTable<DateTimeParts> for DateTimeParts {
vortex_panic!(Compute: "must decode TemporalMetadata from extension metadata");
};

if !array.as_ref().is_valid(index, ctx)? {
if !state.is_valid(index, ctx)? {
return Ok(Scalar::null(DType::Extension(ext)));
}

let days: i32 = array
.days()
.execute_scalar(index, ctx)?
.as_primitive()
.as_::<i32>()
.vortex_expect("days fits in i32");
let seconds: i32 = array
.seconds()
.execute_scalar(index, ctx)?
.as_primitive()
.as_::<i32>()
.vortex_expect("seconds fits in i32");
let subseconds: i32 = array
.subseconds()
.execute_scalar(index, ctx)?
.as_primitive()
.as_::<i32>()
.vortex_expect("subseconds fits in i32");
let days = part_at(state, DateTimePartsSlots::DAYS, "days", index, ctx)?;
let seconds = part_at(state, DateTimePartsSlots::SECONDS, "seconds", index, ctx)?;
let subseconds = part_at(
state,
DateTimePartsSlots::SUBSECONDS,
"subseconds",
index,
ctx,
)?;

let ts = timestamp::combine(
TimestampParts {
Expand All @@ -72,4 +66,30 @@ impl OperationsVTable<DateTimeParts> for DateTimeParts {
Scalar::primitive(ts, ext.storage_dtype().nullability()),
))
}

fn scalar_at(
array: ArrayView<'_, DateTimeParts>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
Self::probe_scalar(&mut ProbeState::once(array), index, ctx)
}
}

/// Reads one timestamp part out of `slot`, through the probe so a repeated read keeps the
/// child's preparation.
fn part_at(
state: &mut ProbeState<'_, DateTimeParts>,
slot: usize,
name: &'static str,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<i32> {
Ok(state
.slot(slot)?
.ok_or_else(|| vortex_err!("DateTimeParts {name} slot is missing"))?
.execute_scalar(index, ctx)?
.as_primitive()
.as_::<i32>()
.vortex_expect("timestamp part fits in i32"))
}
17 changes: 17 additions & 0 deletions encodings/fastlanes/src/for_/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,23 @@ fn drifting_u32(len: u32) -> PrimitiveArray {
PrimitiveArray::from_iter((0..len).map(|i| (i / 1024) * 1_000_000 + i % 100))
}

#[test]
fn repeated_probe_across_sliced_chunks() -> VortexResult<()> {
let mut ctx = SESSION.create_execution_ctx();
let values = drifting_u32(3000);
let encoded = FoR::encode_chunked(values.clone(), &mut ctx)?;
let values = values.into_array();
let sliced = encoded.into_array().slice(1000..2100)?;
let mut probe = sliced.repeated_probe();
for index in [1048, 24, 0, 1099, 23, 1047, 24] {
assert_eq!(
probe.execute_scalar(index, &mut ctx)?,
values.execute_scalar(index + 1000, &mut ctx)?
);
}
Ok(())
}

#[rstest]
#[case::empty(PrimitiveArray::from_iter(Vec::<u32>::new()))]
#[case::one(PrimitiveArray::from_iter([7u32]))]
Expand Down
27 changes: 22 additions & 5 deletions encodings/fastlanes/src/for_/vtable/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,28 +3,37 @@

use vortex_array::ArrayView;
use vortex_array::ExecutionCtx;
use vortex_array::ProbeState;
use vortex_array::match_each_integer_ptype;
use vortex_array::scalar::Scalar;
use vortex_array::vtable::OperationsVTable;
use vortex_error::VortexExpect;
use vortex_error::VortexResult;
use vortex_error::vortex_err;

use super::FoR;
use crate::FL_CHUNK_SIZE;
use crate::for_::array::FoRArrayExt;
use crate::for_::array::FoRArraySlotsExt;
use crate::for_::array::FoRSlots;
impl OperationsVTable<FoR> for FoR {
type ProbeState = ();

fn scalar_at(
array: ArrayView<'_, FoR>,
fn probe_scalar(
state: &mut ProbeState<'_, FoR>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
let encoded_pvalue = array.encoded().execute_scalar(index, ctx)?;
let array = state.array();
let encoded_pvalue = state
.slot(FoRSlots::ENCODED)?
.ok_or_else(|| vortex_err!("FoR encoded slot is missing"))?
.execute_scalar(index, ctx)?;
let encoded_pvalue = encoded_pvalue.as_primitive();
let chunk = (usize::from(array.offset()) + index) / FL_CHUNK_SIZE;
let reference = array.references().execute_scalar(chunk, ctx)?;
let reference = state
.slot(FoRSlots::REFERENCES)?
.ok_or_else(|| vortex_err!("FoR references slot is missing"))?
.execute_scalar(chunk, ctx)?;
let reference = reference.as_primitive();

Ok(match_each_integer_ptype!(array.ptype(), |P| {
Expand All @@ -41,6 +50,14 @@ impl OperationsVTable<FoR> for FoR {
.unwrap_or_else(|| Scalar::null(array.dtype().clone()))
}))
}

fn scalar_at(
array: ArrayView<'_, FoR>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
Self::probe_scalar(&mut ProbeState::once(array), index, ctx)
}
}

#[cfg(test)]
Expand Down
36 changes: 32 additions & 4 deletions encodings/fsst/src/ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,9 @@

use vortex_array::ArrayView;
use vortex_array::ExecutionCtx;
use vortex_array::IntoArray;
use vortex_array::ProbeState;
use vortex_array::RepeatedArrayProbe;
use vortex_array::arrays::varbin::varbin_scalar;
use vortex_array::scalar::Scalar;
use vortex_array::vtable::OperationsVTable;
Expand All @@ -13,18 +16,43 @@ use vortex_error::VortexResult;
use crate::FSST;
use crate::FSSTArrayExt;

/// The codes array is rebuilt from the codes buffer and the offsets slot on every read, so a
/// repeated probe keeps one probe over it and the offsets child keeps its preparation.
#[derive(Default)]
pub struct FsstProbeState {
codes: Option<RepeatedArrayProbe>,
}

impl OperationsVTable<FSST> for FSST {
type ProbeState = ();
type ProbeState = FsstProbeState;

fn scalar_at(
array: ArrayView<'_, FSST>,
fn probe_scalar(
state: &mut ProbeState<'_, FSST>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
let compressed = array.codes().execute_scalar(index, ctx)?;
if !state.is_valid(index, ctx)? {
return Ok(Scalar::null(state.array().dtype().clone()));
}
let array = state.array();
let compressed = match state.retained() {
None => array.codes().into_array().execute_scalar(index, ctx)?,
Some(retained) => retained
.codes
.get_or_insert_with(|| RepeatedArrayProbe::new(array.codes().into_array()))
.execute_scalar(index, ctx)?,
};
let binary_datum = compressed.as_binary().value().vortex_expect("non-null");

let decoded_buffer = ByteBuffer::from(array.decompressor().decompress(binary_datum));
Ok(varbin_scalar(decoded_buffer, array.dtype()))
}

fn scalar_at(
array: ArrayView<'_, FSST>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
Self::probe_scalar(&mut ProbeState::once(array), index, ctx)
}
}
20 changes: 17 additions & 3 deletions encodings/zigzag/src/array.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,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::buffer::BufferHandle;
Expand All @@ -33,6 +34,7 @@ use vortex_error::VortexExpect;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
use vortex_error::vortex_ensure;
use vortex_error::vortex_err;
use vortex_error::vortex_panic;
use vortex_session::VortexSession;
use vortex_session::registry::CachedId;
Expand Down Expand Up @@ -233,12 +235,16 @@ impl Default for ZigZagData {
impl OperationsVTable<ZigZag> for ZigZag {
type ProbeState = ();

fn scalar_at(
array: ArrayView<'_, ZigZag>,
fn probe_scalar(
state: &mut ProbeState<'_, ZigZag>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
let scalar = array.encoded().execute_scalar(index, ctx)?;
let array = state.array();
let scalar = state
.slot(ZigZagSlots::ENCODED)?
.ok_or_else(|| vortex_err!("ZigZag encoded slot is missing"))?
.execute_scalar(index, ctx)?;
if scalar.is_null() {
return scalar.primitive_reinterpret_cast(ZigZagArrayExt::ptype(&array));
}
Expand All @@ -255,6 +261,14 @@ impl OperationsVTable<ZigZag> for ZigZag {
)
}))
}

fn scalar_at(
array: ArrayView<'_, ZigZag>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
Self::probe_scalar(&mut ProbeState::once(array), index, ctx)
}
}

impl ValidityChild<ZigZag> for ZigZag {
Expand Down
19 changes: 17 additions & 2 deletions vortex-array/src/arrays/chunked/vtable/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,24 +2,39 @@
// SPDX-FileCopyrightText: Copyright the Vortex contributors

use vortex_error::VortexResult;
use vortex_error::vortex_err;

use crate::ExecutionCtx;
use crate::array::ArrayView;
use crate::array::OperationsVTable;
use crate::array::ProbeState;
use crate::arrays::Chunked;
use crate::arrays::chunked::ChunkedArrayExt;
use crate::arrays::chunked::ChunkedSlots;
use crate::scalar::Scalar;

impl OperationsVTable<Chunked> for Chunked {
type ProbeState = ();

fn probe_scalar(
state: &mut ProbeState<'_, Chunked>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
let (chunk_index, chunk_offset) = state.array().find_chunk_idx(index)?;
let slot = ChunkedSlots::CHUNKS_OFFSET + chunk_index;
state
.slot(slot)?
.ok_or_else(|| vortex_err!("Chunked chunk slot {slot} is missing"))?
.execute_scalar(chunk_offset, ctx)
}

fn scalar_at(
array: ArrayView<'_, Chunked>,
index: usize,
ctx: &mut ExecutionCtx,
) -> VortexResult<Scalar> {
let (chunk_index, chunk_offset) = array.find_chunk_idx(index)?;
array.chunk(chunk_index).execute_scalar(chunk_offset, ctx)
Self::probe_scalar(&mut ProbeState::once(array), index, ctx)
}
}

Expand Down
Loading
Loading