This is an automated email from the ASF dual-hosted git repository.
etseidl pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new a01d14f744 fix(parquet): preserve distinct-count invalidation across
page flushes (#11264)
a01d14f744 is described below
commit a01d14f74402607aac0afe73a3bc69801fbe0158
Author: Hippolyte Barraud <[email protected]>
AuthorDate: Mon Sep 28 19:33:24 2026 -0400
fix(parquet): preserve distinct-count invalidation across page flushes
(#11264)
# Which issue does this PR close?
- Closes #11263.
# Rationale for this change
Caller-supplied distinct counts describe an individual write, while
column metadata describes the entire column chunk. Once earlier data
makes the column-wide count unknown, a later write cannot safely restore
that count using only its own batch statistics.
# What changes are included in this PR?
- Use column and page metrics to detect prior data instead of the
page-local encoder value count.
- Accept a caller-supplied distinct count only when there is no prior
data in the column chunk.
- Update the regression test to force a page boundary and verify that
the column-wide distinct count remains absent while page statistics
remain correct.
# Are these changes tested?
Yes. The updated regression covers distinct-count invalidation across a
page flush and checks that page-level statistics are preserved.
# Are there any user-facing changes?
There are no public API changes.
Signed-off-by: Hippolyte Barraud <[email protected]>
---
parquet/src/column/writer/mod.rs | 34 +++++++++++++++-------------------
1 file changed, 15 insertions(+), 19 deletions(-)
diff --git a/parquet/src/column/writer/mod.rs b/parquet/src/column/writer/mod.rs
index 18def152f8..c0eef944d0 100644
--- a/parquet/src/column/writer/mod.rs
+++ b/parquet/src/column/writer/mod.rs
@@ -617,12 +617,11 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a,
E> {
update_max(&self.descr, max, &mut
self.column_metrics.max_column_value);
}
- // We can only set the distinct count if there are no other writes
- if self.encoder.num_values() == 0 {
- self.column_metrics.column_distinct_count = distinct_count;
- } else {
- self.column_metrics.column_distinct_count = None;
- }
+ // Encoder counts reset per page; row metrics retain column-wide
history.
+ let has_prior_data = self.column_metrics.total_rows_written != 0
+ || self.page_metrics.num_buffered_values != 0;
+ self.column_metrics.column_distinct_count =
+ if has_prior_data { None } else { distinct_count };
let mut values_offset = 0;
let mut levels_offset = 0;
@@ -2817,6 +2816,7 @@ mod tests {
let props = Arc::new(
WriterProperties::builder()
.set_write_page_header_statistics(true)
+ .set_data_page_row_count_limit(4)
.build(),
);
let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0,
0, props);
@@ -2850,22 +2850,18 @@ mod tests {
.unwrap();
let pages = reader.collect::<Result<Vec<_>>>().unwrap();
- assert_eq!(pages.len(), 2);
+ assert_eq!(pages.len(), 3);
assert_eq!(pages[0].page_type(), PageType::DICTIONARY_PAGE);
assert_eq!(pages[1].page_type(), PageType::DATA_PAGE);
-
- let page_statistics = pages[1].statistics().unwrap();
- assert_eq!(
- page_statistics.min_bytes_opt().unwrap(),
- 1_i32.to_le_bytes()
- );
- assert_eq!(
- page_statistics.max_bytes_opt().unwrap(),
- 7_i32.to_le_bytes()
- );
- assert_eq!(page_statistics.null_count_opt(), Some(0));
- assert!(page_statistics.distinct_count_opt().is_none());
+ assert_eq!(pages[2].page_type(), PageType::DATA_PAGE);
+ for (page, min, max) in [(&pages[1], 1_i32, 4_i32), (&pages[2], 5_i32,
7_i32)] {
+ let stats = page.statistics().unwrap();
+ assert_eq!(stats.min_bytes_opt().unwrap(), min.to_le_bytes());
+ assert_eq!(stats.max_bytes_opt().unwrap(), max.to_le_bytes());
+ assert_eq!(stats.null_count_opt(), Some(0));
+ assert!(stats.distinct_count_opt().is_none());
+ }
}
#[test]