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