Skip to main content

vortex_fsst/compute/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright the Vortex contributors
3
4mod byte_length;
5mod cast;
6mod compare;
7mod filter;
8mod like;
9
10use vortex_array::ArrayRef;
11use vortex_array::ArrayView;
12use vortex_array::ExecutionCtx;
13use vortex_array::IntoArray;
14use vortex_array::arrays::dict::TakeExecute;
15use vortex_array::arrays::varbin::take_varbin;
16use vortex_array::builtins::ArrayBuiltins;
17use vortex_array::scalar::Scalar;
18use vortex_error::VortexResult;
19
20use crate::FSST;
21use crate::FSSTArrayExt;
22use crate::FSSTArraySlotsExt;
23
24impl TakeExecute for FSST {
25    fn take(
26        array: ArrayView<'_, Self>,
27        indices: &ArrayRef,
28        ctx: &mut ExecutionCtx,
29    ) -> VortexResult<Option<ArrayRef>> {
30        Ok(Some(
31            FSST::try_new_with_symbol_table(
32                array
33                    .dtype()
34                    .clone()
35                    .union_nullability(indices.dtype().nullability()),
36                array.symbol_table(),
37                take_varbin(array.codes().as_view(), indices, ctx)?,
38                array
39                    .uncompressed_lengths()
40                    .take(indices.clone())?
41                    .fill_null(Scalar::zero_value(
42                        &array.uncompressed_lengths_dtype().clone(),
43                    ))?,
44                ctx,
45            )?
46            .into_array(),
47        ))
48    }
49}
50
51#[cfg(test)]
52mod tests {
53    use rstest::rstest;
54    use vortex_array::ExecutionCtx;
55    use vortex_array::IntoArray;
56    use vortex_array::VortexSessionExecute;
57    use vortex_array::array_session;
58    use vortex_array::arrays::PrimitiveArray;
59    use vortex_array::arrays::VarBinArray;
60    use vortex_array::compute::conformance::consistency::test_array_consistency;
61    use vortex_array::compute::conformance::take::test_take_conformance;
62    use vortex_array::dtype::DType;
63    use vortex_array::dtype::Nullability;
64    use vortex_error::VortexResult;
65
66    use crate::FSSTArray;
67    use crate::fsst_compress;
68    use crate::fsst_train_compressor;
69
70    #[test]
71    fn test_take_null() -> VortexResult<()> {
72        let mut ctx = array_session().create_execution_ctx();
73        let arr =
74            VarBinArray::from_iter([Some("h")], DType::Utf8(Nullability::NonNullable)).into_array();
75        let compr = fsst_train_compressor(&arr, &mut ctx)?;
76        let fsst = fsst_compress(&arr, &compr, &mut ctx)?;
77
78        let idx1: PrimitiveArray = (0..1).collect();
79
80        assert_eq!(
81            fsst.take(idx1.into_array())?.dtype(),
82            &DType::Utf8(Nullability::NonNullable)
83        );
84
85        let idx2: PrimitiveArray = PrimitiveArray::from_option_iter(vec![Some(0)]);
86
87        assert_eq!(
88            fsst.take(idx2.into_array())?.dtype(),
89            &DType::Utf8(Nullability::Nullable)
90        );
91        Ok(())
92    }
93
94    #[rstest]
95    #[case(VarBinArray::from_iter(
96        ["hello world", "testing fsst", "compression test", "data array", "vortex encoding"].map(Some),
97        DType::Utf8(Nullability::NonNullable),
98    ))]
99    #[case(VarBinArray::from_iter(
100        [Some("hello"), None, Some("world"), Some("test"), None],
101        DType::Utf8(Nullability::Nullable),
102    ))]
103    #[case(VarBinArray::from_iter(
104        ["single element"].map(Some),
105        DType::Utf8(Nullability::NonNullable),
106    ))]
107    fn test_take_fsst_conformance(#[case] varbin: VarBinArray) -> VortexResult<()> {
108        let mut ctx = array_session().create_execution_ctx();
109        let varbin = varbin.into_array();
110        let compressor = fsst_train_compressor(&varbin, &mut ctx)?;
111        let array = fsst_compress(&varbin, &compressor, &mut ctx)?;
112        test_take_conformance(&array.into_array(), &mut ctx);
113        Ok(())
114    }
115
116    type FsstBuilder = fn(&mut ExecutionCtx) -> FSSTArray;
117
118    #[rstest]
119    // Basic string arrays
120    #[case::fsst_simple(|ctx: &mut ExecutionCtx| {
121        let array = VarBinArray::from_iter(
122            ["hello world", "testing fsst", "compression test", "data array", "vortex encoding"].map(Some),
123            DType::Utf8(Nullability::NonNullable),
124        ).into_array();
125        let compressor = fsst_train_compressor(&array, ctx).unwrap();
126        fsst_compress(&array, &compressor, ctx).unwrap()
127    })]
128    // Nullable strings
129    #[case::fsst_nullable(|ctx: &mut ExecutionCtx| {
130        let array = VarBinArray::from_iter(
131            [Some("hello"), None, Some("world"), Some("test"), None],
132            DType::Utf8(Nullability::Nullable),
133        ).into_array();
134        let compressor = fsst_train_compressor(&array, ctx).unwrap();
135        fsst_compress(&array, &compressor, ctx).unwrap()
136    })]
137    // Repetitive patterns (good for FSST compression)
138    #[case::fsst_repetitive(|ctx: &mut ExecutionCtx| {
139        let array = VarBinArray::from_iter(
140            ["http://example.com", "http://test.com", "http://vortex.dev", "http://data.org"].map(Some),
141            DType::Utf8(Nullability::NonNullable),
142        ).into_array();
143        let compressor = fsst_train_compressor(&array, ctx).unwrap();
144        fsst_compress(&array, &compressor, ctx).unwrap()
145    })]
146    // Edge cases
147    #[case::fsst_single(|ctx: &mut ExecutionCtx| {
148        let array = VarBinArray::from_iter(
149            ["single element"].map(Some),
150            DType::Utf8(Nullability::NonNullable),
151        ).into_array();
152        let compressor = fsst_train_compressor(&array, ctx).unwrap();
153        fsst_compress(&array, &compressor, ctx).unwrap()
154    })]
155    #[case::fsst_empty_strings(|ctx: &mut ExecutionCtx| {
156        let array = VarBinArray::from_iter(
157            ["", "test", "", "hello", ""].map(Some),
158            DType::Utf8(Nullability::NonNullable),
159        ).into_array();
160        let compressor = fsst_train_compressor(&array, ctx).unwrap();
161        fsst_compress(&array, &compressor, ctx).unwrap()
162    })]
163    // Large arrays
164    #[case::fsst_large(|ctx: &mut ExecutionCtx| {
165        let data: Vec<Option<&str>> = (0..1500)
166            .map(|i| Some(match i % 10 {
167                0 => "https://www.example.com/page",
168                1 => "https://www.test.org/data",
169                2 => "https://www.vortex.dev/docs",
170                3 => "https://www.github.com/apache/arrow",
171                4 => "https://www.rust-lang.org/learn",
172                5 => "SELECT * FROM table WHERE id = ",
173                6 => "INSERT INTO users (name, email) VALUES",
174                7 => "UPDATE records SET status = 'active'",
175                8 => "DELETE FROM logs WHERE timestamp < ",
176                _ => "CREATE TABLE data (id INT, value TEXT)",
177            }))
178            .collect();
179        let array = VarBinArray::from_iter(data, DType::Utf8(Nullability::NonNullable)).into_array();
180        let compressor = fsst_train_compressor(&array, ctx).unwrap();
181        fsst_compress(&array, &compressor, ctx).unwrap()
182    })]
183
184    fn test_fsst_consistency(#[case] build: FsstBuilder) {
185        let mut ctx = array_session().create_execution_ctx();
186        let array = build(&mut ctx);
187        test_array_consistency(
188            &array.into_array(),
189            &mut array_session().create_execution_ctx(),
190        );
191    }
192}