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