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;
19use vortex_compressor::scheme::SchemeId;
20use vortex_error::VortexResult;
21use vortex_fastlanes::Delta;
22use vortex_fastlanes::RLE;
23use vortex_fastlanes::RLEArrayExt;
24use vortex_fastlanes::RLEArraySlotsExt;
25
26use super::DeltaScheme;
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    let compressed_indices = {
63        let rle_indices_primitive = rle_array
64            .indices()
65            .clone()
66            .execute::<PrimitiveArray>(exec_ctx)?
67            .narrow(exec_ctx)?;
68        let rle_indices = rle_indices_primitive.into_array();
69        if compressor.has_scheme(DeltaScheme::default().id()) {
70            try_compress_delta(
71                compressor,
72                &rle_indices,
73                &compress_ctx,
74                scheme.id(),
75                1,
76                exec_ctx,
77            )?
78        } else {
79            compressor.compress_child(&rle_indices, &compress_ctx, scheme.id(), 1, exec_ctx)?
80        }
81    };
82
83    let rle_offsets_primitive = rle_array
84        .values_idx_offsets()
85        .clone()
86        .execute::<PrimitiveArray>(exec_ctx)?
87        .narrow(exec_ctx)?;
88    let compressed_offsets = compressor.compress_child(
89        &rle_offsets_primitive.into_array(),
90        &compress_ctx,
91        scheme.id(),
92        2,
93        exec_ctx,
94    )?;
95
96    // SAFETY: Recursive compression doesn't affect the invariants.
97    unsafe {
98        Ok(RLE::new_unchecked(
99            compressed_values,
100            compressed_indices,
101            compressed_offsets,
102            rle_array.offset(),
103            rle_array.len(),
104        )
105        .into_array())
106    }
107}
108
109pub(crate) fn try_compress_delta(
110    compressor: &CascadingCompressor,
111    child: &ArrayRef,
112    parent_ctx: &CompressorContext,
113    parent_id: SchemeId,
114    child_index: usize,
115    exec_ctx: &mut ExecutionCtx,
116) -> VortexResult<ArrayRef> {
117    let child_primitive = child.clone().execute::<PrimitiveArray>(exec_ctx)?;
118    let (bases, deltas) = vortex_fastlanes::delta_compress(&child_primitive, exec_ctx)?;
119
120    let compressed_bases = compressor.compress_child(
121        &bases.into_array(),
122        parent_ctx,
123        parent_id,
124        child_index,
125        exec_ctx,
126    )?;
127    let compressed_deltas = compressor.compress_child(
128        &deltas.into_array(),
129        parent_ctx,
130        parent_id,
131        child_index,
132        exec_ctx,
133    )?;
134
135    Delta::try_new(compressed_bases, compressed_deltas, 0, child.len()).map(IntoArray::into_array)
136}
137
138impl Scheme for IntRLEScheme {
139    fn scheme_name(&self) -> &'static str {
140        "vortex.int.rle"
141    }
142
143    fn matches(&self, canonical: &Canonical) -> bool {
144        canonical.dtype().is_int()
145    }
146
147    fn produced_encodings(&self) -> Vec<ArrayId> {
148        vec![RLE.id()]
149    }
150
151    /// Children: values=0, indices=1, offsets=2.
152    fn num_children(&self) -> usize {
153        3
154    }
155
156    fn descendant_exclusions(&self) -> Vec<DescendantExclusion> {
157        rle_descendant_exclusions()
158    }
159
160    fn ancestor_exclusions(&self) -> Vec<AncestorExclusion> {
161        rle_ancestor_exclusions()
162    }
163
164    fn expected_compression_ratio(
165        &self,
166        data: &ArrayAndStats,
167        compress_ctx: CompressorContext,
168        exec_ctx: &mut ExecutionCtx,
169    ) -> CompressionEstimate {
170        // RLE is only useful when we cascade it with another encoding.
171        if compress_ctx.finished_cascading() {
172            return CompressionEstimate::Verdict(EstimateVerdict::Skip);
173        }
174        if data.integer_stats(exec_ctx).average_run_length() < RUN_LENGTH_THRESHOLD {
175            return CompressionEstimate::Verdict(EstimateVerdict::Skip);
176        }
177
178        CompressionEstimate::Deferred(DeferredEstimate::Sample)
179    }
180
181    fn compress(
182        &self,
183        compressor: &CascadingCompressor,
184        data: &ArrayAndStats,
185        compress_ctx: CompressorContext,
186        exec_ctx: &mut ExecutionCtx,
187    ) -> VortexResult<ArrayRef> {
188        rle_compress(self, compressor, data, compress_ctx, exec_ctx)
189    }
190}