From ba6aca6b89026633082f1bb6525f3291c5b6bae8 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 11 Aug 2026 16:10:29 -0400 Subject: [PATCH 1/4] feat(parquet): add row_group_distinct_counts to StatisticsConverter Adds `StatisticsConverter::row_group_distinct_counts()` which reads the `distinct_count` field from parquet row group statistics and returns a `UInt64Array` with one entry per row group (null where the field is absent). This mirrors the existing `row_group_null_counts` and `row_group_nan_counts` methods and unblocks DataFusion from populating `ColumnStatistics::distinct_count` via `statistics_from_parquet_metadata`. Previously this was always `Precision::Absent` for parquet files regardless of whether the file contained the field. Non-arrow-rs writers (parquet-java, parquet-cpp) may already populate `distinct_count` in the file footer, so the read path is immediately useful without any write-side changes. --- parquet/src/arrow/arrow_reader/statistics.rs | 22 ++ parquet/tests/arrow_reader/statistics.rs | 203 ++++++++++++++++++- 2 files changed, 224 insertions(+), 1 deletion(-) diff --git a/parquet/src/arrow/arrow_reader/statistics.rs b/parquet/src/arrow/arrow_reader/statistics.rs index 7ff4a3a4e102..e01fe14e1af2 100644 --- a/parquet/src/arrow/arrow_reader/statistics.rs +++ b/parquet/src/arrow/arrow_reader/statistics.rs @@ -1821,6 +1821,28 @@ impl<'a> StatisticsConverter<'a> { Ok(UInt64Array::from_iter(nan_counts)) } + /// Extract the distinct counts from row group statistics in [`RowGroupMetaData`] + /// + /// See docs on [`Self::row_group_mins`] for details + pub fn row_group_distinct_counts(&self, metadatas: I) -> Result + where + I: IntoIterator, + { + let Some(parquet_index) = self.parquet_column_index else { + let num_row_groups = metadatas.into_iter().count(); + return Ok(UInt64Array::from_iter(std::iter::repeat_n( + None, + num_row_groups, + ))); + }; + + let distinct_counts = metadatas + .into_iter() + .map(|x| x.column(parquet_index).statistics()) + .map(|s| s.and_then(|s| s.distinct_count_opt())); + Ok(UInt64Array::from_iter(distinct_counts)) + } + /// Extract the minimum values from Data Page statistics. /// /// In Parquet files, in addition to the Column Chunk level statistics diff --git a/parquet/tests/arrow_reader/statistics.rs b/parquet/tests/arrow_reader/statistics.rs index 0bc61ba108e5..93b6a782d51f 100644 --- a/parquet/tests/arrow_reader/statistics.rs +++ b/parquet/tests/arrow_reader/statistics.rs @@ -2630,10 +2630,13 @@ mod test { Int32Array, Int64Array, RecordBatch, StringArray, StructArray, TimestampNanosecondArray, new_empty_array, }; - use arrow_schema::{DataType, SchemaRef, TimeUnit}; + use arrow_schema::{DataType, Field, SchemaRef, TimeUnit}; use bytes::Bytes; use parquet::arrow::parquet_column; + use parquet::data_type::{ByteArray, ByteArrayType, Int32Type}; use parquet::file::metadata::{ParquetMetaData, RowGroupMetaData}; + use parquet::file::writer::SerializedFileWriter; + use parquet::schema::parser::parse_message_type; use std::path::PathBuf; use std::sync::Arc; // TODO error cases (with parquet statistics that are mismatched in expected type) @@ -3222,4 +3225,202 @@ mod test { .collect(); Arc::new(array) } + + // Roundtrip test: write 10 unique UTF-8 strings to a single row group with + // distinct_count=10 stored in the statistics, then read it back via + // StatisticsConverter::row_group_distinct_counts and assert we get 10. + // Uses the low-level writer because ArrowWriter does not yet write distinct_count. + #[test] + fn test_row_group_distinct_counts_utf8_roundtrip() { + let unique_string_values: Vec = (0..10) + .map(|index| ByteArray::from(format!("value_{index}").into_bytes())) + .collect(); + + let parquet_schema = Arc::new( + parse_message_type("message schema { REQUIRED BYTE_ARRAY col (UTF8); }").unwrap(), + ); + let writer_properties = Arc::new( + WriterProperties::builder() + .set_statistics_enabled(EnabledStatistics::Chunk) + .build(), + ); + + let mut file_buffer: Vec = Vec::new(); + let mut file_writer = + SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties) + .unwrap(); + + let mut row_group_writer = file_writer.next_row_group().unwrap(); + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + let min_value = unique_string_values.first().unwrap().clone(); + let max_value = unique_string_values.last().unwrap().clone(); + let expected_distinct_count = unique_string_values.len() as u64; + column_writer + .typed::() + .write_batch_with_statistics( + &unique_string_values, + None, + None, + Some(&min_value), + Some(&max_value), + Some(expected_distinct_count), + ) + .unwrap(); + column_writer.close().unwrap(); + row_group_writer.close().unwrap(); + file_writer.close().unwrap(); + + let parquet_bytes = Bytes::from(file_buffer); + let reader_builder = ParquetRecordBatchReaderBuilder::try_new(parquet_bytes).unwrap(); + let file_metadata = reader_builder.metadata().clone(); + let arrow_schema = reader_builder.schema().clone(); + let parquet_schema_descriptor = file_metadata.file_metadata().schema_descr(); + + let statistics_converter = + StatisticsConverter::try_new("col", &arrow_schema, parquet_schema_descriptor).unwrap(); + let distinct_counts = statistics_converter + .row_group_distinct_counts(file_metadata.row_groups().iter()) + .unwrap(); + + assert_eq!( + distinct_counts, + UInt64Array::from(vec![Some(expected_distinct_count)]), + "expected distinct_count of 10 unique string values" + ); + } + + #[test] + fn test_row_group_distinct_counts() { + let parquet_schema = Arc::new( + parse_message_type("message schema { REQUIRED INT32 col; }").unwrap(), + ); + let writer_properties = Arc::new( + WriterProperties::builder() + .set_statistics_enabled(EnabledStatistics::Chunk) + .build(), + ); + + let mut file_buffer: Vec = Vec::new(); + let mut file_writer = + SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties) + .unwrap(); + + // row group 0: 4 distinct values + let mut row_group_writer = file_writer.next_row_group().unwrap(); + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch_with_statistics(&[1, 2, 3, 4], None, None, Some(&1), Some(&4), Some(4)) + .unwrap(); + column_writer.close().unwrap(); + row_group_writer.close().unwrap(); + + // row group 1: 2 distinct values + let mut row_group_writer = file_writer.next_row_group().unwrap(); + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch_with_statistics( + &[10, 10, 20, 20], + None, + None, + Some(&10), + Some(&20), + Some(2), + ) + .unwrap(); + column_writer.close().unwrap(); + row_group_writer.close().unwrap(); + + file_writer.close().unwrap(); + + let parquet_bytes = Bytes::from(file_buffer); + let reader_builder = ParquetRecordBatchReaderBuilder::try_new(parquet_bytes).unwrap(); + let file_metadata = reader_builder.metadata().clone(); + let arrow_schema = reader_builder.schema().clone(); + let parquet_schema_descriptor = file_metadata.file_metadata().schema_descr(); + + let statistics_converter = + StatisticsConverter::try_new("col", &arrow_schema, parquet_schema_descriptor).unwrap(); + let distinct_counts = statistics_converter + .row_group_distinct_counts(file_metadata.row_groups().iter()) + .unwrap(); + + assert_eq!( + distinct_counts, + UInt64Array::from(vec![Some(4), Some(2)]), + "distinct counts should match the values injected via write_batch_with_statistics" + ); + } + + // Verifies that iteration continues across all row groups even when one is + // missing distinct_count. Row group 1 has no distinct_count written, but + // row groups 0 and 2 do — the result should be [Some(2), null, Some(5)], + // not a short-circuit to all-nulls. + #[test] + fn test_row_group_distinct_counts_absent() { + let parquet_schema = Arc::new( + parse_message_type("message schema { REQUIRED INT32 col; }").unwrap(), + ); + let writer_properties = Arc::new( + WriterProperties::builder() + .set_statistics_enabled(EnabledStatistics::Chunk) + .build(), + ); + + let mut file_buffer: Vec = Vec::new(); + let mut file_writer = + SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties) + .unwrap(); + + // row group 0: distinct_count present + let mut row_group_writer = file_writer.next_row_group().unwrap(); + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch_with_statistics(&[1, 2], None, None, Some(&1), Some(&2), Some(2)) + .unwrap(); + column_writer.close().unwrap(); + row_group_writer.close().unwrap(); + + // row group 1: distinct_count absent — should appear as null in output + let mut row_group_writer = file_writer.next_row_group().unwrap(); + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch_with_statistics(&[3, 4], None, None, Some(&3), Some(&4), None) + .unwrap(); + column_writer.close().unwrap(); + row_group_writer.close().unwrap(); + + // row group 2: distinct_count present — iteration must reach here despite row group 1 + let mut row_group_writer = file_writer.next_row_group().unwrap(); + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch_with_statistics(&[5, 6, 7, 8, 9], None, None, Some(&5), Some(&9), Some(5)) + .unwrap(); + column_writer.close().unwrap(); + row_group_writer.close().unwrap(); + + file_writer.close().unwrap(); + + let parquet_bytes = Bytes::from(file_buffer); + let reader_builder = ParquetRecordBatchReaderBuilder::try_new(parquet_bytes).unwrap(); + let file_metadata = reader_builder.metadata().clone(); + let arrow_schema = reader_builder.schema().clone(); + let parquet_schema_descriptor = file_metadata.file_metadata().schema_descr(); + + let statistics_converter = + StatisticsConverter::try_new("col", &arrow_schema, parquet_schema_descriptor).unwrap(); + let distinct_counts = statistics_converter + .row_group_distinct_counts(file_metadata.row_groups().iter()) + .unwrap(); + + assert_eq!( + distinct_counts, + UInt64Array::from(vec![Some(2), None, Some(5)]), + "a missing distinct_count in one row group should not affect the others" + ); + } } From 22ae8e9aed09933a1454f37c883b6e9159091125 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 11 Aug 2026 16:14:52 -0400 Subject: [PATCH 2/4] trim --- parquet/tests/arrow_reader/statistics.rs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/parquet/tests/arrow_reader/statistics.rs b/parquet/tests/arrow_reader/statistics.rs index 93b6a782d51f..549863550d66 100644 --- a/parquet/tests/arrow_reader/statistics.rs +++ b/parquet/tests/arrow_reader/statistics.rs @@ -3226,10 +3226,9 @@ mod test { Arc::new(array) } - // Roundtrip test: write 10 unique UTF-8 strings to a single row group with - // distinct_count=10 stored in the statistics, then read it back via - // StatisticsConverter::row_group_distinct_counts and assert we get 10. - // Uses the low-level writer because ArrowWriter does not yet write distinct_count. + // Verifies that distinct_count is correctly read back from UTF-8 column statistics. + // Uses the low-level writer to inject a known distinct_count into the parquet footer + // since ArrowWriter does not yet write this field. #[test] fn test_row_group_distinct_counts_utf8_roundtrip() { let unique_string_values: Vec = (0..10) @@ -3353,10 +3352,7 @@ mod test { ); } - // Verifies that iteration continues across all row groups even when one is - // missing distinct_count. Row group 1 has no distinct_count written, but - // row groups 0 and 2 do — the result should be [Some(2), null, Some(5)], - // not a short-circuit to all-nulls. + // Verifies that a missing distinct_count in one row group does not affect the others. #[test] fn test_row_group_distinct_counts_absent() { let parquet_schema = Arc::new( From 14c12473f45d0c2779f3aa012640402a4930a685 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 11 Aug 2026 16:26:03 -0400 Subject: [PATCH 3/4] fix lint --- parquet/tests/arrow_reader/statistics.rs | 19 +++++++------------ 1 file changed, 7 insertions(+), 12 deletions(-) diff --git a/parquet/tests/arrow_reader/statistics.rs b/parquet/tests/arrow_reader/statistics.rs index 549863550d66..e5d3d73799cc 100644 --- a/parquet/tests/arrow_reader/statistics.rs +++ b/parquet/tests/arrow_reader/statistics.rs @@ -3246,8 +3246,7 @@ mod test { let mut file_buffer: Vec = Vec::new(); let mut file_writer = - SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties) - .unwrap(); + SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties).unwrap(); let mut row_group_writer = file_writer.next_row_group().unwrap(); let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); @@ -3290,9 +3289,8 @@ mod test { #[test] fn test_row_group_distinct_counts() { - let parquet_schema = Arc::new( - parse_message_type("message schema { REQUIRED INT32 col; }").unwrap(), - ); + let parquet_schema = + Arc::new(parse_message_type("message schema { REQUIRED INT32 col; }").unwrap()); let writer_properties = Arc::new( WriterProperties::builder() .set_statistics_enabled(EnabledStatistics::Chunk) @@ -3301,8 +3299,7 @@ mod test { let mut file_buffer: Vec = Vec::new(); let mut file_writer = - SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties) - .unwrap(); + SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties).unwrap(); // row group 0: 4 distinct values let mut row_group_writer = file_writer.next_row_group().unwrap(); @@ -3355,9 +3352,8 @@ mod test { // Verifies that a missing distinct_count in one row group does not affect the others. #[test] fn test_row_group_distinct_counts_absent() { - let parquet_schema = Arc::new( - parse_message_type("message schema { REQUIRED INT32 col; }").unwrap(), - ); + let parquet_schema = + Arc::new(parse_message_type("message schema { REQUIRED INT32 col; }").unwrap()); let writer_properties = Arc::new( WriterProperties::builder() .set_statistics_enabled(EnabledStatistics::Chunk) @@ -3366,8 +3362,7 @@ mod test { let mut file_buffer: Vec = Vec::new(); let mut file_writer = - SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties) - .unwrap(); + SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties).unwrap(); // row group 0: distinct_count present let mut row_group_writer = file_writer.next_row_group().unwrap(); From 2d5b7a32817c2b7a8f3ad198349cab64669d92cd Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 11 Aug 2026 16:30:38 -0400 Subject: [PATCH 4/4] trim test --- parquet/tests/arrow_reader/statistics.rs | 62 ------------------------ 1 file changed, 62 deletions(-) diff --git a/parquet/tests/arrow_reader/statistics.rs b/parquet/tests/arrow_reader/statistics.rs index e5d3d73799cc..e0740d38e93e 100644 --- a/parquet/tests/arrow_reader/statistics.rs +++ b/parquet/tests/arrow_reader/statistics.rs @@ -3287,68 +3287,6 @@ mod test { ); } - #[test] - fn test_row_group_distinct_counts() { - let parquet_schema = - Arc::new(parse_message_type("message schema { REQUIRED INT32 col; }").unwrap()); - let writer_properties = Arc::new( - WriterProperties::builder() - .set_statistics_enabled(EnabledStatistics::Chunk) - .build(), - ); - - let mut file_buffer: Vec = Vec::new(); - let mut file_writer = - SerializedFileWriter::new(&mut file_buffer, parquet_schema, writer_properties).unwrap(); - - // row group 0: 4 distinct values - let mut row_group_writer = file_writer.next_row_group().unwrap(); - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch_with_statistics(&[1, 2, 3, 4], None, None, Some(&1), Some(&4), Some(4)) - .unwrap(); - column_writer.close().unwrap(); - row_group_writer.close().unwrap(); - - // row group 1: 2 distinct values - let mut row_group_writer = file_writer.next_row_group().unwrap(); - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch_with_statistics( - &[10, 10, 20, 20], - None, - None, - Some(&10), - Some(&20), - Some(2), - ) - .unwrap(); - column_writer.close().unwrap(); - row_group_writer.close().unwrap(); - - file_writer.close().unwrap(); - - let parquet_bytes = Bytes::from(file_buffer); - let reader_builder = ParquetRecordBatchReaderBuilder::try_new(parquet_bytes).unwrap(); - let file_metadata = reader_builder.metadata().clone(); - let arrow_schema = reader_builder.schema().clone(); - let parquet_schema_descriptor = file_metadata.file_metadata().schema_descr(); - - let statistics_converter = - StatisticsConverter::try_new("col", &arrow_schema, parquet_schema_descriptor).unwrap(); - let distinct_counts = statistics_converter - .row_group_distinct_counts(file_metadata.row_groups().iter()) - .unwrap(); - - assert_eq!( - distinct_counts, - UInt64Array::from(vec![Some(4), Some(2)]), - "distinct counts should match the values injected via write_batch_with_statistics" - ); - } - // Verifies that a missing distinct_count in one row group does not affect the others. #[test] fn test_row_group_distinct_counts_absent() {