From 72cb558394b51044e338b7daafc0cab32df4ffaf Mon Sep 17 00:00:00 2001 From: Sean <2579935313@qq.com> Date: Sat, 3 Oct 2026 13:33:42 +0800 Subject: [PATCH] feat(parquet): prune explicitly ordered INT96 timestamps --- .../core/tests/parquet/file_statistics.rs | 99 +++++ .../datasource-parquet/src/file_format.rs | 31 +- datafusion/datasource-parquet/src/metadata.rs | 83 +++- .../src/statistics_order_tests.rs | 399 +++++++++++++++++- 4 files changed, 579 insertions(+), 33 deletions(-) diff --git a/datafusion/core/tests/parquet/file_statistics.rs b/datafusion/core/tests/parquet/file_statistics.rs index f6d733ec69720..25afe94692e12 100644 --- a/datafusion/core/tests/parquet/file_statistics.rs +++ b/datafusion/core/tests/parquet/file_statistics.rs @@ -48,6 +48,105 @@ use datafusion_physical_plan::filter::FilterExec; use datafusion_physical_plan::statistics::{StatisticsArgs, StatisticsContext}; use tempfile::tempdir; +#[tokio::test] +async fn int96_format_statistics_use_configured_resolution() { + use arrow::datatypes::TimeUnit; + use datafusion::datasource::file_format::FileFormat; + use datafusion_common::ScalarValue; + use datafusion_common::config::TableParquetOptions; + use object_store::{ObjectStore, ObjectStoreExt, memory::InMemory, path::Path}; + use parquet::data_type::{Int96, Int96Type}; + use parquet::file::properties::WriterProperties; + use parquet::file::writer::SerializedFileWriter; + use parquet::schema::parser::parse_message_type; + + let values = [0_i64, 10_000_000_000].map(|seconds| { + let nanos = seconds.rem_euclid(86_400) as u64 * 1_000_000_000; + let mut value = Int96::new(); + value.set_data( + nanos as u32, + (nanos >> 32) as u32, + (seconds.div_euclid(86_400) + 2_440_588) as u32, + ); + value + }); + let parquet_schema = + Arc::new(parse_message_type("message test { REQUIRED INT96 ts; }").unwrap()); + let mut bytes = Vec::new(); + let mut writer = SerializedFileWriter::new( + &mut bytes, + parquet_schema, + Arc::new(WriterProperties::default()), + ) + .unwrap(); + let mut row_group = writer.next_row_group().unwrap(); + let mut column = row_group.next_column().unwrap().unwrap(); + column + .typed::() + .write_batch(&values, None, None) + .unwrap(); + column.close().unwrap(); + row_group.close().unwrap(); + writer.close().unwrap(); + + let store: Arc = Arc::new(InMemory::new()); + let path = Path::from("int96.parquet"); + store.put(&path, bytes.into()).await.unwrap(); + let object = store.head(&path).await.unwrap(); + let state = SessionContext::new().state(); + let mut options = TableParquetOptions::default(); + options.global.coerce_int96 = Some("us".into()); + options.global.coerce_int96_tz = Some("UTC".into()); + let configured = ParquetFormat::default().with_options(options); + let schema = configured + .infer_schema(&state, &store, std::slice::from_ref(&object)) + .await + .unwrap(); + assert_eq!( + schema.field(0).data_type(), + &DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())) + ); + + for (format, coerced) in [(configured, true), (ParquetFormat::default(), false)] { + let statistics = format + .infer_stats(&state, &store, schema.clone(), &object) + .await + .unwrap(); + let combined = format + .infer_stats_and_ordering(&state, &store, schema.clone(), &object) + .await + .unwrap(); + for statistics in [statistics, combined.statistics] { + assert_eq!(statistics.num_rows, Precision::Exact(2)); + assert_eq!( + statistics.column_statistics[0].null_count, + Precision::Exact(0) + ); + for (bound, value) in [ + (&statistics.column_statistics[0].min_value, 0), + ( + &statistics.column_statistics[0].max_value, + 10_000_000_000_000_000, + ), + ] { + assert_eq!( + bound, + &if coerced { + Precision::Exact(ScalarValue::TimestampMicrosecond( + Some(value), + Some("UTC".into()), + )) + } else { + // A microsecond table schema does not change a default + // reader's wrapping nanosecond conversion. + Precision::Absent + } + ); + } + } + } +} + #[tokio::test] async fn check_stats_precision_with_filter_pushdown() { let testdata = datafusion::test_util::parquet_test_data(); diff --git a/datafusion/datasource-parquet/src/file_format.rs b/datafusion/datasource-parquet/src/file_format.rs index dc0f9a7be434e..18ef696a124ae 100644 --- a/datafusion/datasource-parquet/src/file_format.rs +++ b/datafusion/datasource-parquet/src/file_format.rs @@ -437,6 +437,19 @@ impl FileFormat for ParquetFormat { .with_metadata_size_hint(self.metadata_size_hint()) .with_decryption_properties(file_decryption_properties) .with_file_metadata_cache(Some(file_metadata_cache)) + .with_coerce_int96( + self.coerce_int96() + .map(|unit| parse_coerce_int96_string(&unit)) + .transpose()?, + ) + .with_coerce_int96_tz( + self.options + .global + .coerce_int96_tz + .as_deref() + .map(parse_coerce_int96_tz_string) + .transpose()?, + ) .fetch_statistics(&table_schema) .await } @@ -480,10 +493,20 @@ impl FileFormat for ParquetFormat { .with_file_metadata_cache(Some(file_metadata_cache)) .fetch_metadata() .await?; - let statistics = DFParquetMetadata::statistics_from_parquet_metadata( - &metadata, - &table_schema, - )?; + let statistics = + DFParquetMetadata::statistics_from_parquet_metadata_with_coercion( + &metadata, + &table_schema, + self.coerce_int96() + .map(|unit| parse_coerce_int96_string(&unit)) + .transpose()?, + self.options + .global + .coerce_int96_tz + .as_deref() + .map(parse_coerce_int96_tz_string) + .transpose()?, + )?; let ordering = crate::metadata::ordering_from_parquet_metadata(&metadata, &table_schema)?; Ok( diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index 51dc98e4f2f02..85739e69b8fcf 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -71,24 +71,24 @@ fn requires_unsigned_byte_array_order(column: &ColumnDescriptor) -> bool { /// Arrow's string and binary comparisons. Even the modern bounds cannot be /// interpreted without the corresponding footer `column_orders` entry. /// Signed logical types, such as decimals, retain their existing behavior. -/// Columns with undefined sort orders, such as `INT96`, never have usable -/// min/max bounds regardless of their physical type. The `INT96` check is -/// defensive because parquet-rs does not currently expose those bounds. +/// Columns with undefined sort orders never have usable min/max bounds. +/// `INT96` timestamp bounds require an explicit `INT96_TIMESTAMP_ORDER` +/// footer entry; the schema's timestamp sort order alone is not sufficient. pub(crate) fn has_untrusted_min_max_order( parquet_schema: &SchemaDescriptor, column_orders: Option<&[ColumnOrder]>, parquet_column_index: usize, ) -> bool { let column = parquet_schema.column(parquet_column_index); - // As of arrow 60, INT96 columns report `SortOrder::INT96_TIMESTAMP` - // rather than `UNDEFINED`; keep treating their min/max as untrusted. - // until - if matches!( - column.sort_order(), - SortOrder::UNDEFINED | SortOrder::INT96_TIMESTAMP - ) { + if column.sort_order() == SortOrder::UNDEFINED { return true; } + if column.sort_order() == SortOrder::INT96_TIMESTAMP { + return column_orders + .and_then(|orders| orders.get(parquet_column_index)) + .copied() + != Some(ColumnOrder::INT96_TIMESTAMP_ORDER); + } requires_unsigned_byte_array_order(&column) && (column.sort_order() != SortOrder::UNSIGNED || column_orders @@ -222,7 +222,7 @@ impl<'a> DFParquetMetadata<'a> { } /// Set the [`TimeUnit`] that INT96 timestamp columns should be coerced - /// to when reading the schema. + /// to when reading the schema and statistics. /// /// INT96 in Parquet has no defined unit or timezone, so leaving this /// `None` reads INT96 columns as nanosecond timestamps with no timezone @@ -464,7 +464,12 @@ impl<'a> DFParquetMetadata<'a> { /// the statistics in the metadata using [`Self::statistics_from_parquet_metadata`] pub async fn fetch_statistics(&self, table_schema: &SchemaRef) -> Result { let metadata = self.fetch_metadata().await?; - Self::statistics_from_parquet_metadata(&metadata, table_schema) + Self::statistics_from_parquet_metadata_with_coercion( + &metadata, + table_schema, + self.coerce_int96, + self.coerce_int96_tz.clone(), + ) } /// Convert statistics in [`ParquetMetaData`] into [`Statistics`] using [`StatisticsConverter`] @@ -499,6 +504,10 @@ impl<'a> DFParquetMetadata<'a> { /// 2. The column is in arrow schema, but not in parquet schema due to schema revolution, min/max values are set to Precision::Exact(null) /// - Null counts are set to Precision::Exact(num_rows) (conservatively assuming all values could be null) /// + /// INT96 timestamps use the reader's default nanosecond resolution. Use + /// [`Self::fetch_statistics`] with [`Self::with_coerce_int96`] when the + /// reader is configured to decode INT96 at another resolution. + /// /// # Byte Size Calculation: /// /// - For primitive types with known fixed size, exact byte size is calculated as (byte width * number of rows) @@ -507,6 +516,21 @@ impl<'a> DFParquetMetadata<'a> { pub fn statistics_from_parquet_metadata( metadata: &ParquetMetaData, logical_file_schema: &SchemaRef, + ) -> Result { + Self::statistics_from_parquet_metadata_with_coercion( + metadata, + logical_file_schema, + None, + None, + ) + } + + /// Convert statistics using the same INT96 settings as the data reader. + pub(crate) fn statistics_from_parquet_metadata_with_coercion( + metadata: &ParquetMetaData, + logical_file_schema: &SchemaRef, + coerce_int96: Option, + coerce_int96_tz: Option>, ) -> Result { let row_groups_metadata = metadata.row_groups(); @@ -540,6 +564,21 @@ impl<'a> DFParquetMetadata<'a> { physical_file_schema = merged; } + // Match the reader's configured INT96 resolution, rather than assuming + // that the table's timestamp type is the type physically read. In + // particular, converting through nanoseconds can wrap wider timestamps. + if let Some(unit) = coerce_int96 + && let Some(coerced) = Int96Coercer::new( + file_metadata.schema_descr(), + &physical_file_schema, + &unit, + ) + .with_timezone(coerce_int96_tz) + .coerce() + { + physical_file_schema = coerced; + } + statistics.column_statistics = if has_statistics { let (mut max_accs, mut min_accs) = @@ -567,11 +606,21 @@ impl<'a> DFParquetMetadata<'a> { stats_converter.with_missing_null_counts_as_zero(false); let parquet_index = stats_converter.parquet_column_index(); if parquet_index.is_some_and(|index| { - has_untrusted_min_max_order( - file_metadata.schema_descr(), - file_metadata.column_orders().map(Vec::as_slice), - index, - ) + // A schema cast is not an INT96 read coercion. + // Keep bounds unknown if the actual reader unit + // differs from the table's timestamp type. + (file_metadata + .schema_descr() + .column(index) + .physical_type() + == PhysicalType::INT96 + && stats_converter.arrow_field().data_type() + != field.data_type()) + || has_untrusted_min_max_order( + file_metadata.schema_descr(), + file_metadata.column_orders().map(Vec::as_slice), + index, + ) }) || has_untrusted_byte_array_stats( file_metadata.schema_descr(), parquet_index, diff --git a/datafusion/datasource-parquet/src/statistics_order_tests.rs b/datafusion/datasource-parquet/src/statistics_order_tests.rs index 93e5f63ad2313..6d093260fe52a 100644 --- a/datafusion/datasource-parquet/src/statistics_order_tests.rs +++ b/datafusion/datasource-parquet/src/statistics_order_tests.rs @@ -15,13 +15,13 @@ // specific language governing permissions and limitations // under the License. -//! Regression tests for interpreting Parquet byte-array statistics orders. +//! Regression tests for interpreting Parquet statistics orders. use std::io::Write; use std::sync::Arc; use arrow::array::{BooleanArray, record_batch}; -use arrow::datatypes::{DataType, Field, Schema, SchemaRef}; +use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit}; use bytes::Bytes; use datafusion_common::pruning::{PrunableStatistics, PruningStatistics}; use datafusion_common::stats::Precision; @@ -32,9 +32,9 @@ use datafusion_physical_expr::planner::logical2physical; use datafusion_physical_plan::metrics::{Count, ExecutionPlanMetricsSet}; use datafusion_pruning::{MAX_IN_LIST_SIZE, PruningPredicate, PruningPredicateBuilder}; use parquet::arrow::ArrowWriter; -use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; +use parquet::arrow::arrow_reader::{ArrowReaderOptions, ParquetRecordBatchReaderBuilder}; use parquet::basic::{ColumnOrder, LogicalType, SortOrder, Type as PhysicalType}; -use parquet::data_type::{ByteArray, FixedLenByteArray}; +use parquet::data_type::{ByteArray, FixedLenByteArray, Int96, Int96Type}; use parquet::file::metadata::page_index::{PageIndex, PageIndexBuilder}; use parquet::file::metadata::{ ColumnChunkMetaData, ColumnIndexBuilder, FileMetaData, OffsetIndexBuilder, @@ -43,7 +43,8 @@ use parquet::file::metadata::{ }; use parquet::file::properties::{EnabledStatistics, WriterProperties}; use parquet::file::statistics::Statistics as ParquetStatistics; -use parquet::file::writer::TrackedWrite; +use parquet::file::writer::{SerializedFileWriter, TrackedWrite}; +use parquet::schema::parser::parse_message_type; use parquet::schema::types::{SchemaDescriptor, Type as ParquetType}; use crate::RowGroupAccessPlanFilter; @@ -64,6 +65,7 @@ struct TestFile { bytes: Bytes, schema: SchemaRef, metadata: Arc, + int96_coercion: Option<(TimeUnit, Option>)>, } impl TestFile { @@ -110,6 +112,7 @@ impl TestFile { bytes: original, schema, metadata: Arc::new(metadata), + int96_coercion: None, }; } @@ -215,6 +218,7 @@ impl TestFile { bytes, schema, metadata: Arc::new(metadata), + int96_coercion: None, } } @@ -228,8 +232,21 @@ impl TestFile { } fn statistics(&self) -> Statistics { - DFParquetMetadata::statistics_from_parquet_metadata(&self.metadata, &self.schema) + if let Some((unit, timezone)) = &self.int96_coercion { + DFParquetMetadata::statistics_from_parquet_metadata_with_coercion( + &self.metadata, + &self.schema, + Some(*unit), + timezone.clone(), + ) .unwrap() + } else { + DFParquetMetadata::statistics_from_parquet_metadata( + &self.metadata, + &self.schema, + ) + .unwrap() + } } fn file_matches(&self, predicate: &PruningPredicate) -> bool { @@ -278,10 +295,12 @@ impl TestFile { .row_groups .into_iter() .flat_map(|rg| { - let mut builder = - ParquetRecordBatchReaderBuilder::try_new(self.bytes.clone()) - .unwrap() - .with_row_groups(vec![rg.selection.row_group_index()]); + let mut builder = ParquetRecordBatchReaderBuilder::try_new_with_options( + self.bytes.clone(), + ArrowReaderOptions::new().with_schema(Arc::clone(&self.schema)), + ) + .unwrap() + .with_row_groups(vec![rg.selection.row_group_index()]); if let Some(selection) = rg.selection.selection() { builder = builder.with_row_selection(selection.clone()); } @@ -744,8 +763,357 @@ fn signed_decimal_byte_array_statistics_remain_usable() { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Int96Order { + Timestamp, + Missing, + TypeDefined, + Unknown, +} + +fn int96(seconds: i64) -> Int96 { + let day = 2_440_588 + seconds.div_euclid(86_400); + let nanos = seconds.rem_euclid(86_400) as u64 * 1_000_000_000; + let mut value = Int96::new(); + value.set_data(nanos as u32, (nanos >> 32) as u32, day as u32); + value +} + +impl TestFile { + fn int96( + order: Int96Order, + unit: TimeUnit, + timezone: Option>, + row_groups: &[Vec>], + ) -> Self { + let parquet_schema = + Arc::new(parse_message_type("message schema { optional int96 s; }").unwrap()); + let properties = Arc::new( + WriterProperties::builder() + .set_data_page_row_count_limit(2) + .set_write_batch_size(2) + .set_dictionary_enabled(false) + .set_statistics_enabled(EnabledStatistics::Page) + .build(), + ); + let mut bytes = Vec::new(); + let mut writer = + SerializedFileWriter::new(&mut bytes, parquet_schema, properties).unwrap(); + for values in row_groups { + let mut row_group = writer.next_row_group().unwrap(); + let mut column = row_group.next_column().unwrap().unwrap(); + let levels = values + .iter() + .map(|v| i16::from(v.is_some())) + .collect::>(); + let values = values.iter().flatten().copied().collect::>(); + column + .typed::() + .write_batch(&values, Some(&levels), None) + .unwrap(); + column.close().unwrap(); + row_group.close().unwrap(); + } + writer.close().unwrap(); + + // Replace only the footer's column_orders field. The real data and + // page indexes stay identical, so every order variant reads the same rows. + let end = bytes.len() - 8; + let encoded_orders = [0x19, 0x1c, 0x3c, 0, 0, 0]; + let start = end - encoded_orders.len(); + assert_eq!(&bytes[start..end], &encoded_orders); + if order == Int96Order::Missing { + let metadata_start = footer_start(&bytes); + bytes.drain(start..end - 1); + let end = bytes.len() - 8; + let metadata_len = (end - metadata_start) as u32; + bytes[end..end + 4].copy_from_slice(&metadata_len.to_le_bytes()); + } else if order == Int96Order::TypeDefined { + bytes[start + 2] = 0x1c; + } else if order == Int96Order::Unknown { + bytes[start + 2] = 0x4c; + } + let bytes = Bytes::from(bytes); + let metadata = read_metadata(&bytes); + let expected_order = match order { + Int96Order::Timestamp => ColumnOrder::INT96_TIMESTAMP_ORDER, + Int96Order::Missing => ColumnOrder::UNDEFINED, + Int96Order::TypeDefined => { + ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNDEFINED) + } + Int96Order::Unknown => ColumnOrder::UNKNOWN, + }; + assert_eq!(metadata.file_metadata().column_order(0), expected_order); + let schema = Schema::new(vec![Field::new( + "s", + DataType::Timestamp(TimeUnit::Nanosecond, None), + true, + )]); + let schema = crate::Int96Coercer::new( + metadata.file_metadata().schema_descr(), + &schema, + &unit, + ) + .with_timezone(timezone.clone()) + .coerce() + .unwrap(); + Self { + bytes, + schema: Arc::new(schema), + metadata: Arc::new(metadata), + int96_coercion: Some((unit, timezone)), + } + } + + fn timestamp(&self, seconds: i64) -> ScalarValue { + let DataType::Timestamp(unit, timezone) = self.schema.field(0).data_type() else { + unreachable!() + }; + match unit { + TimeUnit::Second => { + ScalarValue::TimestampSecond(Some(seconds), timezone.clone()) + } + TimeUnit::Millisecond => { + ScalarValue::TimestampMillisecond(Some(seconds * 1_000), timezone.clone()) + } + TimeUnit::Microsecond => ScalarValue::TimestampMicrosecond( + Some(seconds * 1_000_000), + timezone.clone(), + ), + TimeUnit::Nanosecond => ScalarValue::TimestampNanosecond( + Some(seconds * 1_000_000_000), + timezone.clone(), + ), + } + } +} + #[test] -fn undefined_int96_order_is_never_trusted() { +fn int96_timestamp_order_prunes_files_row_groups_and_pages() { + let row_groups = [ + vec![ + Some(int96(-86_400)), + Some(int96(-1)), + Some(int96(0)), + Some(int96(1)), + ], + vec![ + Some(int96(86_400)), + Some(int96(86_401)), + Some(int96(172_800)), + Some(int96(172_801)), + ], + vec![None; 4], + ]; + for unit in [ + TimeUnit::Second, + TimeUnit::Millisecond, + TimeUnit::Microsecond, + TimeUnit::Nanosecond, + ] { + for timezone in [None, Some(Arc::from("UTC"))] { + for order in [ + Int96Order::Timestamp, + Int96Order::Missing, + Int96Order::TypeDefined, + Int96Order::Unknown, + ] { + let file = TestFile::int96(order, unit, timezone.clone(), &row_groups); + let trusted = order == Int96Order::Timestamp; + let statistics = file.statistics(); + assert_eq!(statistics.num_rows, Precision::Exact(12)); + assert_eq!( + statistics.column_statistics[0].null_count, + Precision::Exact(4) + ); + assert_eq!( + statistics.column_statistics[0].min_value, + if trusted { + Precision::Exact(file.timestamp(-86_400)) + } else { + Precision::Absent + }, + "order={order:?}, unit={unit:?}, timezone={timezone:?}" + ); + assert_eq!( + statistics.column_statistics[0].max_value, + if trusted { + Precision::Exact(file.timestamp(172_801)) + } else { + Precision::Absent + } + ); + + let (_, absent) = + file.predicate(&col("s").eq(lit(file.timestamp(259_200)))); + assert_eq!(file.file_matches(&absent), !trusted); + let (physical, predicate) = + file.predicate(&col("s").eq(lit(file.timestamp(0)))); + let all = ParquetAccessPlan::new_all(file.metadata.num_row_groups()); + assert_eq!(file.matching_rows(&physical, all.clone()), 1); + assert!(file.file_matches(&predicate)); + let row_groups = file.row_group_plan(&predicate); + assert_eq!( + row_groups.row_group_indexes(), + if trusted { vec![0] } else { vec![0, 1] } + ); + assert_eq!(file.matching_rows(&physical, row_groups.clone()), 1); + + // Start page pruning from all row groups, independently of row-group pruning. + let page_metrics = metrics(); + let pages = + PagePruningAccessPlanFilter::new(&physical, Arc::clone(&file.schema)) + .prune_plan_with_page_index( + all, + &file.schema, + file.metadata.file_metadata().schema_descr(), + &file.metadata, + &page_metrics, + ); + assert_eq!( + page_metrics.page_index_rows_pruned.pruned(), + if trusted { 10 } else { 4 } + ); + assert_eq!(file.matching_rows(&physical, pages), 1); + assert_eq!( + file.matching_rows(&physical, file.page_plan(&physical, row_groups)), + 1 + ); + + let mut runtime_pruner = RowGroupPruner::new( + physical, + Arc::clone(&file.schema), + Arc::clone(&file.metadata), + Count::new(), + Count::new(), + MAX_IN_LIST_SIZE, + ); + assert!(!runtime_pruner.should_prune(&[0])); + assert_eq!(runtime_pruner.should_prune(&[1]), trusted); + assert!(runtime_pruner.should_prune(&[2])); + } + } + } +} + +#[test] +fn int96_overflow_keeps_wrapped_rows_and_coerced_bounds() { + // The second timestamp wraps to a negative nanosecond value when read. + // Neither endpoint is a safe bound for the resulting nanosecond interval. + let row_groups = [vec![Some(int96(0)), Some(int96(10_000_000_000))]]; + let file = TestFile::int96( + Int96Order::Timestamp, + TimeUnit::Nanosecond, + None, + &row_groups, + ); + let statistics = file.statistics(); + assert_eq!(statistics.column_statistics[0].min_value, Precision::Absent); + assert_eq!(statistics.column_statistics[0].max_value, Precision::Absent); + assert_eq!( + statistics.column_statistics[0].null_count, + Precision::Exact(0) + ); + let (physical, predicate) = file.predicate(&col("s").lt(lit(file.timestamp(0)))); + assert!(file.file_matches(&predicate)); + let row_groups = file.row_group_plan(&predicate); + assert_eq!(row_groups.row_group_indexes(), vec![0]); + assert_eq!( + file.matching_rows(&physical, file.page_plan(&physical, row_groups)), + 1 + ); + + for unit in [ + TimeUnit::Second, + TimeUnit::Millisecond, + TimeUnit::Microsecond, + ] { + let file = TestFile::int96( + Int96Order::Timestamp, + unit, + None, + &[vec![Some(int96(0)), Some(int96(10_000_000_000))]], + ); + let statistics = file.statistics(); + assert_eq!( + statistics.column_statistics[0].min_value, + Precision::Exact(file.timestamp(0)) + ); + assert_eq!( + statistics.column_statistics[0].max_value, + Precision::Exact(file.timestamp(10_000_000_000)) + ); + let uncoerced = DFParquetMetadata::statistics_from_parquet_metadata( + &file.metadata, + &file.schema, + ) + .unwrap(); + assert_eq!(uncoerced.column_statistics[0].min_value, Precision::Absent); + assert_eq!(uncoerced.column_statistics[0].max_value, Precision::Absent); + assert_eq!( + uncoerced.column_statistics[0].null_count, + Precision::Exact(0) + ); + let (_, predicate) = file.predicate(&col("s").lt(lit(file.timestamp(0)))); + assert!(!file.file_matches(&predicate)); + assert!( + file.row_group_plan(&predicate) + .row_group_indexes() + .is_empty() + ); + } +} + +#[tokio::test] +async fn int96_fetched_statistics_use_configured_resolution() { + use object_store::ObjectStoreExt; + use object_store::memory::InMemory; + use object_store::path::Path; + + let file = TestFile::int96( + Int96Order::Timestamp, + TimeUnit::Microsecond, + Some(Arc::from("UTC")), + &[vec![Some(int96(0)), Some(int96(10_000_000_000))]], + ); + let store = InMemory::new(); + let path = Path::from("int96.parquet"); + store.put(&path, file.bytes.clone().into()).await.unwrap(); + let object = store.head(&path).await.unwrap(); + let metadata = DFParquetMetadata::new(&store, &object) + .with_coerce_int96(Some(TimeUnit::Microsecond)) + .with_coerce_int96_tz(Some(Arc::from("UTC"))); + let schema = Arc::new(metadata.fetch_schema().await.unwrap()); + assert_eq!(schema, file.schema); + let statistics = metadata.fetch_statistics(&schema).await.unwrap(); + assert_eq!( + statistics.column_statistics[0].min_value, + Precision::Exact(file.timestamp(0)) + ); + assert_eq!( + statistics.column_statistics[0].max_value, + Precision::Exact(file.timestamp(10_000_000_000)) + ); + assert_eq!( + statistics.column_statistics[0].null_count, + Precision::Exact(0) + ); + + // A table schema alone does not opt in to direct microsecond decoding. + let statistics = DFParquetMetadata::new(&store, &object) + .fetch_statistics(&schema) + .await + .unwrap(); + assert_eq!(statistics.column_statistics[0].min_value, Precision::Absent); + assert_eq!(statistics.column_statistics[0].max_value, Precision::Absent); + assert_eq!( + statistics.column_statistics[0].null_count, + Precision::Exact(0) + ); +} + +#[test] +fn int96_requires_explicit_timestamp_order() { let parquet_type = ParquetType::primitive_type_builder("s", PhysicalType::INT96) .build() .unwrap(); @@ -764,7 +1132,8 @@ fn undefined_int96_order_is_never_trusted() { Some(ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNDEFINED)), Some(ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED)), Some(ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNSIGNED)), - Some(ColumnOrder::INT96_TIMESTAMP_ORDER), + Some(ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::INT96_TIMESTAMP)), + Some(ColumnOrder::IEEE_754_TOTAL_ORDER), ] { assert!( has_untrusted_min_max_order( @@ -775,6 +1144,12 @@ fn undefined_int96_order_is_never_trusted() { "order={order:?}", ); } + assert!(has_untrusted_min_max_order(&schema, Some(&[]), 0)); + assert!(!has_untrusted_min_max_order( + &schema, + Some(&[ColumnOrder::INT96_TIMESTAMP_ORDER]), + 0, + )); } #[test]