vortex_btrblocks/schemes/integer/
rle.rs1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
38pub struct IntRLEScheme;
39
40pub(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 #[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 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 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 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}