From 47ea45fd4cd5586fbe1ccac899a0b8249b1dbcf1 Mon Sep 17 00:00:00 2001 From: Nikhil Nayak Date: Fri, 14 Nov 2025 08:08:42 -0800 Subject: [PATCH 1/2] Implemented Multiple Date32 Extractions Per Query --- src/parquet/src/optimizers/lineage_opt.rs | 140 +++++++++++++++------- src/parquet/src/optimizers/mod.rs | 6 +- 2 files changed, 103 insertions(+), 43 deletions(-) diff --git a/src/parquet/src/optimizers/lineage_opt.rs b/src/parquet/src/optimizers/lineage_opt.rs index 344b9895..eba9504d 100644 --- a/src/parquet/src/optimizers/lineage_opt.rs +++ b/src/parquet/src/optimizers/lineage_opt.rs @@ -2,7 +2,7 @@ //! It then attaches the metadata to schema adapter, which is then passed to the physical plan. //! The physical optimizer will move the metadata to the fields of the schema. -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::str::FromStr; use std::sync::{Arc, Mutex, OnceLock}; @@ -46,7 +46,7 @@ impl SupportedIntervalUnit { #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) struct DateExtraction { pub(crate) column: Column, - pub(crate) component: SupportedIntervalUnit, + pub(crate) components: HashSet, } /// Metadata describing a Variant column that participates in a `variant_get`. @@ -59,10 +59,36 @@ pub(crate) struct VariantExtraction { /// Annotation that should be attached to a column in the file schema. #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum ColumnAnnotation { - DatePart(SupportedIntervalUnit), + DatePart(HashSet), VariantPath(String), } +impl ColumnAnnotation { + /// Serialize DatePart units to a comma-separated string. + /// Returns None if this is not a DatePart annotation. + pub(crate) fn serialize_date_part(&self) -> Option { + match self { + ColumnAnnotation::DatePart(units) => { + let mut sorted_units: Vec<&SupportedIntervalUnit> = units.iter().collect(); + // Sort by a consistent order: Year, Month, Day + sorted_units.sort_by_key(|unit| match unit { + SupportedIntervalUnit::Year => 0, + SupportedIntervalUnit::Month => 1, + SupportedIntervalUnit::Day => 2, + }); + Some( + sorted_units + .iter() + .map(|unit| unit.metadata_value()) + .collect::>() + .join(","), + ) + } + ColumnAnnotation::VariantPath(_) => None, + } + } +} + /// Logical optimizer that analyses the logical plan to detect columns that /// are only used via compatible `EXTRACT` or `variant_get` projections. #[derive(Debug, Default)] @@ -499,24 +525,40 @@ impl TableColumnUsage { let mut extractions = Vec::new(); for (key, stats) in self.usage.iter() { if matches!(stats.data_type, DataType::Date32) { - // Check if every usage's first operation is Extract with the same unit - let first_unit = stats.usages.first().and_then(|usage| { - if let Some(Operation::Extract(unit)) = usage.first() { - Some(unit) - } else { - None + // Collect all extract units from paths where the first n operations are all extracts + let mut all_units = HashSet::new(); + let mut all_paths_valid = true; + + for usage in &stats.usages { + // Collect all Extract units from the leading sequence of extracts + let mut path_units = HashSet::new(); + for op in usage { + match op { + Operation::Extract(unit) => { + path_units.insert(unit); + } + _ => { + // Stop at first non-extract operation + break; + } + } } - }); - if let Some(first_unit) = first_unit { - let all_matches = stats.usages.iter().all(|usage| { - matches!(usage.first(), Some(Operation::Extract(unit)) if unit == first_unit) - }); - if all_matches { - extractions.push(DateExtraction { - column: key.to_column(), - component: *first_unit, - }); + + if path_units.is_empty() { + // This path doesn't start with Extract, so skip this column + all_paths_valid = false; + break; } + + // Union the units from this path into the overall set + all_units.extend(path_units); + } + + if all_paths_valid && !all_units.is_empty() { + extractions.push(DateExtraction { + column: key.to_column(), + components: all_units, + }); } } } @@ -554,7 +596,7 @@ fn build_annotation_map( for extraction in date_findings { annotations.insert( ColumnKey::from_column(&extraction.column), - ColumnAnnotation::DatePart(extraction.component), + ColumnAnnotation::DatePart(extraction.components.clone()), ); } for extraction in variant_findings { @@ -1034,10 +1076,19 @@ mod tests { let expected_field_metadata = expected .iter() .map(|extraction| { - ( - extraction.column.name().to_string(), - extraction.component.metadata_value().to_string(), - ) + let mut sorted_units: Vec<&SupportedIntervalUnit> = + extraction.components.iter().collect(); + sorted_units.sort_by_key(|unit| match unit { + SupportedIntervalUnit::Year => 0, + SupportedIntervalUnit::Month => 1, + SupportedIntervalUnit::Day => 2, + }); + let metadata_value = sorted_units + .iter() + .map(|unit| unit.metadata_value()) + .collect::>() + .join(","); + (extraction.column.name().to_string(), metadata_value) }) .collect::>(); assert_eq!(field_metadata_map, expected_field_metadata); @@ -1048,20 +1099,14 @@ mod tests { general_test( "SELECT EXTRACT(YEAR FROM table_a.date) AS year, EXTRACT(DAY FROM table_b.date) AS day FROM table_a INNER JOIN table_b ON table_a.event_ts = table_b.event_ts", vec![ - DateExtraction { column: Column::new(Some("table_a"), "date"), component: SupportedIntervalUnit::Year }, - DateExtraction { column: Column::new(Some("table_b"), "date"), component: SupportedIntervalUnit::Day }, - ], - ) - .await; - } - - #[tokio::test] - async fn single_table_multiple_extracts() { - general_test( - "SELECT EXTRACT(YEAR FROM date_copy) AS year, EXTRACT(DAY FROM date) AS day FROM table_a", - vec![ - DateExtraction { column: Column::new(Some("table_a"), "date"), component: SupportedIntervalUnit::Day }, - DateExtraction { column: Column::new(Some("table_a"), "date_copy"), component: SupportedIntervalUnit::Year }, + DateExtraction { + column: Column::new(Some("table_a"), "date"), + components: HashSet::from([SupportedIntervalUnit::Year]), + }, + DateExtraction { + column: Column::new(Some("table_b"), "date"), + components: HashSet::from([SupportedIntervalUnit::Day]), + }, ], ) .await; @@ -1079,7 +1124,7 @@ mod tests { ]; let expected = vec![DateExtraction { column: Column::new(Some("table_a"), "date"), - component: SupportedIntervalUnit::Day, + components: HashSet::from([SupportedIntervalUnit::Day]), }]; for sql in statements { general_test(sql, expected.clone()).await; @@ -1125,11 +1170,9 @@ mod tests { #[tokio::test] async fn inconsistent_extracts_are_ignored() { let statements = vec![ - "SELECT EXTRACT(DAY FROM date) AS day, EXTRACT(MONTH FROM date) AS month FROM table_a", "SELECT EXTRACT(DAY FROM date + INTERVAL '1 day') AS day FROM table_a", "SELECT date FROM table_a", "SELECT EXTRACT(DAY FROM table_a.date) AS day FROM table_a INNER JOIN table_b ON table_a.date = table_b.date", - "SELECT (SELECT MAX(EXTRACT(DAY FROM date)) FROM table_a) AS max_day, (SELECT MIN(EXTRACT(Month FROM date)) FROM table_a) AS min_day", "SELECT EXTRACT(YEAR FROM event_ts) AS year FROM table_a", // todo: time stamp is not supported yet. ]; @@ -1138,6 +1181,21 @@ mod tests { } } + #[tokio::test] + async fn single_table_multiple_extracts() { + let statements = vec![ + "SELECT EXTRACT(DAY FROM date) AS day, EXTRACT(MONTH FROM date) AS month FROM table_a", + "SELECT (SELECT MAX(EXTRACT(DAY FROM date)) FROM table_a) AS max_day, (SELECT MIN(EXTRACT(Month FROM date)) FROM table_a) AS min_day", + ]; + let expected = vec![DateExtraction { + column: Column::new(Some("table_a"), "date"), + components: HashSet::from([SupportedIntervalUnit::Month, SupportedIntervalUnit::Day]), + }]; + for sql in statements { + general_test(sql, expected.clone()).await; + } + } + #[tokio::test] async fn variant_get_metadata_is_propagated() { let temp_dir = TempDir::new().unwrap(); diff --git a/src/parquet/src/optimizers/mod.rs b/src/parquet/src/optimizers/mod.rs index bd89f1ba..2d88600c 100644 --- a/src/parquet/src/optimizers/mod.rs +++ b/src/parquet/src/optimizers/mod.rs @@ -90,9 +90,11 @@ pub fn rewrite_data_source_plan( { let (metadata_key, metadata_value): (&str, String) = match annotation { - ColumnAnnotation::DatePart(unit) => ( + ColumnAnnotation::DatePart(_) => ( DATE_MAPPING_METADATA_KEY, - unit.metadata_value().to_string(), + annotation + .serialize_date_part() + .expect("DatePart should serialize"), ), ColumnAnnotation::VariantPath(path) => { (VARIANT_MAPPING_METADATA_KEY, path) From ad2a236443ef2d1ae3fa3654d8aa98eb73bf172d Mon Sep 17 00:00:00 2001 From: Nikhil Nayak Date: Wed, 19 Nov 2025 19:00:37 -0800 Subject: [PATCH 2/2] Implemented multiple DatePart caching (part of #393) --- ...ests__date_optimizer__date_extraction.snap | 12 +- ...date_optimizer__date_extraction_case2.snap | 12 +- ...__date_optimizer__date_extraction_day.snap | 12 +- ...date_optimizer__date_extraction_month.snap | 10 +- src/storage/bench/squeeze_date32.rs | 9 +- src/storage/src/cache/core.rs | 105 ++- src/storage/src/cache/expressions.rs | 166 ++++- src/storage/src/cache/mod.rs | 2 +- src/storage/src/liquid_array/mod.rs | 2 +- .../src/liquid_array/primitive_array.rs | 20 +- .../src/liquid_array/squeezed_date32_array.rs | 614 +++++++++++++----- src/storage/src/liquid_array/utils.rs | 19 +- 12 files changed, 741 insertions(+), 242 deletions(-) diff --git a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction.snap b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction.snap index a172b6a4..7ce4c5df 100644 --- a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction.snap +++ b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction.snap @@ -4,13 +4,13 @@ expression: stats --- entries.total: 123 entries.after_first_run: 123 -entries.memory.arrow: 4 -entries.memory.liquid: 38 -entries.memory.hybrid_liquid: 81 +entries.memory.arrow: 2 +entries.memory.liquid: 26 +entries.memory.hybrid_liquid: 95 entries.disk.liquid: 0 entries.disk.arrow: 0 -usage.memory_bytes: 1032554 -usage.disk_bytes: 1081512 +usage.memory_bytes: 976214 +usage.disk_bytes: 1268440 runtime.get_arrow_array_calls: 0 runtime.get_with_selection_calls: 123 runtime.get_with_predicate_calls: 0 @@ -18,4 +18,4 @@ runtime.get_predicate_hybrid_success: 0 runtime.get_predicate_hybrid_needs_io: 0 runtime.get_predicate_hybrid_unsupported: 0 runtime.try_read_liquid_calls: 0 -runtime.hit_date32_expression_calls: 81 +runtime.hit_date32_expression_calls: 95 diff --git a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_case2.snap b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_case2.snap index a6121b0b..26d63cb4 100644 --- a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_case2.snap +++ b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_case2.snap @@ -4,13 +4,13 @@ expression: stats --- entries.total: 123 entries.after_first_run: 123 -entries.memory.arrow: 4 -entries.memory.liquid: 38 -entries.memory.hybrid_liquid: 81 +entries.memory.arrow: 2 +entries.memory.liquid: 26 +entries.memory.hybrid_liquid: 95 entries.disk.liquid: 0 entries.disk.arrow: 0 -usage.memory_bytes: 1032554 -usage.disk_bytes: 1081512 +usage.memory_bytes: 976214 +usage.disk_bytes: 1268440 runtime.get_arrow_array_calls: 0 runtime.get_with_selection_calls: 246 runtime.get_with_predicate_calls: 0 @@ -18,4 +18,4 @@ runtime.get_predicate_hybrid_success: 0 runtime.get_predicate_hybrid_needs_io: 0 runtime.get_predicate_hybrid_unsupported: 0 runtime.try_read_liquid_calls: 0 -runtime.hit_date32_expression_calls: 162 +runtime.hit_date32_expression_calls: 190 diff --git a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_day.snap b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_day.snap index a172b6a4..7ce4c5df 100644 --- a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_day.snap +++ b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_day.snap @@ -4,13 +4,13 @@ expression: stats --- entries.total: 123 entries.after_first_run: 123 -entries.memory.arrow: 4 -entries.memory.liquid: 38 -entries.memory.hybrid_liquid: 81 +entries.memory.arrow: 2 +entries.memory.liquid: 26 +entries.memory.hybrid_liquid: 95 entries.disk.liquid: 0 entries.disk.arrow: 0 -usage.memory_bytes: 1032554 -usage.disk_bytes: 1081512 +usage.memory_bytes: 976214 +usage.disk_bytes: 1268440 runtime.get_arrow_array_calls: 0 runtime.get_with_selection_calls: 123 runtime.get_with_predicate_calls: 0 @@ -18,4 +18,4 @@ runtime.get_predicate_hybrid_success: 0 runtime.get_predicate_hybrid_needs_io: 0 runtime.get_predicate_hybrid_unsupported: 0 runtime.try_read_liquid_calls: 0 -runtime.hit_date32_expression_calls: 81 +runtime.hit_date32_expression_calls: 95 diff --git a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_month.snap b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_month.snap index 9fa6a142..a172b6a4 100644 --- a/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_month.snap +++ b/src/local/src/tests/snapshots/liquid_cache_local__tests__date_optimizer__date_extraction_month.snap @@ -5,12 +5,12 @@ expression: stats entries.total: 123 entries.after_first_run: 123 entries.memory.arrow: 4 -entries.memory.liquid: 46 -entries.memory.hybrid_liquid: 73 +entries.memory.liquid: 38 +entries.memory.hybrid_liquid: 81 entries.disk.liquid: 0 entries.disk.arrow: 0 -usage.memory_bytes: 1023346 -usage.disk_bytes: 974696 +usage.memory_bytes: 1032554 +usage.disk_bytes: 1081512 runtime.get_arrow_array_calls: 0 runtime.get_with_selection_calls: 123 runtime.get_with_predicate_calls: 0 @@ -18,4 +18,4 @@ runtime.get_predicate_hybrid_success: 0 runtime.get_predicate_hybrid_needs_io: 0 runtime.get_predicate_hybrid_unsupported: 0 runtime.try_read_liquid_calls: 0 -runtime.hit_date32_expression_calls: 73 +runtime.hit_date32_expression_calls: 81 diff --git a/src/storage/bench/squeeze_date32.rs b/src/storage/bench/squeeze_date32.rs index 56e0061d..857fbae3 100644 --- a/src/storage/bench/squeeze_date32.rs +++ b/src/storage/bench/squeeze_date32.rs @@ -1,9 +1,10 @@ use arrow::array::{Array, ArrayRef, cast::AsArray}; +use arrow::compute::DatePart; use arrow::datatypes::Date32Type; use clap::Parser; use datafusion::prelude::*; use futures::StreamExt; -use liquid_cache_storage::liquid_array::{Date32Field, LiquidPrimitiveArray, SqueezedDate32Array}; +use liquid_cache_storage::liquid_array::{LiquidPrimitiveArray, SqueezedDate32Array}; #[global_allocator] static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc; @@ -76,9 +77,9 @@ async fn run_for_column(ctx: &SessionContext, col: &str, limit: Option) { let liquid = LiquidPrimitiveArray::::from_arrow_array(prim.clone()); total_liquid_bytes += liquid.get_array_memory_size(); - let squeezed_year = SqueezedDate32Array::from_liquid_date32(&liquid, Date32Field::Year); - let squeezed_month = SqueezedDate32Array::from_liquid_date32(&liquid, Date32Field::Month); - let squeezed_day = SqueezedDate32Array::from_liquid_date32(&liquid, Date32Field::Day); + let squeezed_year = SqueezedDate32Array::from_liquid_date32(&liquid, DatePart::Year); + let squeezed_month = SqueezedDate32Array::from_liquid_date32(&liquid, DatePart::Month); + let squeezed_day = SqueezedDate32Array::from_liquid_date32(&liquid, DatePart::Day); total_year_bytes += squeezed_year.get_array_memory_size(); total_month_bytes += squeezed_month.get_array_memory_size(); diff --git a/src/storage/src/cache/core.rs b/src/storage/src/cache/core.rs index 395d223b..80572152 100644 --- a/src/storage/src/cache/core.rs +++ b/src/storage/src/cache/core.rs @@ -781,16 +781,32 @@ impl CacheStorage { CachedData::MemoryHybridLiquid(array) => { if let Some(CacheExpression::ExtractDate32 { field }) = expression_ref && let Some(squeezed) = array.as_any().downcast_ref::() - && squeezed.field() == *field { - let component = Arc::new(squeezed.to_component_date32()) as ArrayRef; - self.runtime_stats.incr_hit_date32_expression(); - if let Some(selection) = selection { - let selection_array = BooleanArray::new(selection.clone(), None); - let filtered = arrow::compute::filter(&component, &selection_array).ok()?; - return Some(filtered); + for part in field.iter() { + if let Some(component) = squeezed.to_component_date32(part) { + let component = Arc::new(component) as ArrayRef; + self.runtime_stats.incr_hit_date32_expression(); + if let Some(selection) = selection { + let selection_array = BooleanArray::new(selection.clone(), None); + let filtered = + arrow::compute::filter(&component, &selection_array).ok()?; + return Some(filtered); + } + return Some(component); + } + + if let Some(component) = squeezed.compute_derived(part) { + let component = Arc::new(component) as ArrayRef; + self.runtime_stats.incr_hit_date32_expression(); + if let Some(selection) = selection { + let selection_array = BooleanArray::new(selection.clone(), None); + let filtered = + arrow::compute::filter(&component, &selection_array).ok()?; + return Some(filtered); + } + return Some(component); + } } - return Some(component); } if let Some(selection) = selection { match array.filter(selection) { @@ -1125,15 +1141,14 @@ impl<'a> IntoFuture for EvaluatePredicate<'a> { mod tests { use super::*; use crate::cache::{ - CacheEntry, CacheExpression, + CacheEntry, CacheExpression, DatePartSet, cache_policies::{CachePolicy, LruPolicy}, utils::{create_cache_store, create_test_array, create_test_arrow_array}, }; - use crate::liquid_array::{ - Date32Field, LiquidHybridArrayRef, LiquidPrimitiveArray, SqueezedDate32Array, - }; + use crate::liquid_array::{LiquidHybridArrayRef, LiquidPrimitiveArray, SqueezedDate32Array}; use crate::sync::thread; use arrow::array::{Array, Date32Array, Int32Array}; + use arrow::compute::DatePart; use arrow::datatypes::Date32Type; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -1206,7 +1221,7 @@ mod tests { let date_values = Date32Array::from(vec![Some(0), Some(365), None, Some(730)]); let liquid = LiquidPrimitiveArray::::from_arrow_array(date_values.clone()); - let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, Date32Field::Year); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePart::Year); let hybrid: LiquidHybridArrayRef = Arc::new(squeezed); store @@ -1215,7 +1230,7 @@ mod tests { let registry = store.expression_registry(); let expr_id = registry - .register(CacheExpression::extract_date32(Date32Field::Year)) + .register(CacheExpression::extract_date32(DatePart::Year)) .expect("register expr"); let result = store .get(&entry_id) @@ -1235,6 +1250,68 @@ mod tests { assert_eq!(result.value(3), 1972); } + #[tokio::test] + async fn test_get_arrow_array_with_expression_extracts_year_month() { + let store = create_cache_store(1 << 20, Box::new(LruPolicy::new())); + let entry_id = EntryID::from(42); + + // Dates: Jan 1 1970, Jan 1 1971, null, Jan 1 1972 + let date_values = Date32Array::from(vec![Some(0), Some(397), None, Some(790)]); + let liquid = LiquidPrimitiveArray::::from_arrow_array(date_values.clone()); + + // Cache both Year and Month components together + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + let hybrid: LiquidHybridArrayRef = Arc::new(squeezed); + + store + .insert_inner(entry_id, CacheEntry::memory_hybrid_liquid(hybrid.clone())) + .await; + + let registry = store.expression_registry(); + + // Test 1: Extract Year from the cached YearMonth data + let year_expr_id = registry + .register(CacheExpression::extract_date32(DatePart::Year)) + .expect("register year expr"); + let year_result = store + .get(&entry_id) + .with_expression_hint(year_expr_id) + .read() + .await + .expect("year array present"); + + let year_result = year_result + .as_any() + .downcast_ref::() + .expect("date32 year result"); + assert_eq!(year_result.len(), 4); + assert_eq!(year_result.value(0), 1970); + assert_eq!(year_result.value(1), 1971); + assert!(year_result.is_null(2)); + assert_eq!(year_result.value(3), 1972); + + // Test 2: Extract Month from the same cached YearMonth data + let month_expr_id = registry + .register(CacheExpression::extract_date32(DatePart::Month)) + .expect("register month expr"); + let month_result = store + .get(&entry_id) + .with_expression_hint(month_expr_id) + .read() + .await + .expect("month array present"); + + let month_result = month_result + .as_any() + .downcast_ref::() + .expect("date32 month result"); + assert_eq!(month_result.len(), 4); + assert_eq!(month_result.value(0), 1); // January + assert_eq!(month_result.value(1), 2); // February + assert!(month_result.is_null(2)); + assert_eq!(month_result.value(3), 3); // March + } + #[tokio::test] async fn test_cache_advice_strategies() { // Comprehensive test of all three advice types diff --git a/src/storage/src/cache/expressions.rs b/src/storage/src/cache/expressions.rs index 75a63aea..23fbe106 100644 --- a/src/storage/src/cache/expressions.rs +++ b/src/storage/src/cache/expressions.rs @@ -4,15 +4,141 @@ use std::collections::HashMap; use std::num::NonZeroU16; use std::sync::{Arc, RwLock}; -use crate::liquid_array::Date32Field; +use arrow::compute::DatePart; +use std::hash::{Hash, Hasher}; + +/// A set of date parts that can be extracted from a `Date32` array. +/// Does not include YearMonthDay because this should never be extracted as components. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DatePartSet { + /// The year part. + Year, + /// The month part. + Month, + /// The day part. + Day, + /// The year and month parts. + YearMonth, + /// The year and day parts. + YearDay, + /// The month and day parts. + MonthDay, +} + +impl DatePartSet { + /// Iterate over the date parts in the set. + pub fn iter(&self) -> DatePartSetIter { + let mut parts = [DatePart::Year; 3]; + let mut len = 0; + let push = |part: DatePart, parts: &mut [DatePart; 3], len: &mut usize| { + parts[*len] = part; + *len += 1; + }; + + match self { + Self::Year => push(DatePart::Year, &mut parts, &mut len), + Self::Month => push(DatePart::Month, &mut parts, &mut len), + Self::Day => push(DatePart::Day, &mut parts, &mut len), + Self::YearMonth => { + push(DatePart::Year, &mut parts, &mut len); + push(DatePart::Month, &mut parts, &mut len); + } + Self::YearDay => { + push(DatePart::Year, &mut parts, &mut len); + push(DatePart::Day, &mut parts, &mut len); + } + Self::MonthDay => { + push(DatePart::Month, &mut parts, &mut len); + push(DatePart::Day, &mut parts, &mut len); + } + } + + DatePartSetIter { parts, len, idx: 0 } + } + + /// Check if the set contains the given date part. + pub fn contains(&self, part: DatePart) -> bool { + match self { + Self::Year => part == DatePart::Year, + Self::Month => part == DatePart::Month, + Self::Day => part == DatePart::Day, + Self::YearMonth => part == DatePart::Year || part == DatePart::Month, + Self::YearDay => part == DatePart::Year || part == DatePart::Day, + Self::MonthDay => part == DatePart::Month || part == DatePart::Day, + } + } + + /// Check if the set contains the year part. + pub fn has_year(&self) -> bool { + matches!(self, Self::Year | Self::YearMonth | Self::YearDay) + } + + /// Check if the set contains the month part. + pub fn has_month(&self) -> bool { + matches!(self, Self::Month | Self::YearMonth | Self::MonthDay) + } + + /// Check if the set contains the day part. + pub fn has_day(&self) -> bool { + matches!(self, Self::Day | Self::YearDay | Self::MonthDay) + } + + /// Get the length of the set. + pub fn len(&self) -> usize { + match self { + Self::Year | Self::Month | Self::Day => 1, + _ => 2, + } + } + + /// Check if the set is empty. + pub fn is_empty(&self) -> bool { + false + } + + /// Check if the set is a superset of the other set. + pub fn is_superset_of(&self, other: DatePartSet) -> bool { + other.iter().all(|part| self.contains(part)) + } +} + +pub struct DatePartSetIter { + parts: [DatePart; 3], + len: usize, + idx: usize, +} + +impl Iterator for DatePartSetIter { + type Item = DatePart; + + fn next(&mut self) -> Option { + if self.idx >= self.len { + return None; + } + let part = self.parts[self.idx]; + self.idx += 1; + Some(part) + } +} + +impl From for DatePartSet { + fn from(value: DatePart) -> Self { + match value { + DatePart::Year => DatePartSet::Year, + DatePart::Month => DatePartSet::Month, + DatePart::Day => DatePartSet::Day, + _ => panic!("Unsupported DatePart variant for caching: {:?}", value), + } + } +} /// Experimental expression descriptor for cache lookups. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq)] pub enum CacheExpression { /// Extract a specific component (YEAR/MONTH/DAY) from a `Date32` column. ExtractDate32 { - /// Component to extract (YEAR/MONTH/DAY). - field: Date32Field, + /// Components to extract (YEAR/MONTH/DAY). + field: DatePartSet, }, /// Extract a field from a variant column via `variant_get`. VariantGet { @@ -21,10 +147,27 @@ pub enum CacheExpression { }, } +impl Hash for CacheExpression { + fn hash(&self, state: &mut H) { + match self { + Self::ExtractDate32 { field } => { + 0u8.hash(state); + std::mem::discriminant(field).hash(state); + } + Self::VariantGet { path } => { + 1u8.hash(state); + path.hash(state); + } + } + } +} + impl CacheExpression { /// Build an extract expression for a `Date32` column. - pub fn extract_date32(field: Date32Field) -> Self { - Self::ExtractDate32 { field } + pub fn extract_date32(field: impl Into) -> Self { + Self::ExtractDate32 { + field: field.into(), + } } /// Build a variant-get expression for the provided dotted path. @@ -38,16 +181,19 @@ impl CacheExpression { pub fn try_from_date_part_str(value: &str) -> Option { let upper = value.to_ascii_uppercase(); let field = match upper.as_str() { - "YEAR" => Date32Field::Year, - "MONTH" => Date32Field::Month, - "DAY" => Date32Field::Day, + "YEAR" => DatePartSet::Year, + "MONTH" => DatePartSet::Month, + "DAY" => DatePartSet::Day, + "YEAR,MONTH" => DatePartSet::YearMonth, + "YEAR,DAY" => DatePartSet::YearDay, + "MONTH,DAY" => DatePartSet::MonthDay, _ => return None, }; Some(Self::ExtractDate32 { field }) } /// Return the requested `Date32` component when this is an extract expression. - pub fn as_date32_field(&self) -> Option { + pub fn as_date32_field(&self) -> Option { match self { Self::ExtractDate32 { field } => Some(*field), Self::VariantGet { .. } => None, diff --git a/src/storage/src/cache/mod.rs b/src/storage/src/cache/mod.rs index 5cbcce45..a1622cf2 100644 --- a/src/storage/src/cache/mod.rs +++ b/src/storage/src/cache/mod.rs @@ -18,7 +18,7 @@ pub use core::{ BlockingIoContext, CacheStorage, CacheStorageBuilder, DefaultIoContext, EvaluatePredicate, Get, Insert, IoContext, }; -pub use expressions::{CacheExpression, ExpressionId, ExpressionRegistry}; +pub use expressions::{CacheExpression, DatePartSet, ExpressionId, ExpressionRegistry}; pub use stats::{CacheStats, RuntimeStats, RuntimeStatsSnapshot}; pub use transcode::transcode_liquid_inner; pub use utils::{EntryID, LiquidCompressorStates}; diff --git a/src/storage/src/liquid_array/mod.rs b/src/storage/src/liquid_array/mod.rs index c40c7a0d..e8de03dd 100644 --- a/src/storage/src/liquid_array/mod.rs +++ b/src/storage/src/liquid_array/mod.rs @@ -41,7 +41,7 @@ pub use primitive_array::{ LiquidI64Array, LiquidPrimitiveArray, LiquidPrimitiveDeltaArray, LiquidPrimitiveType, LiquidU8Array, LiquidU16Array, LiquidU32Array, LiquidU64Array, }; -pub use squeezed_date32_array::{Date32Field, SqueezedDate32Array}; +pub use squeezed_date32_array::SqueezedDate32Array; pub use variant_array::VariantExtractedArray; use crate::cache::CacheExpression; diff --git a/src/storage/src/liquid_array/primitive_array.rs b/src/storage/src/liquid_array/primitive_array.rs index 70996b84..74e78a42 100644 --- a/src/storage/src/liquid_array/primitive_array.rs +++ b/src/storage/src/liquid_array/primitive_array.rs @@ -17,15 +17,13 @@ use fastlanes::BitPacking; use num_traits::{AsPrimitive, FromPrimitive}; use super::LiquidDataType; -use crate::cache::CacheExpression; +use crate::cache::{CacheExpression, DatePartSet}; use crate::liquid_array::hybrid_primitive_array::{ LiquidPrimitiveClampedArray, LiquidPrimitiveQuantizedArray, }; use crate::liquid_array::ipc::{LiquidIPCHeader, PhysicalTypeMarker, get_physical_type_id}; use crate::liquid_array::raw::BitPackedArray; -use crate::liquid_array::{ - Date32Field, LiquidArray, LiquidHybridArrayRef, PrimitiveKind, SqueezedDate32Array, -}; +use crate::liquid_array::{LiquidArray, LiquidHybridArrayRef, PrimitiveKind, SqueezedDate32Array}; use crate::utils::get_bit_width; use arrow::datatypes::ArrowNativeType; use bytes::Bytes; @@ -402,7 +400,7 @@ where // Special handle for Date32 arrays with component extraction support. let field = expression_hint .and_then(|expr| expr.as_date32_field()) - .unwrap_or(Date32Field::Year); + .unwrap_or(DatePartSet::Year); return Some(( Arc::new(SqueezedDate32Array::from_liquid_date32(self, field)) as LiquidHybridArrayRef, @@ -885,10 +883,10 @@ mod tests { let registry = ExpressionRegistry::new(); let expr_month = registry - .register(CacheExpression::extract_date32(Date32Field::Month)) + .register(CacheExpression::extract_date32(DatePartSet::Month)) .expect("register month"); let expr_year = registry - .register(CacheExpression::extract_date32(Date32Field::Year)) + .register(CacheExpression::extract_date32(DatePartSet::Year)) .expect("register year"); let tracker = ExpressionHintTracker::new(); @@ -906,7 +904,7 @@ mod tests { .downcast_ref::() .expect("expected squeezed date32 array"); - assert_eq!(squeezed.field(), Date32Field::Month); + assert_eq!(squeezed.field(), DatePartSet::Month); } #[test] @@ -917,10 +915,10 @@ mod tests { let registry = ExpressionRegistry::new(); let expr_year = registry - .register(CacheExpression::extract_date32(Date32Field::Year)) + .register(CacheExpression::extract_date32(DatePartSet::Year)) .expect("register year"); let expr_day = registry - .register(CacheExpression::extract_date32(Date32Field::Day)) + .register(CacheExpression::extract_date32(DatePartSet::Day)) .expect("register day"); let tracker = ExpressionHintTracker::new(); tracker.record_expression(expr_year); @@ -937,7 +935,7 @@ mod tests { .downcast_ref::() .expect("expected squeezed date32 array"); - assert_eq!(squeezed.field(), Date32Field::Day); + assert_eq!(squeezed.field(), DatePartSet::Day); } test_roundtrip!( diff --git a/src/storage/src/liquid_array/squeezed_date32_array.rs b/src/storage/src/liquid_array/squeezed_date32_array.rs index f4beed52..427a5c04 100644 --- a/src/storage/src/liquid_array/squeezed_date32_array.rs +++ b/src/storage/src/liquid_array/squeezed_date32_array.rs @@ -1,4 +1,4 @@ -use arrow::array::{ArrayRef, BooleanArray, PrimitiveArray, cast::AsArray}; +use arrow::array::{Array, ArrayRef, BooleanArray, PrimitiveArray, cast::AsArray}; use arrow::buffer::{BooleanBuffer, ScalarBuffer}; use arrow::datatypes::{ArrowPrimitiveType, Date32Type, Int32Type, UInt32Type}; use arrow_schema::DataType; @@ -7,211 +7,234 @@ use std::sync::Arc; use super::LiquidArray; use super::primitive_array::LiquidPrimitiveArray; use super::{IoRange, LiquidArrayRef, LiquidDataType, LiquidHybridArray}; +use crate::cache::DatePartSet; use crate::liquid_array::LiquidPrimitiveType; use crate::liquid_array::raw::BitPackedArray; use crate::utils::get_bit_width; +use arrow::compute::DatePart; -/// Which component to extract from a `Date32` (days since UNIX epoch). -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -pub enum Date32Field { - /// Year component - Year, - /// Month component - Month, - /// Day component - Day, -} - -/// A bit-packed array that stores a single extracted component (YEAR/MONTH/DAY) -/// from a `Date32` array. -/// -/// Values are stored as unsigned offsets from `reference_value`, using the same -/// bit-packing machinery as primitive arrays. #[derive(Debug, Clone)] -pub struct SqueezedDate32Array { - field: Date32Field, +struct ComponentChunk { + part: DatePart, bit_packed: BitPackedArray, - /// The minimum extracted value used as reference for offsetting. reference_value: i32, } -impl SqueezedDate32Array { - /// Build a squeezed representation (YEAR/MONTH/DAY) from a `LiquidPrimitiveArray`. - pub fn from_liquid_date32( - array: &LiquidPrimitiveArray, - field: Date32Field, - ) -> Self { - // Decode the logical Date32 array (i32: days since epoch) from the liquid array. - let arrow_array: PrimitiveArray = - array.to_arrow_array().as_primitive::().clone(); +impl ComponentChunk { + fn len(&self) -> usize { + self.bit_packed.len() + } - let (_dt, values, nulls) = arrow_array.into_parts(); - - // Compute min and max for the extracted component, skipping nulls. - let mut has_value = false; - let mut min_component: i32 = i32::MAX; - let mut max_component: i32 = i32::MIN; - - // Fast path: if all nulls, return a null bit-packed array of the same length. - if let Some(nulls_buf) = &nulls - && nulls_buf.null_count() == values.len() - { - return Self { - field, - bit_packed: BitPackedArray::new_null_array(values.len()), - reference_value: 0, - }; - } + fn get_array_memory_size(&self) -> usize { + self.bit_packed.get_array_memory_size() + std::mem::size_of::() + } - for (idx, &days) in values.iter().enumerate() { - if let Some(nulls_buf) = &nulls - && nulls_buf.is_null(idx) - { - continue; - } - let (year, month, day) = ymd_from_epoch_days(days); - let comp = match field { - Date32Field::Year => year, - Date32Field::Month => month as i32, - Date32Field::Day => day as i32, - }; - has_value = true; - if comp < min_component { - min_component = comp; - } - if comp > max_component { - max_component = comp; - } + fn new_null(part: DatePart, len: usize) -> Self { + Self { + part, + bit_packed: BitPackedArray::new_null_array(len), + reference_value: 0, } + } + + fn from_components(part: DatePart, comps_array: PrimitiveArray) -> Self { + let len = comps_array.len(); + let min_component = + arrow::compute::kernels::aggregate::min(&comps_array).unwrap_or(i32::MAX); + let max_component = + arrow::compute::kernels::aggregate::max(&comps_array).unwrap_or(i32::MIN); - // If no non-null values found, return an all-null structure (defensive) + let has_value = min_component != i32::MAX && max_component != i32::MIN; if !has_value { - return Self { - field, - bit_packed: BitPackedArray::new_null_array(values.len()), - reference_value: 0, - }; + return Self::new_null(part, len); } - // Compute bit width from the value range. let max_offset = (max_component as i64 - min_component as i64) as u64; let bit_width = get_bit_width(max_offset); - // Build unsigned offsets for packing; placeholders are fine for nulls. - let offsets: ScalarBuffer<::Native> = - ScalarBuffer::from_iter((0..values.len()).map(|idx| { - if nulls.as_ref().is_some_and(|n| n.is_null(idx)) { - 0u32 - } else { - let (year, month, day) = ymd_from_epoch_days(values[idx]); - let comp = match field { - Date32Field::Year => year, - Date32Field::Month => month as i32, - Date32Field::Day => day as i32, - }; - (comp - min_component) as u32 - } - })); + let (_dt, comps_values, comps_nulls) = comps_array.into_parts(); + let offsets: ScalarBuffer = ScalarBuffer::from_iter( + comps_values + .iter() + .map(|&v| (v.saturating_sub(min_component)) as u32), + ); - let unsigned_array = PrimitiveArray::::new(offsets, nulls); + let unsigned_array = PrimitiveArray::::new(offsets, comps_nulls); let bit_packed = BitPackedArray::from_primitive(unsigned_array, bit_width); Self { - field, + part, bit_packed, reference_value: min_component, } } - /// Length of the array. - pub fn len(&self) -> usize { - self.bit_packed.len() - } - - /// Whether the array has no elements. - pub fn is_empty(&self) -> bool { - self.len() == 0 - } - - /// Memory size of the bit-packed representation plus reference value. - pub fn get_array_memory_size(&self) -> usize { - self.bit_packed.get_array_memory_size() + std::mem::size_of::() - } - - /// The extracted component type. - pub fn field(&self) -> Date32Field { - self.field - } - - /// Convert back to an Arrow `Int32` array representing the extracted component values. - /// Useful for verification or future pushdown logic. - pub fn to_component_date32(&self) -> PrimitiveArray { + fn to_component_array(&self) -> PrimitiveArray { let unsigned: PrimitiveArray = self.bit_packed.to_primitive(); let (_dt, values, nulls) = unsigned.into_parts(); - let ref_v = self.reference_value; let signed_values: ScalarBuffer<::Native> = - ScalarBuffer::from_iter(values.iter().map(|&v| (v as i32).saturating_add(ref_v))); + ScalarBuffer::from_iter( + values + .iter() + .map(|&v| (v as i32).saturating_add(self.reference_value)), + ); PrimitiveArray::::new(signed_values, nulls) } - /// Lossy reconstruction to Arrow `Date32` (days since epoch). - /// - /// Mapping used: - /// - Year: (year, 1, 1) - /// - Month: (1970, month, 1) - /// - Day: (1970, 1, day) - pub fn to_arrow_date32_lossy(&self) -> PrimitiveArray { + fn to_lossy_date32(&self) -> PrimitiveArray { let unsigned: PrimitiveArray = self.bit_packed.to_primitive(); let (_dt, values, nulls) = unsigned.into_parts(); - let ref_v = self.reference_value; let days_values: ScalarBuffer<::Native> = ScalarBuffer::from_iter(values.iter().enumerate().map(|(i, &off)| { if nulls.as_ref().is_some_and(|n| n.is_null(i)) { 0i32 } else { - match self.field { - Date32Field::Year => { + match self.part { + DatePart::Year => { let y = ref_v + off as i32; ymd_to_epoch_days(y, 1, 1) } - Date32Field::Month => { + DatePart::Month => { let m = (ref_v + off as i32) as u32; ymd_to_epoch_days(1970, m, 1) } - Date32Field::Day => { + DatePart::Day => { let d = (ref_v + off as i32) as u32; ymd_to_epoch_days(1970, 1, d) } + _ => unreachable!("Stored field should be Year/Month/Day"), } } })); - PrimitiveArray::::new(days_values, nulls) } + + fn derive_with_arrow(&self, requested: DatePart) -> PrimitiveArray { + let lossy = self.to_lossy_date32(); + let derived_array = arrow::compute::date_part(&lossy, requested).unwrap(); + let derived_int32 = derived_array.as_primitive::(); + let (_dt, derived_values, derived_nulls) = derived_int32.clone().into_parts(); + let derived_date32_values: ScalarBuffer<::Native> = + ScalarBuffer::from_iter(derived_values.iter().copied()); + PrimitiveArray::::new(derived_date32_values, derived_nulls) + } } -/// Convert days since UNIX epoch (1970-01-01) to (year, month, day) in the -/// proleptic Gregorian calendar using a branchless integer algorithm. -fn ymd_from_epoch_days(days_since_epoch: i32) -> (i32, u32, u32) { - // Based on Howard Hinnant's civil_from_days algorithm. - let z = days_since_epoch as i64 + 719_468; // shift to civil epoch - let era = if z >= 0 { - z / 146_097 - } else { - (z - 146_096) / 146_097 - }; - let doe = z - era * 146_097; // [0, 146096] - let yoe = (doe - doe / 1_460 + doe / 36_524 - doe / 146_096) / 365; // [0, 399] - let mut y = yoe + era * 400; - let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); // [0, 365] - let mp = (5 * doy + 2) / 153; // [0, 11] - let d = (doy - (153 * mp + 2) / 5) + 1; // [1, 31] - let m = mp + if mp < 10 { 3 } else { -9 }; // [1, 12] - if m <= 2 { - y += 1; - } - (y as i32, m as u32, d as u32) +/// A bit-packed array that stores extracted components (YEAR/MONTH/DAY) +/// from a `Date32` array. +/// +/// Each component is stored as unsigned offsets from its component-specific `reference_value`. +#[derive(Debug, Clone)] +pub struct SqueezedDate32Array { + parts: DatePartSet, + chunks: Vec, +} + +impl SqueezedDate32Array { + /// Build a squeezed representation (YEAR/MONTH/DAY) from a `LiquidPrimitiveArray`. + pub fn from_liquid_date32( + array: &LiquidPrimitiveArray, + field: impl Into, + ) -> Self { + let parts = field.into(); + let arrow_array: PrimitiveArray = + array.to_arrow_array().as_primitive::().clone(); + let len = arrow_array.len(); + + let all_null = arrow_array.null_count() == len && len > 0; + let chunks = parts + .iter() + .map(|part| { + if all_null { + ComponentChunk::new_null(part, len) + } else { + Self::build_chunk(part, &arrow_array) + } + }) + .collect(); + + Self { parts, chunks } + } + + fn build_chunk(part: DatePart, arrow_array: &PrimitiveArray) -> ComponentChunk { + let comps = arrow::compute::date_part(arrow_array, part).unwrap(); + let comps_array = comps.as_primitive::().clone(); + ComponentChunk::from_components(part, comps_array) + } + + fn chunk(&self, part: DatePart) -> Option<&ComponentChunk> { + self.chunks.iter().find(|chunk| chunk.part == part) + } + + fn primary_chunk(&self) -> Option<&ComponentChunk> { + self.parts.iter().next().and_then(|part| self.chunk(part)) + } + + /// Length of the array. + pub fn len(&self) -> usize { + self.chunks.first().map(|chunk| chunk.len()).unwrap_or(0) + } + + /// Whether the array has no elements. + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + /// Memory size of the bit-packed representation plus reference values. + pub fn get_array_memory_size(&self) -> usize { + self.chunks + .iter() + .map(|chunk| chunk.get_array_memory_size()) + .sum() + } + + /// The extracted component set. + pub fn field(&self) -> DatePartSet { + self.parts + } + + /// Whether the array can compute the requested derived field. + /// Quarter can be derived from Month, DayOfYear from Month+Day. + pub fn can_derive_field(&self, requested: DatePart) -> bool { + if self.chunk(requested).is_some() { + return true; + } + match requested { + DatePart::Quarter => self.chunk(DatePart::Month).is_some(), + DatePart::DayOfYear => { + self.chunk(DatePart::Month).is_some() && self.chunk(DatePart::Day).is_some() + } + _ => false, + } + } + + /// Compute a DatePart from the stored components. + /// Returns Date32Type array (same pattern as to_component_date32) for consistency. + pub fn compute_derived(&self, requested: DatePart) -> Option> { + if let Some(chunk) = self.chunk(requested) { + return Some(chunk.to_component_array()); + } + + match requested { + DatePart::Quarter => self + .chunk(DatePart::Month) + .map(|chunk| chunk.derive_with_arrow(DatePart::Quarter)), + _ => None, + } + } + + /// Convert back to an Arrow `Int32` array representing the extracted component values. + pub fn to_component_date32(&self, part: DatePart) -> Option> { + self.chunk(part).map(|chunk| chunk.to_component_array()) + } + + /// Lossy reconstruction to Arrow `Date32` (days since epoch) using the primary component. + pub fn to_arrow_date32_lossy(&self) -> PrimitiveArray { + self.primary_chunk() + .map(|chunk| chunk.to_lossy_date32()) + .unwrap_or_else(|| PrimitiveArray::::new_null(0)) + } } /// Convert a date (year, month, day) in proleptic Gregorian calendar to @@ -260,32 +283,38 @@ impl LiquidHybridArray for SqueezedDate32Array { } fn filter(&self, selection: &BooleanBuffer) -> Result { - let unsigned_array: PrimitiveArray = self.bit_packed.to_primitive(); + let Some(chunk) = self.primary_chunk() else { + return Ok(Arc::new(PrimitiveArray::::new_null(0))); + }; + + let unsigned_array: PrimitiveArray = chunk.bit_packed.to_primitive(); let selection = BooleanArray::new(selection.clone(), None); let filtered_values = arrow::compute::kernels::filter::filter(&unsigned_array, &selection).unwrap(); let filtered_values = filtered_values.as_primitive::().clone(); // Reconstruct lossy Date32 directly from filtered offsets let (_dt, values, nulls) = filtered_values.into_parts(); - let ref_v = self.reference_value; + let ref_v = chunk.reference_value; let days_values: ScalarBuffer<::Native> = ScalarBuffer::from_iter(values.iter().enumerate().map(|(i, &off)| { if nulls.as_ref().is_some_and(|n| n.is_null(i)) { 0i32 } else { - match self.field { - Date32Field::Year => { + match chunk.part { + DatePart::Year => { let y = ref_v + off as i32; ymd_to_epoch_days(y, 1, 1) } - Date32Field::Month => { + DatePart::Month => { let m = (ref_v + off as i32) as u32; ymd_to_epoch_days(1970, m, 1) } - Date32Field::Day => { + DatePart::Day => { let d = (ref_v + off as i32) as u32; ymd_to_epoch_days(1970, 1, d) } + + _ => unreachable!("Stored field should be Year/Month/Day"), } } })); @@ -313,6 +342,7 @@ impl LiquidHybridArray for SqueezedDate32Array { #[cfg(test)] mod tests { use super::*; + use crate::cache::DatePartSet; use arrow::array::PrimitiveArray; use std::sync::Arc; @@ -326,20 +356,60 @@ mod tests { assert_eq!(a_ref.as_ref(), b_ref.as_ref()); } - fn extract(field: Date32Field, input: Vec>) -> PrimitiveArray { + fn extract(field: DatePart, input: Vec>) -> PrimitiveArray { let arr = dates(&input); let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, field); - squeezed.to_component_date32() + squeezed + .to_component_date32(field) + .expect("component should exist") } - fn lossy(field: Date32Field, input: Vec>) -> PrimitiveArray { + fn lossy(field: DatePart, input: Vec>) -> PrimitiveArray { let arr = dates(&input); let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, field); squeezed.to_arrow_date32_lossy() } + #[test] + fn test_year_month_components() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 1, 1)), + Some(ymd_to_epoch_days(1981, 12, 31)), + None, + Some(ymd_to_epoch_days(1999, 7, 4)), + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + + assert_eq!(squeezed.field(), DatePartSet::YearMonth); + let years = squeezed + .to_component_date32(DatePart::Year) + .expect("year chunk"); + let months = squeezed + .to_component_date32(DatePart::Month) + .expect("month chunk"); + assert!(squeezed.to_component_date32(DatePart::Day).is_none()); + + let expected_years = + PrimitiveArray::::from(vec![Some(1970), Some(1981), None, Some(1999)]); + let expected_months = + PrimitiveArray::::from(vec![Some(1), Some(12), None, Some(7)]); + + assert_prim_eq(years, expected_years); + assert_prim_eq(months, expected_months); + + let quarters = squeezed + .compute_derived(DatePart::Quarter) + .expect("quarter derivation"); + let expected_quarters = + PrimitiveArray::::from(vec![Some(1), Some(4), None, Some(3)]); + assert_prim_eq(quarters, expected_quarters); + } + #[test] fn test_extraction_correctness() { // YEAR @@ -351,7 +421,7 @@ mod tests { ]; let expected = PrimitiveArray::::from(vec![Some(1969), Some(1970), Some(1971), None]); - assert_prim_eq(extract(Date32Field::Year, input), expected); + assert_prim_eq(extract(DatePart::Year, input), expected); // MONTH let input = vec![ @@ -361,7 +431,7 @@ mod tests { None, ]; let expected = PrimitiveArray::::from(vec![Some(1), Some(2), Some(12), None]); - assert_prim_eq(extract(Date32Field::Month, input), expected); + assert_prim_eq(extract(DatePart::Month, input), expected); // DAY let input = vec![ @@ -371,7 +441,7 @@ mod tests { None, ]; let expected = PrimitiveArray::::from(vec![Some(1), Some(31), Some(1), None]); - assert_prim_eq(extract(Date32Field::Day, input), expected); + assert_prim_eq(extract(DatePart::Day, input), expected); } #[test] @@ -387,7 +457,7 @@ mod tests { Some(ymd_to_epoch_days(2000, 1, 1)), None, ]); - assert_prim_eq(lossy(Date32Field::Year, input), expected); + assert_prim_eq(lossy(DatePart::Year, input), expected); // MONTH → (1970,m,1) let input = vec![ @@ -400,7 +470,7 @@ mod tests { Some(ymd_to_epoch_days(1970, 12, 1)), None, ]); - assert_prim_eq(lossy(Date32Field::Month, input), expected); + assert_prim_eq(lossy(DatePart::Month, input), expected); // DAY → (1970,1,d) let input = vec![ @@ -413,7 +483,7 @@ mod tests { Some(ymd_to_epoch_days(1970, 1, 5)), None, ]); - assert_prim_eq(lossy(Date32Field::Day, input), expected); + assert_prim_eq(lossy(DatePart::Day, input), expected); } #[test] @@ -427,12 +497,13 @@ mod tests { None, ]; - for &field in &[Date32Field::Year, Date32Field::Month, Date32Field::Day] { + for &field in &[DatePart::Year, DatePart::Month, DatePart::Day] { let comp1 = extract(field, input.clone()); let lossy_dt = lossy(field, input.clone()); let liquid2 = LiquidPrimitiveArray::::from_arrow_array(lossy_dt); - let comp2 = - SqueezedDate32Array::from_liquid_date32(&liquid2, field).to_component_date32(); + let comp2 = SqueezedDate32Array::from_liquid_date32(&liquid2, field) + .to_component_date32(field) + .expect("component should exist"); assert_prim_eq(comp1, comp2); } } @@ -441,7 +512,7 @@ mod tests { fn test_all_nulls_behavior() { let input = vec![None, None, None]; - for &field in &[Date32Field::Year, Date32Field::Month, Date32Field::Day] { + for &field in &[DatePart::Year, DatePart::Month, DatePart::Day] { let comp = extract(field, input.clone()); let expected_comp = PrimitiveArray::::from(vec![None, None, None]); assert_prim_eq(comp, expected_comp); @@ -451,4 +522,213 @@ mod tests { assert_prim_eq(lossy_dt, expected_dt); } } + + #[test] + fn test_month_day_components() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 1, 1)), + Some(ymd_to_epoch_days(1981, 12, 31)), + None, + Some(ymd_to_epoch_days(1999, 7, 4)), + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::MonthDay); + + assert_eq!(squeezed.field(), DatePartSet::MonthDay); + let months = squeezed + .to_component_date32(DatePart::Month) + .expect("month chunk"); + let days = squeezed + .to_component_date32(DatePart::Day) + .expect("day chunk"); + assert!(squeezed.to_component_date32(DatePart::Year).is_none()); + + let expected_months = + PrimitiveArray::::from(vec![Some(1), Some(12), None, Some(7)]); + let expected_days = + PrimitiveArray::::from(vec![Some(1), Some(31), None, Some(4)]); + + assert_prim_eq(months, expected_months); + assert_prim_eq(days, expected_days); + } + + #[test] + fn test_year_day_components() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 1, 1)), + Some(ymd_to_epoch_days(1981, 12, 31)), + None, + Some(ymd_to_epoch_days(1999, 7, 4)), + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearDay); + + assert_eq!(squeezed.field(), DatePartSet::YearDay); + let years = squeezed + .to_component_date32(DatePart::Year) + .expect("year chunk"); + let days = squeezed + .to_component_date32(DatePart::Day) + .expect("day chunk"); + assert!(squeezed.to_component_date32(DatePart::Month).is_none()); + + let expected_years = + PrimitiveArray::::from(vec![Some(1970), Some(1981), None, Some(1999)]); + let expected_days = + PrimitiveArray::::from(vec![Some(1), Some(31), None, Some(4)]); + + assert_prim_eq(years, expected_years); + assert_prim_eq(days, expected_days); + } + + #[test] + fn test_can_derive_field() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 3, 15)), + Some(ymd_to_epoch_days(1981, 7, 4)), + None, + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + + // YearMonth can derive Quarter (from Month) + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + assert!(squeezed.can_derive_field(DatePart::Year)); + assert!(squeezed.can_derive_field(DatePart::Month)); + assert!(squeezed.can_derive_field(DatePart::Quarter)); + assert!(!squeezed.can_derive_field(DatePart::Day)); + assert!(!squeezed.can_derive_field(DatePart::DayOfYear)); + + // Month only can derive Quarter + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::Month); + assert!(squeezed.can_derive_field(DatePart::Month)); + assert!(squeezed.can_derive_field(DatePart::Quarter)); + assert!(!squeezed.can_derive_field(DatePart::Year)); + assert!(!squeezed.can_derive_field(DatePart::Day)); + + // MonthDay can derive Quarter + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::MonthDay); + assert!(squeezed.can_derive_field(DatePart::Month)); + assert!(squeezed.can_derive_field(DatePart::Day)); + assert!(squeezed.can_derive_field(DatePart::Quarter)); + assert!(!squeezed.can_derive_field(DatePart::Year)); + } + + #[test] + fn test_cannot_compute_derived_fields() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 3, 15)), + Some(ymd_to_epoch_days(1981, 7, 4)), + None, + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + + // YearMonth cannot derive Day + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + assert!(squeezed.compute_derived(DatePart::Day).is_none()); + assert!(squeezed.compute_derived(DatePart::Year).is_some()); + assert!(squeezed.compute_derived(DatePart::Month).is_some()); + + // Month only cannot derive Year or Day + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::Month); + assert!(squeezed.compute_derived(DatePart::Year).is_none()); + assert!(squeezed.compute_derived(DatePart::Day).is_none()); + assert!(squeezed.compute_derived(DatePart::Quarter).is_some()); + } + + #[test] + fn test_filter_multi_part() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 1, 1)), + Some(ymd_to_epoch_days(1981, 12, 31)), + Some(ymd_to_epoch_days(1999, 7, 4)), + None, + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + + // Filter: keep first and third elements + let selection = BooleanBuffer::from(vec![true, false, true, false]); + let filtered = squeezed.filter(&selection).expect("filter should succeed"); + let filtered_arr = filtered + .as_any() + .downcast_ref::>() + .expect("should be Date32 array"); + + assert_eq!(filtered_arr.len(), 2); + // Lossy reconstruction uses primary component (Year), so should be (1970,1,1) and (1999,1,1) + assert_eq!(filtered_arr.value(0), ymd_to_epoch_days(1970, 1, 1)); + assert_eq!(filtered_arr.value(1), ymd_to_epoch_days(1999, 1, 1)); + } + + #[test] + fn test_memory_size_multi_part() { + let input = vec![ + Some(ymd_to_epoch_days(1970, 1, 1)), + Some(ymd_to_epoch_days(1981, 12, 31)), + Some(ymd_to_epoch_days(1999, 7, 4)), + ]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + + // Single part + let single = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::Year); + let single_size = single.get_array_memory_size(); + + // Multi-part should be larger + let multi = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + let multi_size = multi.get_array_memory_size(); + + assert!( + multi_size > single_size, + "multi-part should use more memory" + ); + } + + #[test] + fn test_empty_array() { + let input: Vec> = vec![]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + + assert_eq!(squeezed.len(), 0); + assert!(squeezed.is_empty()); + } + + #[test] + fn test_all_nulls_multi_part() { + let input = vec![None, None, None]; + + let arr = dates(&input); + let liquid = LiquidPrimitiveArray::::from_arrow_array(arr); + let squeezed = SqueezedDate32Array::from_liquid_date32(&liquid, DatePartSet::YearMonth); + + assert_eq!(squeezed.len(), 3); + assert_eq!(squeezed.field(), DatePartSet::YearMonth); + + // Should still be able to get components (all null) + let years = squeezed + .to_component_date32(DatePart::Year) + .expect("year chunk should exist"); + assert_eq!(years.null_count(), 3); + assert_eq!(years.len(), 3); + + let months = squeezed + .to_component_date32(DatePart::Month) + .expect("month chunk should exist"); + assert_eq!(months.null_count(), 3); + assert_eq!(months.len(), 3); + } } diff --git a/src/storage/src/liquid_array/utils.rs b/src/storage/src/liquid_array/utils.rs index ab6f78b4..580e2638 100644 --- a/src/storage/src/liquid_array/utils.rs +++ b/src/storage/src/liquid_array/utils.rs @@ -1,7 +1,6 @@ use std::sync::atomic::{AtomicU64, Ordering}; use crate::cache::ExpressionId; - #[cfg(test)] pub(crate) fn gen_test_decimal_array( data_type: arrow_schema::DataType, @@ -111,10 +110,8 @@ impl Clone for ExpressionHintTracker { #[cfg(test)] mod tests { use super::ExpressionHintTracker; - use crate::{ - cache::{CacheExpression, ExpressionRegistry}, - liquid_array::Date32Field, - }; + use crate::cache::{CacheExpression, ExpressionRegistry}; + use arrow::compute::DatePart; #[test] fn majority_date32_field_returns_most_frequent_field() { @@ -122,10 +119,10 @@ mod tests { let registry = ExpressionRegistry::new(); let year = registry - .register(CacheExpression::extract_date32(Date32Field::Year)) + .register(CacheExpression::extract_date32(DatePart::Year)) .expect("register year"); let month = registry - .register(CacheExpression::extract_date32(Date32Field::Month)) + .register(CacheExpression::extract_date32(DatePart::Month)) .expect("register month"); tracker.record_expression(year); @@ -137,7 +134,7 @@ mod tests { let expr = registry.get(majority).expect("resolve expression"); assert_eq!( expr.as_ref(), - &CacheExpression::extract_date32(Date32Field::Year) + &CacheExpression::extract_date32(DatePart::Year) ); } @@ -147,10 +144,10 @@ mod tests { let registry = ExpressionRegistry::new(); let year = registry - .register(CacheExpression::extract_date32(Date32Field::Year)) + .register(CacheExpression::extract_date32(DatePart::Year)) .expect("register year"); let month = registry - .register(CacheExpression::extract_date32(Date32Field::Month)) + .register(CacheExpression::extract_date32(DatePart::Month)) .expect("register month"); tracker.record_expression(year); @@ -162,7 +159,7 @@ mod tests { let expr = registry.get(majority).expect("resolve expression"); assert_eq!( expr.as_ref(), - &CacheExpression::extract_date32(Date32Field::Month) + &CacheExpression::extract_date32(DatePart::Month) ); } }