Skip to content

Commit 72cb558

Browse files
committed
feat(parquet): prune explicitly ordered INT96 timestamps
1 parent 801b017 commit 72cb558

4 files changed

Lines changed: 579 additions & 33 deletions

File tree

‎datafusion/core/tests/parquet/file_statistics.rs‎

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,105 @@ use datafusion_physical_plan::filter::FilterExec;
4848
use datafusion_physical_plan::statistics::{StatisticsArgs, StatisticsContext};
4949
use tempfile::tempdir;
5050

51+
#[tokio::test]
52+
async fn int96_format_statistics_use_configured_resolution() {
53+
use arrow::datatypes::TimeUnit;
54+
use datafusion::datasource::file_format::FileFormat;
55+
use datafusion_common::ScalarValue;
56+
use datafusion_common::config::TableParquetOptions;
57+
use object_store::{ObjectStore, ObjectStoreExt, memory::InMemory, path::Path};
58+
use parquet::data_type::{Int96, Int96Type};
59+
use parquet::file::properties::WriterProperties;
60+
use parquet::file::writer::SerializedFileWriter;
61+
use parquet::schema::parser::parse_message_type;
62+
63+
let values = [0_i64, 10_000_000_000].map(|seconds| {
64+
let nanos = seconds.rem_euclid(86_400) as u64 * 1_000_000_000;
65+
let mut value = Int96::new();
66+
value.set_data(
67+
nanos as u32,
68+
(nanos >> 32) as u32,
69+
(seconds.div_euclid(86_400) + 2_440_588) as u32,
70+
);
71+
value
72+
});
73+
let parquet_schema =
74+
Arc::new(parse_message_type("message test { REQUIRED INT96 ts; }").unwrap());
75+
let mut bytes = Vec::new();
76+
let mut writer = SerializedFileWriter::new(
77+
&mut bytes,
78+
parquet_schema,
79+
Arc::new(WriterProperties::default()),
80+
)
81+
.unwrap();
82+
let mut row_group = writer.next_row_group().unwrap();
83+
let mut column = row_group.next_column().unwrap().unwrap();
84+
column
85+
.typed::<Int96Type>()
86+
.write_batch(&values, None, None)
87+
.unwrap();
88+
column.close().unwrap();
89+
row_group.close().unwrap();
90+
writer.close().unwrap();
91+
92+
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
93+
let path = Path::from("int96.parquet");
94+
store.put(&path, bytes.into()).await.unwrap();
95+
let object = store.head(&path).await.unwrap();
96+
let state = SessionContext::new().state();
97+
let mut options = TableParquetOptions::default();
98+
options.global.coerce_int96 = Some("us".into());
99+
options.global.coerce_int96_tz = Some("UTC".into());
100+
let configured = ParquetFormat::default().with_options(options);
101+
let schema = configured
102+
.infer_schema(&state, &store, std::slice::from_ref(&object))
103+
.await
104+
.unwrap();
105+
assert_eq!(
106+
schema.field(0).data_type(),
107+
&DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into()))
108+
);
109+
110+
for (format, coerced) in [(configured, true), (ParquetFormat::default(), false)] {
111+
let statistics = format
112+
.infer_stats(&state, &store, schema.clone(), &object)
113+
.await
114+
.unwrap();
115+
let combined = format
116+
.infer_stats_and_ordering(&state, &store, schema.clone(), &object)
117+
.await
118+
.unwrap();
119+
for statistics in [statistics, combined.statistics] {
120+
assert_eq!(statistics.num_rows, Precision::Exact(2));
121+
assert_eq!(
122+
statistics.column_statistics[0].null_count,
123+
Precision::Exact(0)
124+
);
125+
for (bound, value) in [
126+
(&statistics.column_statistics[0].min_value, 0),
127+
(
128+
&statistics.column_statistics[0].max_value,
129+
10_000_000_000_000_000,
130+
),
131+
] {
132+
assert_eq!(
133+
bound,
134+
&if coerced {
135+
Precision::Exact(ScalarValue::TimestampMicrosecond(
136+
Some(value),
137+
Some("UTC".into()),
138+
))
139+
} else {
140+
// A microsecond table schema does not change a default
141+
// reader's wrapping nanosecond conversion.
142+
Precision::Absent
143+
}
144+
);
145+
}
146+
}
147+
}
148+
}
149+
51150
#[tokio::test]
52151
async fn check_stats_precision_with_filter_pushdown() {
53152
let testdata = datafusion::test_util::parquet_test_data();

‎datafusion/datasource-parquet/src/file_format.rs‎

Lines changed: 27 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -437,6 +437,19 @@ impl FileFormat for ParquetFormat {
437437
.with_metadata_size_hint(self.metadata_size_hint())
438438
.with_decryption_properties(file_decryption_properties)
439439
.with_file_metadata_cache(Some(file_metadata_cache))
440+
.with_coerce_int96(
441+
self.coerce_int96()
442+
.map(|unit| parse_coerce_int96_string(&unit))
443+
.transpose()?,
444+
)
445+
.with_coerce_int96_tz(
446+
self.options
447+
.global
448+
.coerce_int96_tz
449+
.as_deref()
450+
.map(parse_coerce_int96_tz_string)
451+
.transpose()?,
452+
)
440453
.fetch_statistics(&table_schema)
441454
.await
442455
}
@@ -480,10 +493,20 @@ impl FileFormat for ParquetFormat {
480493
.with_file_metadata_cache(Some(file_metadata_cache))
481494
.fetch_metadata()
482495
.await?;
483-
let statistics = DFParquetMetadata::statistics_from_parquet_metadata(
484-
&metadata,
485-
&table_schema,
486-
)?;
496+
let statistics =
497+
DFParquetMetadata::statistics_from_parquet_metadata_with_coercion(
498+
&metadata,
499+
&table_schema,
500+
self.coerce_int96()
501+
.map(|unit| parse_coerce_int96_string(&unit))
502+
.transpose()?,
503+
self.options
504+
.global
505+
.coerce_int96_tz
506+
.as_deref()
507+
.map(parse_coerce_int96_tz_string)
508+
.transpose()?,
509+
)?;
487510
let ordering =
488511
crate::metadata::ordering_from_parquet_metadata(&metadata, &table_schema)?;
489512
Ok(

‎datafusion/datasource-parquet/src/metadata.rs‎

Lines changed: 66 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -71,24 +71,24 @@ fn requires_unsigned_byte_array_order(column: &ColumnDescriptor) -> bool {
7171
/// Arrow's string and binary comparisons. Even the modern bounds cannot be
7272
/// interpreted without the corresponding footer `column_orders` entry.
7373
/// Signed logical types, such as decimals, retain their existing behavior.
74-
/// Columns with undefined sort orders, such as `INT96`, never have usable
75-
/// min/max bounds regardless of their physical type. The `INT96` check is
76-
/// defensive because parquet-rs does not currently expose those bounds.
74+
/// Columns with undefined sort orders never have usable min/max bounds.
75+
/// `INT96` timestamp bounds require an explicit `INT96_TIMESTAMP_ORDER`
76+
/// footer entry; the schema's timestamp sort order alone is not sufficient.
7777
pub(crate) fn has_untrusted_min_max_order(
7878
parquet_schema: &SchemaDescriptor,
7979
column_orders: Option<&[ColumnOrder]>,
8080
parquet_column_index: usize,
8181
) -> bool {
8282
let column = parquet_schema.column(parquet_column_index);
83-
// As of arrow 60, INT96 columns report `SortOrder::INT96_TIMESTAMP`
84-
// rather than `UNDEFINED`; keep treating their min/max as untrusted.
85-
// until <https://github.com/apache/datafusion/issues/25484>
86-
if matches!(
87-
column.sort_order(),
88-
SortOrder::UNDEFINED | SortOrder::INT96_TIMESTAMP
89-
) {
83+
if column.sort_order() == SortOrder::UNDEFINED {
9084
return true;
9185
}
86+
if column.sort_order() == SortOrder::INT96_TIMESTAMP {
87+
return column_orders
88+
.and_then(|orders| orders.get(parquet_column_index))
89+
.copied()
90+
!= Some(ColumnOrder::INT96_TIMESTAMP_ORDER);
91+
}
9292
requires_unsigned_byte_array_order(&column)
9393
&& (column.sort_order() != SortOrder::UNSIGNED
9494
|| column_orders
@@ -222,7 +222,7 @@ impl<'a> DFParquetMetadata<'a> {
222222
}
223223

224224
/// Set the [`TimeUnit`] that INT96 timestamp columns should be coerced
225-
/// to when reading the schema.
225+
/// to when reading the schema and statistics.
226226
///
227227
/// INT96 in Parquet has no defined unit or timezone, so leaving this
228228
/// `None` reads INT96 columns as nanosecond timestamps with no timezone
@@ -464,7 +464,12 @@ impl<'a> DFParquetMetadata<'a> {
464464
/// the statistics in the metadata using [`Self::statistics_from_parquet_metadata`]
465465
pub async fn fetch_statistics(&self, table_schema: &SchemaRef) -> Result<Statistics> {
466466
let metadata = self.fetch_metadata().await?;
467-
Self::statistics_from_parquet_metadata(&metadata, table_schema)
467+
Self::statistics_from_parquet_metadata_with_coercion(
468+
&metadata,
469+
table_schema,
470+
self.coerce_int96,
471+
self.coerce_int96_tz.clone(),
472+
)
468473
}
469474

470475
/// Convert statistics in [`ParquetMetaData`] into [`Statistics`] using [`StatisticsConverter`]
@@ -499,6 +504,10 @@ impl<'a> DFParquetMetadata<'a> {
499504
/// 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)
500505
/// - Null counts are set to Precision::Exact(num_rows) (conservatively assuming all values could be null)
501506
///
507+
/// INT96 timestamps use the reader's default nanosecond resolution. Use
508+
/// [`Self::fetch_statistics`] with [`Self::with_coerce_int96`] when the
509+
/// reader is configured to decode INT96 at another resolution.
510+
///
502511
/// # Byte Size Calculation:
503512
///
504513
/// - 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> {
507516
pub fn statistics_from_parquet_metadata(
508517
metadata: &ParquetMetaData,
509518
logical_file_schema: &SchemaRef,
519+
) -> Result<Statistics> {
520+
Self::statistics_from_parquet_metadata_with_coercion(
521+
metadata,
522+
logical_file_schema,
523+
None,
524+
None,
525+
)
526+
}
527+
528+
/// Convert statistics using the same INT96 settings as the data reader.
529+
pub(crate) fn statistics_from_parquet_metadata_with_coercion(
530+
metadata: &ParquetMetaData,
531+
logical_file_schema: &SchemaRef,
532+
coerce_int96: Option<TimeUnit>,
533+
coerce_int96_tz: Option<Arc<str>>,
510534
) -> Result<Statistics> {
511535
let row_groups_metadata = metadata.row_groups();
512536

@@ -540,6 +564,21 @@ impl<'a> DFParquetMetadata<'a> {
540564
physical_file_schema = merged;
541565
}
542566

567+
// Match the reader's configured INT96 resolution, rather than assuming
568+
// that the table's timestamp type is the type physically read. In
569+
// particular, converting through nanoseconds can wrap wider timestamps.
570+
if let Some(unit) = coerce_int96
571+
&& let Some(coerced) = Int96Coercer::new(
572+
file_metadata.schema_descr(),
573+
&physical_file_schema,
574+
&unit,
575+
)
576+
.with_timezone(coerce_int96_tz)
577+
.coerce()
578+
{
579+
physical_file_schema = coerced;
580+
}
581+
543582
statistics.column_statistics =
544583
if has_statistics {
545584
let (mut max_accs, mut min_accs) =
@@ -567,11 +606,21 @@ impl<'a> DFParquetMetadata<'a> {
567606
stats_converter.with_missing_null_counts_as_zero(false);
568607
let parquet_index = stats_converter.parquet_column_index();
569608
if parquet_index.is_some_and(|index| {
570-
has_untrusted_min_max_order(
571-
file_metadata.schema_descr(),
572-
file_metadata.column_orders().map(Vec::as_slice),
573-
index,
574-
)
609+
// A schema cast is not an INT96 read coercion.
610+
// Keep bounds unknown if the actual reader unit
611+
// differs from the table's timestamp type.
612+
(file_metadata
613+
.schema_descr()
614+
.column(index)
615+
.physical_type()
616+
== PhysicalType::INT96
617+
&& stats_converter.arrow_field().data_type()
618+
!= field.data_type())
619+
|| has_untrusted_min_max_order(
620+
file_metadata.schema_descr(),
621+
file_metadata.column_orders().map(Vec::as_slice),
622+
index,
623+
)
575624
}) || has_untrusted_byte_array_stats(
576625
file_metadata.schema_descr(),
577626
parquet_index,

0 commit comments

Comments
 (0)