Skip to main content

vortex_fastlanes/delta/vtable/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright the Vortex contributors
3
4use std::hash::Hash;
5use std::hash::Hasher;
6
7use prost::Message;
8use vortex_array::Array;
9use vortex_array::ArrayEq;
10use vortex_array::ArrayHash;
11use vortex_array::ArrayId;
12use vortex_array::ArrayParts;
13use vortex_array::ArrayRef;
14use vortex_array::ArrayView;
15use vortex_array::EqMode;
16use vortex_array::ExecutionCtx;
17use vortex_array::ExecutionResult;
18use vortex_array::IntoArray;
19use vortex_array::arrays::PrimitiveArray;
20use vortex_array::buffer::BufferHandle;
21use vortex_array::dtype::DType;
22use vortex_array::dtype::PType;
23use vortex_array::serde::ArrayChildren;
24use vortex_array::vtable::VTable;
25use vortex_error::VortexResult;
26use vortex_error::vortex_ensure;
27use vortex_error::vortex_err;
28use vortex_error::vortex_panic;
29use vortex_session::VortexSession;
30use vortex_session::registry::CachedId;
31
32use crate::DeltaData;
33use crate::delta::array::DeltaArrayExt;
34use crate::delta::array::DeltaArraySlotsExt;
35use crate::delta::array::DeltaSlots;
36use crate::delta::array::DeltaSlotsView;
37use crate::delta::array::delta_decompress::delta_decompress;
38use crate::delta::array::lane_count;
39use crate::delta_compress;
40
41mod operations;
42mod rules;
43mod slice;
44mod validity;
45
46/// A [`Delta`]-encoded Vortex array.
47pub type DeltaArray = Array<Delta>;
48
49#[derive(Clone, prost::Message)]
50#[repr(C)]
51pub struct DeltaMetadata {
52    #[prost(uint64, tag = "1")]
53    deltas_len: u64,
54    #[prost(uint32, tag = "2")]
55    offset: u32, // must be <1024
56}
57
58impl ArrayHash for DeltaData {
59    fn array_hash<H: Hasher>(&self, state: &mut H, _accuracy: EqMode) {
60        self.offset.hash(state);
61    }
62}
63
64impl ArrayEq for DeltaData {
65    fn array_eq(&self, other: &Self, _accuracy: EqMode) -> bool {
66        self.offset == other.offset
67    }
68}
69
70impl VTable for Delta {
71    type TypedArrayData = DeltaData;
72
73    type OperationsVTable = Self;
74    type ValidityVTable = Self;
75
76    fn id(&self) -> ArrayId {
77        static ID: CachedId = CachedId::new("fastlanes.delta");
78        *ID
79    }
80
81    fn validate(
82        &self,
83        data: &Self::TypedArrayData,
84        dtype: &DType,
85        len: usize,
86        slots: &[Option<ArrayRef>],
87    ) -> VortexResult<()> {
88        let delta_slots = DeltaSlotsView::from_slots(slots);
89        validate_parts(
90            delta_slots.bases,
91            delta_slots.deltas,
92            data.offset,
93            dtype,
94            len,
95        )
96    }
97
98    fn nbuffers(_array: ArrayView<'_, Self>) -> usize {
99        0
100    }
101
102    fn buffer(_array: ArrayView<'_, Self>, idx: usize) -> BufferHandle {
103        vortex_panic!("DeltaArray buffer index {idx} out of bounds")
104    }
105
106    fn buffer_name(_array: ArrayView<'_, Self>, _idx: usize) -> Option<String> {
107        None
108    }
109
110    fn with_buffers(
111        &self,
112        array: ArrayView<'_, Self>,
113        buffers: &[BufferHandle],
114    ) -> VortexResult<ArrayParts<Self>> {
115        vortex_array::vtable::with_empty_buffers(self, array, buffers)
116    }
117
118    fn reduce_parent(
119        array: ArrayView<'_, Self>,
120        parent: &ArrayRef,
121        child_idx: usize,
122    ) -> VortexResult<Option<ArrayRef>> {
123        rules::RULES.evaluate(array, parent, child_idx)
124    }
125
126    fn slot_name(_array: ArrayView<'_, Self>, idx: usize) -> String {
127        DeltaSlots::NAMES[idx].to_string()
128    }
129
130    fn serialize(
131        array: ArrayView<'_, Self>,
132        _session: &VortexSession,
133    ) -> VortexResult<Option<Vec<u8>>> {
134        Ok(Some(
135            DeltaMetadata {
136                deltas_len: array.deltas().len() as u64,
137                offset: array.offset() as u32,
138            }
139            .encode_to_vec(),
140        ))
141    }
142
143    fn deserialize(
144        &self,
145        dtype: &DType,
146        len: usize,
147        metadata: &[u8],
148        buffers: &[BufferHandle],
149        children: &dyn ArrayChildren,
150        _session: &VortexSession,
151    ) -> VortexResult<ArrayParts<Self>> {
152        vortex_ensure!(
153            buffers.is_empty(),
154            "DeltaArray expects 0 buffers, got {}",
155            buffers.len()
156        );
157        vortex_ensure!(
158            children.len() == 2,
159            "DeltaArray expects 2 children, got {}",
160            children.len()
161        );
162        let metadata = DeltaMetadata::decode(metadata)?;
163        let ptype = PType::try_from(dtype)?;
164        let lanes = lane_count(ptype);
165
166        // Compute the length of the bases array
167        let deltas_len = usize::try_from(metadata.deltas_len)
168            .map_err(|_| vortex_err!("deltas_len {} overflowed usize", metadata.deltas_len))?;
169        let num_chunks = deltas_len / 1024;
170        let remainder_base_size = if deltas_len % 1024 > 0 { 1 } else { 0 };
171        let bases_len = num_chunks * lanes + remainder_base_size;
172
173        let bases = children.get(0, dtype, bases_len)?;
174        let deltas = children.get(1, dtype, deltas_len)?;
175
176        let data = DeltaData::try_new(metadata.offset as usize)?;
177        let slots = DeltaSlots { bases, deltas }.into_slots();
178        Ok(ArrayParts::new(self.clone(), dtype.clone(), len, data).with_slots(slots))
179    }
180
181    fn execute(array: Array<Self>, ctx: &mut ExecutionCtx) -> VortexResult<ExecutionResult> {
182        Ok(ExecutionResult::done(
183            delta_decompress(&array, ctx)?.into_array(),
184        ))
185    }
186}
187
188#[derive(Clone, Debug)]
189pub struct Delta;
190
191impl Delta {
192    pub fn try_new(
193        bases: ArrayRef,
194        deltas: ArrayRef,
195        offset: usize,
196        len: usize,
197    ) -> VortexResult<DeltaArray> {
198        let dtype = bases.dtype().with_nullability(deltas.dtype().nullability());
199        let data = DeltaData::try_new(offset)?;
200        let slots = DeltaSlots { bases, deltas }.into_slots();
201        Array::try_from_parts(ArrayParts::new(Delta, dtype, len, data).with_slots(slots))
202    }
203
204    /// Compress a primitive array using Delta encoding.
205    pub fn try_from_primitive_array(
206        array: &PrimitiveArray,
207        ctx: &mut ExecutionCtx,
208    ) -> VortexResult<DeltaArray> {
209        let logical_len = array.len();
210        let (bases, deltas) = delta_compress(array, ctx)?;
211        Self::try_new(bases.into_array(), deltas.into_array(), 0, logical_len)
212    }
213}
214
215fn validate_parts(
216    bases: &ArrayRef,
217    deltas: &ArrayRef,
218    offset: usize,
219    dtype: &DType,
220    len: usize,
221) -> VortexResult<()> {
222    vortex_ensure!(
223        offset + len <= deltas.len(),
224        "offset + len, {offset} + {len}, must be less than or equal to the size of deltas: {}",
225        deltas.len()
226    );
227    vortex_ensure!(
228        bases.dtype().eq_ignore_nullability(deltas.dtype()),
229        "DeltaArray: bases and deltas must have the same dtype, got {} and {}",
230        bases.dtype(),
231        deltas.dtype()
232    );
233
234    vortex_ensure!(
235        bases.dtype().is_int(),
236        "DeltaArray: dtype must be an integer, got {}",
237        bases.dtype()
238    );
239
240    let expected_dtype = bases.dtype().with_nullability(deltas.dtype().nullability());
241    vortex_ensure!(
242        dtype == &expected_dtype,
243        "DeltaArray dtype mismatch: expected {expected_dtype}, got {dtype}"
244    );
245
246    let lanes = lane_count(bases.dtype().as_ptype());
247
248    vortex_ensure!(
249        deltas.len().is_multiple_of(1024),
250        "deltas length ({}) must be a multiple of 1024",
251        deltas.len(),
252    );
253    vortex_ensure!(
254        bases.len().is_multiple_of(lanes),
255        "bases length ({}) must be a multiple of LANES ({lanes})",
256        bases.len(),
257    );
258    Ok(())
259}
260
261#[cfg(test)]
262mod tests {
263    use prost::Message;
264    use vortex_array::test_harness::check_metadata;
265
266    use super::DeltaMetadata;
267
268    #[cfg_attr(miri, ignore)]
269    #[test]
270    fn test_delta_metadata() {
271        check_metadata(
272            "delta.metadata",
273            &DeltaMetadata {
274                offset: u32::MAX,
275                deltas_len: u64::MAX,
276            }
277            .encode_to_vec(),
278        );
279    }
280}