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;
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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub struct IntRLEScheme;
40
41pub(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 #[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 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 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 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}