Skip to main content

vortex_btrblocks/schemes/
temporal.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright the Vortex contributors
3
4//! Temporal compression scheme using datetime-part decomposition.
5
6use 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::ExtensionArray;
13use vortex_array::arrays::PrimitiveArray;
14use vortex_array::arrays::TemporalArray;
15use vortex_array::arrays::extension::ExtensionArrayExt;
16use vortex_array::arrays::primitive::PrimitiveArrayExt;
17use vortex_array::dtype::extension::Matcher;
18use vortex_array::extension::datetime::AnyTemporal;
19use vortex_array::extension::datetime::TemporalMetadata;
20use vortex_compressor::scheme::CompressionEstimate;
21use vortex_compressor::scheme::EstimateVerdict;
22use vortex_datetime_parts::DateTimeParts;
23use vortex_datetime_parts::TemporalParts;
24use vortex_datetime_parts::split_temporal;
25use vortex_error::VortexResult;
26
27use crate::ArrayAndStats;
28use crate::CascadingCompressor;
29use crate::CompressorContext;
30use crate::Scheme;
31use crate::SchemeExt;
32
33/// Compression scheme for temporal timestamp arrays via datetime-part decomposition.
34///
35/// Splits timestamps into days, seconds, and subseconds components, compresses each
36/// independently, and wraps the result in a `DateTimePartsArray`.
37#[derive(Debug, Copy, Clone, PartialEq, Eq)]
38pub struct TemporalScheme;
39
40impl Scheme for TemporalScheme {
41    fn scheme_name(&self) -> &'static str {
42        "vortex.ext.temporal"
43    }
44
45    fn matches(&self, canonical: &Canonical) -> bool {
46        let Canonical::Extension(ext) = canonical else {
47            return false;
48        };
49
50        let ext_dtype = ext.ext_dtype();
51
52        matches!(
53            AnyTemporal::try_match(ext_dtype),
54            Some(TemporalMetadata::Timestamp(..))
55        )
56    }
57
58    fn produced_encodings(&self) -> Vec<ArrayId> {
59        vec![DateTimeParts.id()]
60    }
61
62    /// Children: days=0, seconds=1, subseconds=2.
63    fn num_children(&self) -> usize {
64        3
65    }
66
67    fn expected_compression_ratio(
68        &self,
69        _data: &ArrayAndStats,
70        _compress_ctx: CompressorContext,
71        _exec_ctx: &mut ExecutionCtx,
72    ) -> CompressionEstimate {
73        // Temporal compression (splitting into parts) is almost always beneficial.
74        CompressionEstimate::Verdict(EstimateVerdict::AlwaysUse)
75    }
76
77    fn compress(
78        &self,
79        compressor: &CascadingCompressor,
80        data: &ArrayAndStats,
81        compress_ctx: CompressorContext,
82        exec_ctx: &mut ExecutionCtx,
83    ) -> VortexResult<ArrayRef> {
84        let array = data.array().clone();
85        let ext_array = array.execute::<ExtensionArray>(exec_ctx)?;
86        let temporal_array = TemporalArray::try_from(ext_array.into_array())?;
87
88        let dtype = temporal_array.dtype().clone();
89        let TemporalParts {
90            days,
91            seconds,
92            subseconds,
93        } = split_temporal(temporal_array, exec_ctx)?;
94
95        let days_primitive = days.execute::<PrimitiveArray>(exec_ctx)?.narrow(exec_ctx)?;
96        let days = compressor.compress_child(
97            &days_primitive.into_array(),
98            &compress_ctx,
99            self.id(),
100            0,
101            exec_ctx,
102        )?;
103        let seconds_primitive = seconds
104            .execute::<PrimitiveArray>(exec_ctx)?
105            .narrow(exec_ctx)?;
106        let seconds = compressor.compress_child(
107            &seconds_primitive.into_array(),
108            &compress_ctx,
109            self.id(),
110            1,
111            exec_ctx,
112        )?;
113        let subseconds_primitive = subseconds
114            .execute::<PrimitiveArray>(exec_ctx)?
115            .narrow(exec_ctx)?;
116        let subseconds = compressor.compress_child(
117            &subseconds_primitive.into_array(),
118            &compress_ctx,
119            self.id(),
120            2,
121            exec_ctx,
122        )?;
123
124        Ok(DateTimeParts::try_new(dtype, days, seconds, subseconds)?.into_array())
125    }
126}