Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
99 changes: 99 additions & 0 deletions datafusion/core/tests/parquet/file_statistics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<Int96Type>()
.write_batch(&values, None, None)
.unwrap();
column.close().unwrap();
row_group.close().unwrap();
writer.close().unwrap();

let store: Arc<dyn ObjectStore> = 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();
Expand Down
31 changes: 27 additions & 4 deletions datafusion/datasource-parquet/src/file_format.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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(
Expand Down
83 changes: 66 additions & 17 deletions datafusion/datasource-parquet/src/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 <https://github.com/apache/datafusion/issues/25484>
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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<Statistics> {
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`]
Expand Down Expand Up @@ -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)
Expand All @@ -507,6 +516,21 @@ impl<'a> DFParquetMetadata<'a> {
pub fn statistics_from_parquet_metadata(
metadata: &ParquetMetaData,
logical_file_schema: &SchemaRef,
) -> Result<Statistics> {
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<TimeUnit>,
coerce_int96_tz: Option<Arc<str>>,
) -> Result<Statistics> {
let row_groups_metadata = metadata.row_groups();

Expand Down Expand Up @@ -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) =
Expand Down Expand Up @@ -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,
Expand Down
Loading