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;
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#[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 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 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 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 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}