vortex_btrblocks/schemes/string/
fsst.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::VarBin;
14use vortex_array::arrays::VarBinArray;
15use vortex_array::arrays::primitive::PrimitiveArrayExt;
16use vortex_array::arrays::varbin::VarBinArrayExt;
17use vortex_compressor::scheme::CompressionEstimate;
18use vortex_compressor::scheme::DeferredEstimate;
19use vortex_error::VortexResult;
20use vortex_fsst::FSST;
21use vortex_fsst::FSSTArrayExt;
22use vortex_fsst::fsst_compress;
23use vortex_fsst::fsst_train_compressor;
24
25use crate::ArrayAndStats;
26use crate::CascadingCompressor;
27use crate::CompressorContext;
28use crate::Scheme;
29use crate::SchemeExt;
30
31#[derive(Debug, Copy, Clone, PartialEq, Eq)]
38pub struct FSSTScheme;
39
40impl Scheme for FSSTScheme {
41 fn scheme_name(&self) -> &'static str {
42 "vortex.string.fsst"
43 }
44
45 fn matches(&self, canonical: &Canonical) -> bool {
46 canonical.dtype().is_utf8()
47 }
48
49 fn produced_encodings(&self) -> Vec<ArrayId> {
50 vec![FSST.id(), VarBin.id()]
51 }
52
53 fn num_children(&self) -> usize {
55 2
56 }
57
58 fn expected_compression_ratio(
59 &self,
60 _data: &ArrayAndStats,
61 _compress_ctx: CompressorContext,
62 _exec_ctx: &mut ExecutionCtx,
63 ) -> CompressionEstimate {
64 CompressionEstimate::Deferred(DeferredEstimate::Sample)
65 }
66
67 fn compress(
68 &self,
69 compressor: &CascadingCompressor,
70 data: &ArrayAndStats,
71 compress_ctx: CompressorContext,
72 exec_ctx: &mut ExecutionCtx,
73 ) -> VortexResult<ArrayRef> {
74 let utf8 = data.array_as_varbinview().into_owned().into_array();
75 let compressor_fsst = fsst_train_compressor(&utf8, exec_ctx)?;
76 let fsst = fsst_compress(&utf8, &compressor_fsst, exec_ctx)?;
77
78 let uncompressed_lengths_primitive = fsst
79 .uncompressed_lengths()
80 .clone()
81 .execute::<PrimitiveArray>(exec_ctx)?
82 .narrow(exec_ctx)?;
83 let compressed_original_lengths = compressor.compress_child(
84 &uncompressed_lengths_primitive.into_array(),
85 &compress_ctx,
86 self.id(),
87 0,
88 exec_ctx,
89 )?;
90
91 let codes_offsets_primitive = fsst
92 .codes()
93 .offsets()
94 .clone()
95 .execute::<PrimitiveArray>(exec_ctx)?
96 .narrow(exec_ctx)?;
97 let compressed_codes_offsets = compressor.compress_child(
98 &codes_offsets_primitive.into_array(),
99 &compress_ctx,
100 self.id(),
101 1,
102 exec_ctx,
103 )?;
104 let compressed_codes = VarBinArray::try_new(
105 compressed_codes_offsets,
106 fsst.codes().bytes().clone(),
107 fsst.codes().dtype().clone(),
108 fsst.codes().validity()?,
109 )?;
110
111 let fsst = FSST::try_new(
112 fsst.dtype().clone(),
113 fsst.symbols().clone(),
114 fsst.symbol_lengths().clone(),
115 compressed_codes,
116 compressed_original_lengths,
117 exec_ctx,
118 )?;
119
120 Ok(fsst.into_array())
121 }
122}