Skip to main content

vortex_btrblocks/schemes/integer/
rle.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright the Vortex contributors
3
4//! Run-length integer encoding and shared RLE compression helpers.
5
6use vortex_array::ArrayId;
7use vortex_array::ArrayRef;
8use vortex_array::Canonical;
9use vortex_array::ExecutionCtx;
10use vortex_array::IntoArray;
11use vortex_array::VTable;
12use vortex_array::arrays::PrimitiveArray;
13use vortex_array::arrays::primitive::PrimitiveArrayExt;
14use vortex_compressor::scheme::AncestorExclusion;
15use vortex_compressor::scheme::CompressionEstimate;
16use vortex_compressor::scheme::DeferredEstimate;
17use vortex_compressor::scheme::DescendantExclusion;
18use vortex_compressor::scheme::EstimateVerdict;
19#[cfg(feature = "unstable_encodings")]
20use vortex_compressor::scheme::SchemeId;
21use vortex_error::VortexResult;
22#[cfg(feature = "unstable_encodings")]
23use vortex_fastlanes::Delta;
24use vortex_fastlanes::RLE;
25use vortex_fastlanes::RLEArrayExt;
26use vortex_fastlanes::RLEArraySlotsExt;
27
28use super::RUN_LENGTH_THRESHOLD;
29use crate::ArrayAndStats;
30use crate::CascadingCompressor;
31use crate::CompressorContext;
32use crate::Scheme;
33use crate::SchemeExt;
34use crate::schemes::rle_ancestor_exclusions;
35use crate::schemes::rle_descendant_exclusions;
36
37/// RLE scheme for integer arrays.
38#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub struct IntRLEScheme;
40
41/// Shared compression logic for RLE schemes.
42pub(crate) fn rle_compress(
43    scheme: &dyn Scheme,
44    compressor: &CascadingCompressor,
45    data: &ArrayAndStats,
46    compress_ctx: CompressorContext,
47    exec_ctx: &mut ExecutionCtx,
48) -> VortexResult<ArrayRef> {
49    let rle_array = RLE::encode(data.array_as_primitive(), exec_ctx)?;
50
51    let rle_values_primitive = rle_array
52        .values()
53        .clone()
54        .execute::<PrimitiveArray>(exec_ctx)?;
55    let compressed_values = compressor.compress_child(
56        &rle_values_primitive.into_array(),
57        &compress_ctx,
58        scheme.id(),
59        0,
60        exec_ctx,
61    )?;
62
63    // Delta is an unstable encoding, once we deem it stable we can switch over to this always.
64    #[cfg(feature = "unstable_encodings")]
65    let compressed_indices = {
66        let rle_indices_primitive = rle_array
67            .indices()
68            .clone()
69            .execute::<PrimitiveArray>(exec_ctx)?
70            .narrow(exec_ctx)?;
71        try_compress_delta(
72            compressor,
73            &rle_indices_primitive.into_array(),
74            &compress_ctx,
75            scheme.id(),
76            1,
77            exec_ctx,
78        )?
79    };
80
81    #[cfg(not(feature = "unstable_encodings"))]
82    let compressed_indices = {
83        let rle_indices_primitive = rle_array
84            .indices()
85            .clone()
86            .execute::<PrimitiveArray>(exec_ctx)?
87            .narrow(exec_ctx)?;
88        compressor.compress_child(
89            &rle_indices_primitive.into_array(),
90            &compress_ctx,
91            scheme.id(),
92            1,
93            exec_ctx,
94        )?
95    };
96
97    let rle_offsets_primitive = rle_array
98        .values_idx_offsets()
99        .clone()
100        .execute::<PrimitiveArray>(exec_ctx)?
101        .narrow(exec_ctx)?;
102    let compressed_offsets = compressor.compress_child(
103        &rle_offsets_primitive.into_array(),
104        &compress_ctx,
105        scheme.id(),
106        2,
107        exec_ctx,
108    )?;
109
110    // SAFETY: Recursive compression doesn't affect the invariants.
111    unsafe {
112        Ok(RLE::new_unchecked(
113            compressed_values,
114            compressed_indices,
115            compressed_offsets,
116            rle_array.offset(),
117            rle_array.len(),
118        )
119        .into_array())
120    }
121}
122
123#[cfg(feature = "unstable_encodings")]
124pub(crate) fn try_compress_delta(
125    compressor: &CascadingCompressor,
126    child: &ArrayRef,
127    parent_ctx: &CompressorContext,
128    parent_id: SchemeId,
129    child_index: usize,
130    exec_ctx: &mut ExecutionCtx,
131) -> VortexResult<ArrayRef> {
132    let child_primitive = child.clone().execute::<PrimitiveArray>(exec_ctx)?;
133    let (bases, deltas) = vortex_fastlanes::delta_compress(&child_primitive, exec_ctx)?;
134
135    let compressed_bases = compressor.compress_child(
136        &bases.into_array(),
137        parent_ctx,
138        parent_id,
139        child_index,
140        exec_ctx,
141    )?;
142    let compressed_deltas = compressor.compress_child(
143        &deltas.into_array(),
144        parent_ctx,
145        parent_id,
146        child_index,
147        exec_ctx,
148    )?;
149
150    Delta::try_new(compressed_bases, compressed_deltas, 0, child.len()).map(IntoArray::into_array)
151}
152
153impl Scheme for IntRLEScheme {
154    fn scheme_name(&self) -> &'static str {
155        "vortex.int.rle"
156    }
157
158    fn matches(&self, canonical: &Canonical) -> bool {
159        canonical.dtype().is_int()
160    }
161
162    fn produced_encodings(&self) -> Vec<ArrayId> {
163        vec![RLE.id()]
164    }
165
166    /// Children: values=0, indices=1, offsets=2.
167    fn num_children(&self) -> usize {
168        3
169    }
170
171    fn descendant_exclusions(&self) -> Vec<DescendantExclusion> {
172        rle_descendant_exclusions()
173    }
174
175    fn ancestor_exclusions(&self) -> Vec<AncestorExclusion> {
176        rle_ancestor_exclusions()
177    }
178
179    fn expected_compression_ratio(
180        &self,
181        data: &ArrayAndStats,
182        compress_ctx: CompressorContext,
183        exec_ctx: &mut ExecutionCtx,
184    ) -> CompressionEstimate {
185        // RLE is only useful when we cascade it with another encoding.
186        if compress_ctx.finished_cascading() {
187            return CompressionEstimate::Verdict(EstimateVerdict::Skip);
188        }
189        if data.integer_stats(exec_ctx).average_run_length() < RUN_LENGTH_THRESHOLD {
190            return CompressionEstimate::Verdict(EstimateVerdict::Skip);
191        }
192
193        CompressionEstimate::Deferred(DeferredEstimate::Sample)
194    }
195
196    fn compress(
197        &self,
198        compressor: &CascadingCompressor,
199        data: &ArrayAndStats,
200        compress_ctx: CompressorContext,
201        exec_ctx: &mut ExecutionCtx,
202    ) -> VortexResult<ArrayRef> {
203        rle_compress(self, compressor, data, compress_ctx, exec_ctx)
204    }
205}