From 1f5d629ff37a3eae3430cacbb080bd6e8aa2bcd5 Mon Sep 17 00:00:00 2001 From: seidl Date: Tue, 11 Aug 2026 09:51:55 -0700 Subject: [PATCH 1/3] convert page indexes to Vec>> --- parquet/src/arrow/arrow_reader/mod.rs | 7 +- parquet/src/arrow/arrow_reader/statistics.rs | 167 +++++++------- parquet/src/arrow/arrow_writer/mod.rs | 20 +- parquet/src/arrow/in_memory_row_group.rs | 18 +- .../arrow/push_decoder/reader_builder/data.rs | 2 +- .../arrow/push_decoder/reader_builder/mod.rs | 40 ++-- parquet/src/bin/parquet-concat.rs | 8 +- parquet/src/bin/parquet-index.rs | 34 +-- parquet/src/file/metadata/mod.rs | 22 +- parquet/src/file/metadata/parser.rs | 22 +- parquet/src/file/metadata/writer.rs | 73 +----- parquet/src/file/serialized_reader.rs | 218 +++++++++++++----- parquet/src/file/writer.rs | 39 +++- parquet/tests/arrow_reader/io/mod.rs | 5 +- .../tests/arrow_reader/row_filter/async.rs | 2 + parquet/tests/arrow_writer/layout.rs | 10 +- parquet/tests/encryption/encryption_util.rs | 6 +- 17 files changed, 401 insertions(+), 292 deletions(-) diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index 7ab219086fec..c14fe7c03d61 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -1318,7 +1318,12 @@ impl ReaderPageIterator { // To avoid `i[rg_idx][self.column_idx`] panic, we need to filter out empty `i[rg_idx]`. let page_locations = offset_index .filter(|i| !i[rg_idx].is_empty()) - .map(|i| i[rg_idx][self.column_idx].page_locations.clone()); + .map(|i| { + i[rg_idx][self.column_idx] + .as_ref() + .map(|o| o.page_locations.clone()) + }) + .unwrap_or(None); let total_rows = rg.num_rows() as usize; let reader = self.reader.clone(); diff --git a/parquet/src/arrow/arrow_reader/statistics.rs b/parquet/src/arrow/arrow_reader/statistics.rs index 7ff4a3a4e102..9b8a44d2d3c1 100644 --- a/parquet/src/arrow/arrow_reader/statistics.rs +++ b/parquet/src/arrow/arrow_reader/statistics.rs @@ -678,14 +678,14 @@ macro_rules! get_data_page_statistics { $values_iter: ident, $page_statistics: ident ) => {{ - let chunks: Vec<(usize, &ColumnIndexMetaData)> = $iterator.collect(); + let chunks: Vec<(usize, Option<&ColumnIndexMetaData>)> = $iterator.collect(); let capacity: usize = chunks.iter().map(|c| c.0).sum(); match $data_type { DataType::Boolean => { let mut b = BooleanBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::BOOLEAN(index) => { + Some(ColumnIndexMetaData::BOOLEAN(index)) => { for val in index.$values_iter() { b.append_option(val.copied()); } @@ -699,7 +699,7 @@ macro_rules! get_data_page_statistics { let mut b = UInt8Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.and_then(|&x| u8::try_from(x).ok())), @@ -714,7 +714,7 @@ macro_rules! get_data_page_statistics { let mut b = UInt16Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.and_then(|&x| u16::try_from(x).ok())), @@ -729,7 +729,7 @@ macro_rules! get_data_page_statistics { let mut b = UInt32Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|&x| x as u32)), @@ -744,7 +744,7 @@ macro_rules! get_data_page_statistics { let mut b = UInt64Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|&x| x as u64)), @@ -759,7 +759,7 @@ macro_rules! get_data_page_statistics { let mut b = Int8Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.and_then(|&x| i8::try_from(x).ok())), @@ -774,7 +774,7 @@ macro_rules! get_data_page_statistics { let mut b = Int16Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.and_then(|&x| i16::try_from(x).ok())), @@ -789,7 +789,7 @@ macro_rules! get_data_page_statistics { let mut b = Int32Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -804,7 +804,7 @@ macro_rules! get_data_page_statistics { let mut b = Int64Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -819,7 +819,7 @@ macro_rules! get_data_page_statistics { let mut b = Float16Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.and_then(|x| from_bytes_to_f16(x))), @@ -834,7 +834,7 @@ macro_rules! get_data_page_statistics { let mut b = Float32Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::FLOAT(index) => { + Some(ColumnIndexMetaData::FLOAT(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -849,7 +849,7 @@ macro_rules! get_data_page_statistics { let mut b = Float64Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::DOUBLE(index) => { + Some(ColumnIndexMetaData::DOUBLE(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -864,7 +864,7 @@ macro_rules! get_data_page_statistics { let mut b = BinaryBuilder::with_capacity(capacity, capacity * 10); for (len, index) in chunks { match index { - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { for val in index.$values_iter() { b.append_option(val.map(|x| x.as_ref())); } @@ -878,7 +878,7 @@ macro_rules! get_data_page_statistics { let mut b = LargeBinaryBuilder::with_capacity(capacity, capacity * 10); for (len, index) in chunks { match index { - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { for val in index.$values_iter() { b.append_option(val.map(|x| x.as_ref())); } @@ -892,7 +892,7 @@ macro_rules! get_data_page_statistics { let mut b = StringBuilder::with_capacity(capacity, capacity * 10); for (len, index) in chunks { match index { - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { for val in index.$values_iter() { match val { Some(x) => match std::str::from_utf8(x.as_ref()) { @@ -912,7 +912,7 @@ macro_rules! get_data_page_statistics { let mut b = LargeStringBuilder::with_capacity(capacity, capacity * 10); for (len, index) in chunks { match index { - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { for val in index.$values_iter() { match val { Some(x) => match std::str::from_utf8(x.as_ref()) { @@ -937,7 +937,7 @@ macro_rules! get_data_page_statistics { let mut b = TimestampSecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -952,7 +952,7 @@ macro_rules! get_data_page_statistics { let mut b = TimestampMillisecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -967,7 +967,7 @@ macro_rules! get_data_page_statistics { let mut b = TimestampMicrosecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -982,7 +982,7 @@ macro_rules! get_data_page_statistics { let mut b = TimestampNanosecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -999,7 +999,7 @@ macro_rules! get_data_page_statistics { let mut b = Date32Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -1014,7 +1014,7 @@ macro_rules! get_data_page_statistics { let mut b = Date64Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|&x| (x as i64) * 24 * 60 * 60 * 1000)), @@ -1029,7 +1029,7 @@ macro_rules! get_data_page_statistics { let mut b = Date64Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -1044,25 +1044,25 @@ macro_rules! get_data_page_statistics { let mut b = Decimal32Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), ); } - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.and_then(|&x| i32::try_from(x).ok())), ); } - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i32(x.as_ref()))), ); } - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i32(x.as_ref()))), @@ -1077,25 +1077,25 @@ macro_rules! get_data_page_statistics { let mut b = Decimal64Builder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| *x as i64)), ); } - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), ); } - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i64(x.as_ref()))), ); } - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i64(x.as_ref()))), @@ -1110,25 +1110,25 @@ macro_rules! get_data_page_statistics { let mut b = Decimal128Array::builder(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| *x as i128)), ); } - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| *x as i128)), ); } - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i128(x.as_ref()))), ); } - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i128(x.as_ref()))), @@ -1143,25 +1143,25 @@ macro_rules! get_data_page_statistics { let mut b = Decimal256Array::builder(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| i256::from_i128(*x as i128))), ); } - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| i256::from_i128(*x as i128))), ); } - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i256(x.as_ref()))), ); } - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.map(|x| from_bytes_to_i256(x.as_ref()))), @@ -1178,7 +1178,7 @@ macro_rules! get_data_page_statistics { let mut b = Time32SecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -1193,7 +1193,7 @@ macro_rules! get_data_page_statistics { let mut b = Time32MillisecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT32(index) => { + Some(ColumnIndexMetaData::INT32(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -1215,7 +1215,7 @@ macro_rules! get_data_page_statistics { let mut b = Time64MicrosecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -1230,7 +1230,7 @@ macro_rules! get_data_page_statistics { let mut b = Time64NanosecondBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::INT64(index) => { + Some(ColumnIndexMetaData::INT64(index)) => { b.extend_from_iter_option( index.$values_iter() .map(|val| val.copied()), @@ -1250,7 +1250,7 @@ macro_rules! get_data_page_statistics { let mut b = FixedSizeBinaryBuilder::with_capacity(capacity, *size); for (len, index) in chunks { match index { - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index)) => { for val in index.$values_iter() { match val { Some(v) => { @@ -1273,7 +1273,7 @@ macro_rules! get_data_page_statistics { let mut b = StringViewBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { for val in index.$values_iter() { match val { Some(x) => match std::str::from_utf8(x.as_ref()) { @@ -1295,7 +1295,7 @@ macro_rules! get_data_page_statistics { let mut b = BinaryViewBuilder::with_capacity(capacity); for (len, index) in chunks { match index { - ColumnIndexMetaData::BYTE_ARRAY(index) => { + Some(ColumnIndexMetaData::BYTE_ARRAY(index)) => { for val in index.$values_iter() { match val { Some(v) => b.append_value(v.as_ref()), @@ -1361,7 +1361,7 @@ pub(crate) fn min_page_statistics<'a, I>( physical_type: Option, ) -> Result where - I: Iterator, + I: Iterator)>, { get_data_page_statistics!(Min, data_type, iterator, physical_type) } @@ -1374,7 +1374,7 @@ pub(crate) fn max_page_statistics<'a, I>( physical_type: Option, ) -> Result where - I: Iterator, + I: Iterator)>, { get_data_page_statistics!(Max, data_type, iterator, physical_type) } @@ -1385,19 +1385,19 @@ where /// The returned Array is an [`UInt64Array`] pub(crate) fn null_counts_page_statistics<'a, I>(iterator: I) -> Result where - I: Iterator, + I: Iterator)>, { let chunks: Vec<_> = iterator.collect(); let total_capacity: usize = chunks.iter().map(|(len, _)| *len).sum(); let mut values = Vec::with_capacity(total_capacity); let mut nulls = NullBufferBuilder::new(total_capacity); for (len, index) in chunks { - match index.null_counts() { - Some(counts) => { - values.extend(counts.iter().map(|&x| x as u64)); + match index { + Some(index) if index.null_counts().is_some() => { + values.extend(index.null_counts().unwrap().iter().map(|&x| x as u64)); nulls.append_n_non_nulls(len); } - None => { + _ => { values.resize(values.len() + len, 0); nulls.append_n_nulls(len); } @@ -1414,19 +1414,19 @@ where /// The returned Array is an [`UInt64Array`] pub(crate) fn nan_counts_page_statistics<'a, I>(iterator: I) -> Result where - I: Iterator, + I: Iterator)>, { let chunks: Vec<_> = iterator.collect(); let total_capacity: usize = chunks.iter().map(|(len, _)| *len).sum(); let mut values = Vec::with_capacity(total_capacity); let mut nulls = NullBufferBuilder::new(total_capacity); for (len, index) in chunks { - match index.nan_counts() { - Some(counts) => { - values.extend(counts.iter().map(|&x| x as u64)); + match index { + Some(index) if index.nan_counts().is_some() => { + values.extend(index.nan_counts().unwrap().iter().map(|&x| x as u64)); nulls.append_n_non_nulls(len); } - None => { + _ => { values.resize(values.len() + len, 0); nulls.append_n_nulls(len); } @@ -1889,12 +1889,13 @@ impl<'a> StatisticsConverter<'a> { let iter = row_group_indices.into_iter().map(|rg_index| { let column_page_index_per_row_group_per_column = - &column_page_index[*rg_index][parquet_index]; - let num_data_pages = &column_offset_index[*rg_index][parquet_index] - .page_locations() - .len(); + column_page_index[*rg_index][parquet_index].as_ref(); + let num_data_pages = column_offset_index[*rg_index][parquet_index] + .as_ref() + .map(|offset| offset.page_locations().len()) + .unwrap_or(0); - (*num_data_pages, column_page_index_per_row_group_per_column) + (num_data_pages, column_page_index_per_row_group_per_column) }); min_page_statistics(data_type, iter, self.physical_type) @@ -1920,12 +1921,13 @@ impl<'a> StatisticsConverter<'a> { let iter = row_group_indices.into_iter().map(|rg_index| { let column_page_index_per_row_group_per_column = - &column_page_index[*rg_index][parquet_index]; - let num_data_pages = &column_offset_index[*rg_index][parquet_index] - .page_locations() - .len(); + column_page_index[*rg_index][parquet_index].as_ref(); + let num_data_pages = column_offset_index[*rg_index][parquet_index] + .as_ref() + .map(|offset| offset.page_locations().len()) + .unwrap_or(0); - (*num_data_pages, column_page_index_per_row_group_per_column) + (num_data_pages, column_page_index_per_row_group_per_column) }); max_page_statistics(data_type, iter, self.physical_type) @@ -1950,12 +1952,13 @@ impl<'a> StatisticsConverter<'a> { let iter = row_group_indices.into_iter().map(|rg_index| { let column_page_index_per_row_group_per_column = - &column_page_index[*rg_index][parquet_index]; - let num_data_pages = &column_offset_index[*rg_index][parquet_index] - .page_locations() - .len(); + column_page_index[*rg_index][parquet_index].as_ref(); + let num_data_pages = column_offset_index[*rg_index][parquet_index] + .as_ref() + .map(|offset| offset.page_locations().len()) + .unwrap_or(0); - (*num_data_pages, column_page_index_per_row_group_per_column) + (num_data_pages, column_page_index_per_row_group_per_column) }); null_counts_page_statistics(iter) } @@ -1979,12 +1982,13 @@ impl<'a> StatisticsConverter<'a> { let iter = row_group_indices.into_iter().map(|rg_index| { let column_page_index_per_row_group_per_column = - &column_page_index[*rg_index][parquet_index]; - let num_data_pages = &column_offset_index[*rg_index][parquet_index] - .page_locations() - .len(); + column_page_index[*rg_index][parquet_index].as_ref(); + let num_data_pages = column_offset_index[*rg_index][parquet_index] + .as_ref() + .map(|offset| offset.page_locations().len()) + .unwrap_or(0); - (*num_data_pages, column_page_index_per_row_group_per_column) + (num_data_pages, column_page_index_per_row_group_per_column) }); nan_counts_page_statistics(iter) } @@ -2025,7 +2029,10 @@ impl<'a> StatisticsConverter<'a> { let mut row_counts = Vec::new(); let mut nulls = NullBufferBuilder::new(0); for rg_idx in row_group_indices { - let page_locations = &column_offset_index[*rg_idx][parquet_index].page_locations(); + let Some(offset_index) = &column_offset_index[*rg_idx][parquet_index] else { + continue; + }; + let page_locations = offset_index.page_locations(); let row_count_per_page = page_locations .windows(2) diff --git a/parquet/src/arrow/arrow_writer/mod.rs b/parquet/src/arrow/arrow_writer/mod.rs index 4ffdb17c964d..5d054139c2f8 100644 --- a/parquet/src/arrow/arrow_writer/mod.rs +++ b/parquet/src/arrow/arrow_writer/mod.rs @@ -2964,7 +2964,7 @@ mod tests { assert!(reader.metadata().offset_index().is_some()); let offset_indexes = &reader.metadata().offset_index().unwrap()[0]; - let page_locations = offset_indexes[0].page_locations.clone(); + let page_locations = offset_indexes[0].as_ref().unwrap().page_locations.clone(); // We should fallback to PLAIN encoding after the first row and our max page size is 1 bytes // so we expect one dictionary encoded page and then a page per row thereafter. @@ -3372,6 +3372,8 @@ mod tests { if let Some(col_indexes) = file_meta_data.column_index() { for rg_idx in col_indexes { for idx in rg_idx { + assert!(idx.is_some()); + let idx = idx.as_ref().unwrap(); assert!(idx.nan_counts().is_some()); let float_idx = match idx { ColumnIndexMetaData::DOUBLE(idx) => idx, @@ -3449,11 +3451,11 @@ mod tests { assert!(file_meta_data.column_index().is_some()); let col_idx = &file_meta_data.column_index().as_ref().unwrap()[0][0]; - assert_eq!(col_idx.num_pages(), 4); + assert_eq!(col_idx.as_ref().unwrap().num_pages(), 4); // test each page let float_idx = match col_idx { - ColumnIndexMetaData::DOUBLE(idx) => idx, + Some(ColumnIndexMetaData::DOUBLE(idx)) => idx, _ => panic!("expected double statistics"), }; @@ -4947,8 +4949,8 @@ mod tests { assert_eq!(index.len(), 1); assert_eq!(index[0].len(), 2); // 2 columns - assert_eq!(index[0][0].page_locations().len(), 1); // 1 page - assert_eq!(index[0][1].page_locations().len(), 1); // 1 page + assert_eq!(index[0][0].as_ref().unwrap().page_locations().len(), 1); // 1 page + assert_eq!(index[0][1].as_ref().unwrap().page_locations().len(), 1); // 1 page } #[test] @@ -5019,11 +5021,11 @@ mod tests { let a_idx = &column_index[0][0]; assert!( - matches!(a_idx, ColumnIndexMetaData::BYTE_ARRAY(_)), + matches!(a_idx, Some(ColumnIndexMetaData::BYTE_ARRAY(_))), "{a_idx:?}" ); let b_idx = &column_index[0][1]; - assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}"); + assert!(b_idx.is_none(), "{b_idx:?}"); } #[test] @@ -5089,9 +5091,9 @@ mod tests { assert_eq!(column_index[0].len(), 2); // 2 columns let a_idx = &column_index[0][0]; - assert!(matches!(a_idx, ColumnIndexMetaData::NONE), "{a_idx:?}"); + assert!(a_idx.is_none(), "{a_idx:?}"); let b_idx = &column_index[0][1]; - assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}"); + assert!(b_idx.is_none(), "{b_idx:?}"); } #[test] diff --git a/parquet/src/arrow/in_memory_row_group.rs b/parquet/src/arrow/in_memory_row_group.rs index 6c5f013159d5..db523e5be1cc 100644 --- a/parquet/src/arrow/in_memory_row_group.rs +++ b/parquet/src/arrow/in_memory_row_group.rs @@ -30,7 +30,7 @@ use std::sync::Arc; /// An in-memory collection of column chunks #[derive(Debug)] pub(crate) struct InMemoryRowGroup<'a> { - pub(crate) offset_index: Option<&'a [OffsetIndexMetaData]>, + pub(crate) offset_index: Option<&'a [Option]>, /// Column chunks for this row group pub(crate) column_chunks: Vec>>, pub(crate) row_count: usize, @@ -85,7 +85,13 @@ impl InMemoryRowGroup<'_> { // then we need to also fetch a dictionary page. let mut ranges: Vec> = vec![]; let (start, _len) = chunk_meta.byte_range(); - match offset_index[idx].page_locations.first() { + let Some(offset_idx) = offset_index[idx].as_ref() else { + // No offset index for this column, fetch the entire column + ranges.push(start..start + _len); + return ranges; + }; + + match offset_idx.page_locations.first() { Some(first) if first.offset as u64 != start => { ranges.push(start..first.offset as u64); } @@ -96,11 +102,9 @@ impl InMemoryRowGroup<'_> { // (see doc comment for this function for details on `cache_mask`) let use_expanded = cache_mask.map(|m| m.leaf_included(idx)).unwrap_or(false); if use_expanded { - ranges.extend( - expanded_selection.scan_ranges(&offset_index[idx].page_locations), - ); + ranges.extend(expanded_selection.scan_ranges(&offset_idx.page_locations)); } else { - ranges.extend(selection.scan_ranges(&offset_index[idx].page_locations)); + ranges.extend(selection.scan_ranges(&offset_idx.page_locations)); } page_start_offsets.push(ranges.iter().map(|range| range.start).collect()); @@ -203,7 +207,7 @@ impl RowGroups for InMemoryRowGroup<'_> { .offset_index // filter out empty offset indexes (old versions specified Some(vec![]) when no present) .filter(|index| !index.is_empty()) - .map(|index| index[i].page_locations.clone()); + .and_then(|index| index[i].as_ref().map(|idx| idx.page_locations.clone())); let column_chunk_metadata = self.metadata.row_group(self.row_group_idx).column(i); let page_reader = SerializedPageReader::new( data.clone(), diff --git a/parquet/src/arrow/push_decoder/reader_builder/data.rs b/parquet/src/arrow/push_decoder/reader_builder/data.rs index 6fbc2090b06e..4bcbf7e8dc50 100644 --- a/parquet/src/arrow/push_decoder/reader_builder/data.rs +++ b/parquet/src/arrow/push_decoder/reader_builder/data.rs @@ -224,7 +224,7 @@ impl<'a> DataRequestBuilder<'a> { fn get_offset_index( parquet_metadata: &ParquetMetaData, row_group_idx: usize, -) -> Option<&[OffsetIndexMetaData]> { +) -> Option<&[Option]> { parquet_metadata .offset_index() // filter out empty offset indexes (old versions specified Some(vec![]) when no present) diff --git a/parquet/src/arrow/push_decoder/reader_builder/mod.rs b/parquet/src/arrow/push_decoder/reader_builder/mod.rs index ffc10382443f..1a809e518741 100644 --- a/parquet/src/arrow/push_decoder/reader_builder/mod.rs +++ b/parquet/src/arrow/push_decoder/reader_builder/mod.rs @@ -849,7 +849,10 @@ impl RowGroupReaderBuilder { } /// Get the offset index for the specified row group, if any - fn row_group_offset_index(&self, row_group_idx: usize) -> Option<&[OffsetIndexMetaData]> { + fn row_group_offset_index( + &self, + row_group_idx: usize, + ) -> Option<&[Option]> { self.metadata .offset_index() .filter(|index| !index.is_empty()) @@ -881,7 +884,7 @@ impl RowGroupReaderBuilder { fn prepare_selection_for_page_skipping( plan_builder: ReadPlanBuilder, projection_mask: &ProjectionMask, - offset_index: Option<&[OffsetIndexMetaData]>, + offset_index: Option<&[Option]>, total_rows: usize, ) -> ReadPlanBuilder { match plan_builder.resolve_selection_strategy() { @@ -906,7 +909,7 @@ fn prepare_selection_for_page_skipping( fn loaded_row_ranges_for_projection( selection: Option<&RowSelection>, projection_mask: &ProjectionMask, - offset_index: Option<&[OffsetIndexMetaData]>, + offset_index: Option<&[Option]>, total_rows: usize, ) -> Option { let selection = selection?; @@ -916,7 +919,8 @@ fn loaded_row_ranges_for_projection( .iter() .enumerate() .filter_map(|(leaf_idx, column)| { - let pages = column.page_locations(); + let column_metadata = column.as_ref()?; + let pages = column_metadata.page_locations(); (projection_mask.leaf_included(leaf_idx) && !pages.is_empty()).then(|| { RowSelection::from_consecutive_ranges( selection @@ -945,17 +949,19 @@ mod tests { #[test] fn test_loaded_row_ranges_intersect_column_page_boundaries() { - let column = |first_rows: &[i64]| OffsetIndexMetaData { - page_locations: first_rows - .iter() - .enumerate() - .map(|(idx, first_row_index)| PageLocation { - offset: (idx * 10) as i64, - compressed_page_size: 10, - first_row_index: *first_row_index, - }) - .collect(), - unencoded_byte_array_data_bytes: None, + let column = |first_rows: &[i64]| { + Some(OffsetIndexMetaData { + page_locations: first_rows + .iter() + .enumerate() + .map(|(idx, first_row_index)| PageLocation { + offset: (idx * 10) as i64, + compressed_page_size: 10, + first_row_index: *first_row_index, + }) + .collect(), + unencoded_byte_array_data_bytes: None, + }) }; let columns = vec![column(&[0, 4, 8]), column(&[0, 6, 10])]; let selection = RowSelection::from(vec![ @@ -978,7 +984,7 @@ mod tests { #[test] fn test_auto_keeps_mask_when_page_pruning_skips_pages() { - let columns = vec![OffsetIndexMetaData { + let columns = vec![Some(OffsetIndexMetaData { page_locations: [0, 2, 4, 6, 8, 10] .into_iter() .enumerate() @@ -989,7 +995,7 @@ mod tests { }) .collect(), unencoded_byte_array_data_bytes: None, - }]; + })]; let selection = RowSelection::from(vec![ RowSelector::select(1), RowSelector::skip(10), diff --git a/parquet/src/bin/parquet-concat.rs b/parquet/src/bin/parquet-concat.rs index a6f1aef78110..2ef41cba8b30 100644 --- a/parquet/src/bin/parquet-concat.rs +++ b/parquet/src/bin/parquet-concat.rs @@ -109,9 +109,13 @@ impl Args { let mut rg_out = writer.next_row_group()?; for (col_idx, column) in rg.columns().iter().enumerate() { let bloom_filter = read_bloom_filter(column, &input); - let column_index = rg_column_indexes.and_then(|row| row.get(col_idx)).cloned(); + let column_index = rg_column_indexes + .and_then(|row| row.get(col_idx)) + .and_then(|opt| opt.clone()); - let offset_index = rg_offset_indexes.and_then(|row| row.get(col_idx)).cloned(); + let offset_index = rg_offset_indexes + .and_then(|row| row.get(col_idx)) + .and_then(|opt| opt.clone()); let result = ColumnCloseResult { bytes_written: column.compressed_size() as _, diff --git a/parquet/src/bin/parquet-index.rs b/parquet/src/bin/parquet-index.rs index 241fc20533d3..ae99d9841c2b 100644 --- a/parquet/src/bin/parquet-index.rs +++ b/parquet/src/bin/parquet-index.rs @@ -99,24 +99,28 @@ impl Args { ParquetError::General(format!( "No offset index for row group {row_group_idx} column chunk {column_idx}" )) + })?.as_ref().ok_or_else(|| { + ParquetError::General(format!( + "Offset index is None for row group {row_group_idx} column chunk {column_idx}" + )) })?; let row_counts = - compute_row_counts(offset_index.page_locations.as_slice(), row_group.num_rows()); - match &column_indices[column_idx] { - ColumnIndexMetaData::NONE => println!("NO INDEX"), - ColumnIndexMetaData::BOOLEAN(v) => { + compute_row_counts(offset_index.page_locations(), row_group.num_rows()); + match column_indices[column_idx].as_ref() { + None | Some(ColumnIndexMetaData::NONE) => println!("NO INDEX"), + Some(ColumnIndexMetaData::BOOLEAN(v)) => { print_index::(v, offset_index, &row_counts)? } - ColumnIndexMetaData::INT32(v) => print_index(v, offset_index, &row_counts)?, - ColumnIndexMetaData::INT64(v) => print_index(v, offset_index, &row_counts)?, - ColumnIndexMetaData::INT96(v) => print_index(v, offset_index, &row_counts)?, - ColumnIndexMetaData::FLOAT(v) => print_index(v, offset_index, &row_counts)?, - ColumnIndexMetaData::DOUBLE(v) => print_index(v, offset_index, &row_counts)?, - ColumnIndexMetaData::BYTE_ARRAY(v) => { + Some(ColumnIndexMetaData::INT32(v)) => print_index(v, offset_index, &row_counts)?, + Some(ColumnIndexMetaData::INT64(v)) => print_index(v, offset_index, &row_counts)?, + Some(ColumnIndexMetaData::INT96(v)) => print_index(v, offset_index, &row_counts)?, + Some(ColumnIndexMetaData::FLOAT(v)) => print_index(v, offset_index, &row_counts)?, + Some(ColumnIndexMetaData::DOUBLE(v)) => print_index(v, offset_index, &row_counts)?, + Some(ColumnIndexMetaData::BYTE_ARRAY(v)) => { print_bytes_index(v, offset_index, &row_counts)? } - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v)) => { print_bytes_index(v, offset_index, &row_counts)? } } @@ -147,11 +151,11 @@ fn print_index( offset_index: &OffsetIndexMetaData, row_counts: &[i64], ) -> Result<()> { - if column_index.num_pages() as usize != offset_index.page_locations.len() { + if column_index.num_pages() as usize != offset_index.page_locations().len() { return Err(ParquetError::General(format!( "Index length mismatch, got {} and {}", column_index.num_pages(), - offset_index.page_locations.len() + offset_index.page_locations().len() ))); } @@ -186,11 +190,11 @@ fn print_bytes_index( offset_index: &OffsetIndexMetaData, row_counts: &[i64], ) -> Result<()> { - if column_index.num_pages() as usize != offset_index.page_locations.len() { + if column_index.num_pages() as usize != offset_index.page_locations().len() { return Err(ParquetError::General(format!( "Index length mismatch, got {} and {}", column_index.num_pages(), - offset_index.page_locations.len() + offset_index.page_locations().len() ))); } diff --git a/parquet/src/file/metadata/mod.rs b/parquet/src/file/metadata/mod.rs index 874449d992c0..e8461a46fa67 100644 --- a/parquet/src/file/metadata/mod.rs +++ b/parquet/src/file/metadata/mod.rs @@ -141,29 +141,31 @@ pub(crate) use writer::ThriftMetadataWriter; /// documentation]. Each [`ColumnIndex`] holds statistics about all the pages in a /// particular column chunk. /// -/// `column_index[row_group_number][column_number]` holds the +/// `column_index[row_group_number][column_number]` holds the optional /// [`ColumnIndex`] corresponding to column `column_number` of row group -/// `row_group_number`. +/// `row_group_number`. This will be `None` if no index is present for the given +/// column chunk. /// /// For example `column_index[2][3]` holds the [`ColumnIndex`] for the fourth /// column in the third row group of the parquet file. /// /// [PageIndex documentation]: https://github.com/apache/parquet-format/blob/master/PageIndex.md /// [`ColumnIndex`]: crate::file::page_index::column_index::ColumnIndexMetaData -pub type ParquetColumnIndex = Vec>; +pub type ParquetColumnIndex = Vec>>; -/// [`OffsetIndexMetaData`] for each data page of each row group of each column +/// [`OffsetIndexMetaData`] for each column chunk of each row group /// /// This structure is the parsed representation of the [`OffsetIndex`] from the /// Parquet file footer, as described in the Parquet [PageIndex documentation]. /// /// `offset_index[row_group_number][column_number]` holds -/// the [`OffsetIndexMetaData`] corresponding to column -/// `column_number`of row group `row_group_number`. +/// the optional [`OffsetIndexMetaData`] corresponding to column +/// `column_number`of row group `row_group_number`. This will be `None` if no index +/// is present for the given column chunk. /// /// [PageIndex documentation]: https://github.com/apache/parquet-format/blob/master/PageIndex.md /// [`OffsetIndex`]: https://github.com/apache/parquet-format/blob/master/PageIndex.md -pub type ParquetOffsetIndex = Vec>; +pub type ParquetOffsetIndex = Vec>>; /// Parsed metadata for a single Parquet file /// @@ -2117,11 +2119,13 @@ mod tests { offset_index.append_row_count(1); offset_index.append_offset_and_size(2, 3); offset_index.append_unencoded_byte_array_data_bytes(Some(10)); - let offset_index = offset_index.build(); + let offset_index = Some(offset_index.build()); let parquet_meta = ParquetMetaDataBuilder::new(file_metadata) .set_row_groups(row_group_meta) - .set_column_index(Some(vec![vec![ColumnIndexMetaData::BOOLEAN(native_index)]])) + .set_column_index(Some(vec![vec![Some(ColumnIndexMetaData::BOOLEAN( + native_index, + ))]])) .set_offset_index(Some(vec![vec![offset_index]])) .build(); diff --git a/parquet/src/file/metadata/parser.rs b/parquet/src/file/metadata/parser.rs index 9df6bcdd7185..61296527b0dc 100644 --- a/parquet/src/file/metadata/parser.rs +++ b/parquet/src/file/metadata/parser.rs @@ -269,8 +269,9 @@ pub(crate) fn parse_column_index( rg_idx, col_idx, ) + .map(Some) } - None => Ok(ColumnIndexMetaData::NONE), + None => Ok(None), }) .collect::>>() }) @@ -305,23 +306,18 @@ pub(crate) fn parse_offset_index( rg_idx, col_idx, ) + .map(Some) } - None => Err(general_err!("missing offset index")), - }; - - match result { - Ok(index) => row_group_indexes.push(index), - Err(e) => { + None => { if offset_index_policy == PageIndexPolicy::Required { - return Err(e); + Err(general_err!("missing offset index")) } else { - // Invalidate and return - metadata.set_column_index(None); - metadata.set_offset_index(None); - return Ok(()); + Ok(None) } } - } + }; + + row_group_indexes.push(result?); } all_indexes.push(row_group_indexes); } diff --git a/parquet/src/file/metadata/writer.rs b/parquet/src/file/metadata/writer.rs index 4b88077d5857..574522c1fdb1 100644 --- a/parquet/src/file/metadata/writer.rs +++ b/parquet/src/file/metadata/writer.rs @@ -149,23 +149,11 @@ impl<'a, W: Write> ThriftMetadataWriter<'a, W> { .as_ref() .is_some_and(|ci| ci.iter().all(|cii| cii.iter().all(|idx| idx.is_none()))); - // transform from Option>>> to - // Option>> - let column_indexes: Option = if all_none { - None + if all_none { + Ok(None) } else { - column_indexes.map(|ovvi| { - ovvi.into_iter() - .map(|vi| { - vi.into_iter() - .map(|ci| ci.unwrap_or(ColumnIndexMetaData::NONE)) - .collect() - }) - .collect() - }) - }; - - Ok(column_indexes) + Ok(column_indexes) + } } /// Serialize the offset indexes and transform to `Option` @@ -182,18 +170,11 @@ impl<'a, W: Write> ThriftMetadataWriter<'a, W> { .as_ref() .is_some_and(|oi| oi.iter().all(|oii| oii.iter().all(|idx| idx.is_none()))); - let offset_indexes: Option = if all_none { - None + if all_none { + Ok(None) } else { - // FIXME(ets): this will panic if there's a missing index. - offset_indexes.map(|ovvi| { - ovvi.into_iter() - .map(|vi| vi.into_iter().map(|oi| oi.unwrap()).collect()) - .collect() - }) - }; - - Ok(offset_indexes) + Ok(offset_indexes) + } } /// Assembles and writes the final metadata to self.buf @@ -464,8 +445,8 @@ impl<'a, W: Write> ParquetMetaDataWriter<'a, W> { let key_value_metadata = file_metadata.key_value_metadata().cloned(); - let column_indexes = self.convert_column_indexes(); - let offset_indexes = self.convert_offset_index(); + let column_indexes = self.metadata.column_index().cloned(); + let offset_indexes = self.metadata.offset_index().cloned(); let mut encoder = ThriftMetadataWriter::new( &mut self.buf, @@ -491,40 +472,6 @@ impl<'a, W: Write> ParquetMetaDataWriter<'a, W> { Ok(()) } - - fn convert_column_indexes(&self) -> Option>>> { - // TODO(ets): we're converting from ParquetColumnIndex to vec>, - // but then converting back to ParquetColumnIndex in the end. need to unify this. - self.metadata - .column_index() - .map(|row_group_column_indexes| { - (0..self.metadata.row_groups().len()) - .map(|rg_idx| { - let column_indexes = &row_group_column_indexes[rg_idx]; - column_indexes - .iter() - .map(|column_index| Some(column_index.clone())) - .collect() - }) - .collect() - }) - } - - fn convert_offset_index(&self) -> Option>>> { - self.metadata - .offset_index() - .map(|row_group_offset_indexes| { - (0..self.metadata.row_groups().len()) - .map(|rg_idx| { - let offset_indexes = &row_group_offset_indexes[rg_idx]; - offset_indexes - .iter() - .map(|offset_index| Some(offset_index.clone())) - .collect() - }) - .collect() - }) - } } #[derive(Debug, Default)] diff --git a/parquet/src/file/serialized_reader.rs b/parquet/src/file/serialized_reader.rs index fe9a5c04863c..53843fb9a6d5 100644 --- a/parquet/src/file/serialized_reader.rs +++ b/parquet/src/file/serialized_reader.rs @@ -322,7 +322,7 @@ impl FileReader for SerializedFileReader { pub struct SerializedRowGroupReader<'a, R: ChunkReader> { chunk_reader: Arc, metadata: &'a RowGroupMetaData, - offset_index: Option<&'a [OffsetIndexMetaData]>, + offset_index: Option<&'a [Option]>, props: ReaderPropertiesPtr, bloom_filters: Vec>, } @@ -332,7 +332,7 @@ impl<'a, R: ChunkReader> SerializedRowGroupReader<'a, R> { pub fn new( chunk_reader: Arc, metadata: &'a RowGroupMetaData, - offset_index: Option<&'a [OffsetIndexMetaData]>, + offset_index: Option<&'a [Option]>, props: ReaderPropertiesPtr, ) -> Result { let bloom_filters = if props.read_bloom_filter() { @@ -367,7 +367,11 @@ impl RowGroupReader for SerializedRowGroupReader<'_, R fn get_column_page_reader(&self, i: usize) -> Result> { let col = self.metadata.column(i); - let page_locations = self.offset_index.map(|x| x[i].page_locations.clone()); + let page_locations = if let Some(offset_index) = self.offset_index { + offset_index[i].as_ref().map(|oi| oi.page_locations.clone()) + } else { + None + }; let props = Arc::clone(&self.props); Ok(Box::new(SerializedPageReader::new_with_properties( @@ -1761,9 +1765,13 @@ mod tests { let col = row_group.metadata.column(column); - let page_locations = row_group - .offset_index - .map(|x| x[column].page_locations.clone()); + let page_locations = if let Some(offset_index) = row_group.offset_index { + offset_index[column] + .as_ref() + .map(|oi| oi.page_locations.clone()) + } else { + None + }; let props = Arc::clone(&row_group.props); SerializedPageReader::new_with_properties( @@ -2154,7 +2162,7 @@ mod tests { // only one row group assert_eq!(column_index.len(), 1); - let index = if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][0] { + let index = if let Some(ColumnIndexMetaData::BYTE_ARRAY(index)) = &column_index[0][0] { index } else { unreachable!() @@ -2174,7 +2182,7 @@ mod tests { // only one row group assert_eq!(offset_indexes.len(), 1); let offset_index = &offset_indexes[0]; - let page_offset = &offset_index[0].page_locations()[0]; + let page_offset = &offset_index[0].as_ref().unwrap().page_locations()[0]; assert_eq!(4, page_offset.offset); assert_eq!(152, page_offset.compressed_page_size); @@ -2202,164 +2210,256 @@ mod tests { let row_group_metadata = metadata.row_group(0); //col0->id: INT32 UNCOMPRESSED DO:0 FPO:4 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 7299, num_nulls: 0] - assert!(!&column_index[0][0].is_sorted()); - let boundary_order = &column_index[0][0].get_boundary_order(); - assert!(boundary_order.is_some()); - matches!(boundary_order.unwrap(), BoundaryOrder::UNORDERED); - if let ColumnIndexMetaData::INT32(index) = &column_index[0][0] { + let ci = column_index[0][0].as_ref().unwrap(); + assert!(!ci.is_sorted()); + assert!(matches!( + ci.get_boundary_order(), + Some(BoundaryOrder::UNORDERED) + )); + if let ColumnIndexMetaData::INT32(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 0), BoundaryOrder::UNORDERED, ); - assert_eq!(row_group_offset_indexes[0].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[0] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; //col1->bool_col:BOOLEAN UNCOMPRESSED DO:0 FPO:37329 SZ:3022/3022/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: false, max: true, num_nulls: 0] - assert!(&column_index[0][1].is_sorted()); - if let ColumnIndexMetaData::BOOLEAN(index) = &column_index[0][1] { + let ci = column_index[0][1].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::BOOLEAN(index) = ci { assert_eq!(index.num_pages(), 82); - assert_eq!(row_group_offset_indexes[1].page_locations.len(), 82); + assert_eq!( + row_group_offset_indexes[1] + .as_ref() + .unwrap() + .page_locations + .len(), + 82 + ); } else { unreachable!() }; //col2->tinyint_col: INT32 UNCOMPRESSED DO:0 FPO:40351 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0] - assert!(&column_index[0][2].is_sorted()); - if let ColumnIndexMetaData::INT32(index) = &column_index[0][2] { + let ci = column_index[0][2].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::INT32(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 2), BoundaryOrder::ASCENDING, ); - assert_eq!(row_group_offset_indexes[2].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[2] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; //col4->smallint_col: INT32 UNCOMPRESSED DO:0 FPO:77676 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0] - assert!(&column_index[0][3].is_sorted()); - if let ColumnIndexMetaData::INT32(index) = &column_index[0][3] { + let ci = column_index[0][3].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::INT32(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 3), BoundaryOrder::ASCENDING, ); - assert_eq!(row_group_offset_indexes[3].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[3] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; //col5->smallint_col: INT32 UNCOMPRESSED DO:0 FPO:77676 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0] - assert!(&column_index[0][4].is_sorted()); - if let ColumnIndexMetaData::INT32(index) = &column_index[0][4] { + let ci = column_index[0][4].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::INT32(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 4), BoundaryOrder::ASCENDING, ); - assert_eq!(row_group_offset_indexes[4].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[4] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; //col6->bigint_col: INT64 UNCOMPRESSED DO:0 FPO:152326 SZ:71598/71598/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 90, num_nulls: 0] - assert!(!&column_index[0][5].is_sorted()); - if let ColumnIndexMetaData::INT64(index) = &column_index[0][5] { + let ci = column_index[0][5].as_ref().unwrap(); + assert!(!ci.is_sorted()); + if let ColumnIndexMetaData::INT64(index) = ci { check_native_page_index( index, 528, get_row_group_min_max_bytes(row_group_metadata, 5), BoundaryOrder::UNORDERED, ); - assert_eq!(row_group_offset_indexes[5].page_locations.len(), 528); + assert_eq!( + row_group_offset_indexes[5] + .as_ref() + .unwrap() + .page_locations + .len(), + 528 + ); } else { unreachable!() }; //col7->float_col: FLOAT UNCOMPRESSED DO:0 FPO:223924 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: -0.0, max: 9.9, num_nulls: 0] - assert!(&column_index[0][6].is_sorted()); - if let ColumnIndexMetaData::FLOAT(index) = &column_index[0][6] { + let ci = column_index[0][6].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::FLOAT(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 6), BoundaryOrder::ASCENDING, ); - assert_eq!(row_group_offset_indexes[6].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[6] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; //col8->double_col: DOUBLE UNCOMPRESSED DO:0 FPO:261249 SZ:71598/71598/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: -0.0, max: 90.89999999999999, num_nulls: 0] - assert!(!&column_index[0][7].is_sorted()); - if let ColumnIndexMetaData::DOUBLE(index) = &column_index[0][7] { + let ci = column_index[0][7].as_ref().unwrap(); + assert!(!ci.is_sorted()); + if let ColumnIndexMetaData::DOUBLE(index) = ci { check_native_page_index( index, 528, get_row_group_min_max_bytes(row_group_metadata, 7), BoundaryOrder::UNORDERED, ); - assert_eq!(row_group_offset_indexes[7].page_locations.len(), 528); + assert_eq!( + row_group_offset_indexes[7] + .as_ref() + .unwrap() + .page_locations + .len(), + 528 + ); } else { unreachable!() }; //col9->date_string_col: BINARY UNCOMPRESSED DO:0 FPO:332847 SZ:111948/111948/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 01/01/09, max: 12/31/10, num_nulls: 0] - assert!(!&column_index[0][8].is_sorted()); - if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][8] { + let ci = column_index[0][8].as_ref().unwrap(); + assert!(!ci.is_sorted()); + if let ColumnIndexMetaData::BYTE_ARRAY(index) = ci { check_byte_array_page_index( index, 974, get_row_group_min_max_bytes(row_group_metadata, 8), BoundaryOrder::UNORDERED, ); - assert_eq!(row_group_offset_indexes[8].page_locations.len(), 974); + assert_eq!( + row_group_offset_indexes[8] + .as_ref() + .unwrap() + .page_locations + .len(), + 974 + ); } else { unreachable!() }; //col10->string_col: BINARY UNCOMPRESSED DO:0 FPO:444795 SZ:45298/45298/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0] - assert!(&column_index[0][9].is_sorted()); - if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][9] { + let ci = column_index[0][9].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::BYTE_ARRAY(index) = ci { check_byte_array_page_index( index, 352, get_row_group_min_max_bytes(row_group_metadata, 9), BoundaryOrder::ASCENDING, ); - assert_eq!(row_group_offset_indexes[9].page_locations.len(), 352); + assert_eq!( + row_group_offset_indexes[9] + .as_ref() + .unwrap() + .page_locations + .len(), + 352 + ); } else { unreachable!() }; //col11->timestamp_col: INT96 UNCOMPRESSED DO:0 FPO:490093 SZ:111948/111948/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[num_nulls: 0, min/max not defined] - //Notice: min_max values for each page for this col not exits. - assert!(!&column_index[0][10].is_sorted()); - if column_index[0][10] == ColumnIndexMetaData::NONE { - assert_eq!(row_group_offset_indexes[10].page_locations.len(), 974); - } else { - unreachable!() - }; + // this columns lacks an index + assert!(column_index[0][10].is_none()); //col12->year: INT32 UNCOMPRESSED DO:0 FPO:602041 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 2009, max: 2010, num_nulls: 0] - assert!(&column_index[0][11].is_sorted()); - if let ColumnIndexMetaData::INT32(index) = &column_index[0][11] { + let ci = column_index[0][11].as_ref().unwrap(); + assert!(ci.is_sorted()); + if let ColumnIndexMetaData::INT32(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 11), BoundaryOrder::ASCENDING, ); - assert_eq!(row_group_offset_indexes[11].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[11] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; //col13->month: INT32 UNCOMPRESSED DO:0 FPO:639366 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 1, max: 12, num_nulls: 0] - assert!(!&column_index[0][12].is_sorted()); - if let ColumnIndexMetaData::INT32(index) = &column_index[0][12] { + let ci = column_index[0][12].as_ref().unwrap(); + assert!(!ci.is_sorted()); + if let ColumnIndexMetaData::INT32(index) = ci { check_native_page_index( index, 325, get_row_group_min_max_bytes(row_group_metadata, 12), BoundaryOrder::UNORDERED, ); - assert_eq!(row_group_offset_indexes[12].page_locations.len(), 325); + assert_eq!( + row_group_offset_indexes[12] + .as_ref() + .unwrap() + .page_locations + .len(), + 325 + ); } else { unreachable!() }; @@ -2631,7 +2731,7 @@ mod tests { assert_eq!(c.len(), 1); match &c[0] { - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v) => { + Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v)) => { assert_eq!(v.num_pages(), 1); assert_eq!(v.null_count(0).unwrap(), 1); assert_eq!(v.min_value(0).unwrap(), &[0; 11]); @@ -2766,7 +2866,7 @@ mod tests { // test that we got the index matching the row group match pg_idx { - ColumnIndexMetaData::INT32(int_idx) => { + Some(ColumnIndexMetaData::INT32(int_idx)) => { let min = col_stats.min_bytes_opt().unwrap().get_i32_le(); let max = col_stats.max_bytes_opt().unwrap().get_i32_le(); assert_eq!(int_idx.min_value(0), Some(min).as_ref()); @@ -2777,7 +2877,7 @@ mod tests { // check offset index matches too assert_eq!( - off_idx_i.page_locations[0].offset, + off_idx_i.as_ref().unwrap().page_locations[0].offset, metadata.row_group(0).column(0).data_page_offset() ); @@ -2811,7 +2911,7 @@ mod tests { // test that we got the index matching the row group match pg_idx { - ColumnIndexMetaData::INT32(int_idx) => { + Some(ColumnIndexMetaData::INT32(int_idx)) => { let min = col_stats.min_bytes_opt().unwrap().get_i32_le(); let max = col_stats.max_bytes_opt().unwrap().get_i32_le(); assert_eq!(int_idx.min_value(0), Some(min).as_ref()); @@ -2822,7 +2922,7 @@ mod tests { // check offset index matches too assert_eq!( - off_idx_i.page_locations[0].offset, + off_idx_i.as_ref().unwrap().page_locations[0].offset, metadata.row_group(i).column(0).data_page_offset() ); } diff --git a/parquet/src/file/writer.rs b/parquet/src/file/writer.rs index cc2e36b50fd2..9838fd8f1b01 100644 --- a/parquet/src/file/writer.rs +++ b/parquet/src/file/writer.rs @@ -2197,9 +2197,12 @@ mod tests { assert_eq!(column_index[0].len(), 2); // 2 column let a_idx = &column_index[0][0]; - assert!(matches!(a_idx, ColumnIndexMetaData::INT32(_)), "{a_idx:?}"); + assert!( + matches!(a_idx, Some(ColumnIndexMetaData::INT32(_))), + "{a_idx:?}" + ); let b_idx = &column_index[0][1]; - assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}"); + assert!(b_idx.is_none(), "{b_idx:?}"); } #[test] @@ -2279,7 +2282,7 @@ mod tests { let column_index = reader.metadata().column_index().unwrap(); assert_eq!(column_index.len(), 1); assert_eq!(column_index[0].len(), 1); - let col_idx = if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][0] { + let col_idx = if let Some(ColumnIndexMetaData::BYTE_ARRAY(index)) = &column_index[0][0] { assert_eq!(index.num_pages(), 1); index } else { @@ -2294,8 +2297,17 @@ mod tests { let offset_index = reader.metadata().offset_index().unwrap(); assert_eq!(offset_index.len(), 1); assert_eq!(offset_index[0].len(), 1); - assert!(offset_index[0][0].unencoded_byte_array_data_bytes.is_some()); + assert!(offset_index[0][0].is_some()); + assert!( + offset_index[0][0] + .as_ref() + .unwrap() + .unencoded_byte_array_data_bytes + .is_some() + ); let page_sizes = offset_index[0][0] + .as_ref() + .unwrap() .unencoded_byte_array_data_bytes .as_ref() .unwrap(); @@ -2474,7 +2486,7 @@ mod tests { let column_index = reader.metadata().column_index().unwrap(); assert_eq!(column_index.len(), 1); assert_eq!(column_index[0].len(), 1); - let col_idx = if let ColumnIndexMetaData::INT32(index) = &column_index[0][0] { + let col_idx = if let Some(ColumnIndexMetaData::INT32(index)) = &column_index[0][0] { assert_eq!(index.num_pages(), 1); index } else { @@ -2488,7 +2500,13 @@ mod tests { let offset_index = reader.metadata().offset_index().unwrap(); assert_eq!(offset_index.len(), 1); assert_eq!(offset_index[0].len(), 1); - assert!(offset_index[0][0].unencoded_byte_array_data_bytes.is_none()); + assert!( + offset_index[0][0] + .as_ref() + .unwrap() + .unencoded_byte_array_data_bytes + .is_none() + ); } #[test] @@ -2658,8 +2676,11 @@ mod tests { let rg_offset_indexes = offset_indexes.and_then(|oi| oi.get(rg_idx)); let mut rg_out = writer.next_row_group().unwrap(); for (col_idx, column) in rg.columns().iter().enumerate() { - let column_index = rg_column_indexes.and_then(|row| row.get(col_idx)).cloned(); - let offset_index = rg_offset_indexes.and_then(|row| row.get(col_idx)).cloned(); + let column_index = + rg_column_indexes.and_then(|row| row.get(col_idx).and_then(|c| c.clone())); + let offset_index = + rg_offset_indexes.and_then(|row| row.get(col_idx).and_then(|o| o.clone())); + let result = ColumnCloseResult { bytes_written: column.compressed_size() as _, rows_written: rg.num_rows() as _, @@ -2704,7 +2725,7 @@ mod tests { let col_idx = metadata.column_index().expect("column index not present"); let col0 = match &col_idx[0][0] { - ColumnIndexMetaData::INT96(index) => index, + Some(ColumnIndexMetaData::INT96(index)) => index, _ => panic!("expected INT96 stats"), }; let col_min = col0.min_value(0).expect("ColumnIndex min not present"); diff --git a/parquet/tests/arrow_reader/io/mod.rs b/parquet/tests/arrow_reader/io/mod.rs index 7e50b2c4bd9b..add87cc65e5d 100644 --- a/parquet/tests/arrow_reader/io/mod.rs +++ b/parquet/tests/arrow_reader/io/mod.rs @@ -354,7 +354,10 @@ impl TestRowGroups { .enumerate() .map(|(col_idx, col_meta)| { let column_name = col_meta.column_descr().name().to_string(); - let page_locations = offset_index[rg_index][col_idx].page_locations(); + let page_locations = offset_index[rg_index][col_idx] + .as_ref() + .unwrap() + .page_locations(); let dictionary_page_location = col_meta.dictionary_page_offset(); // We can find the byte range of the entire column chunk diff --git a/parquet/tests/arrow_reader/row_filter/async.rs b/parquet/tests/arrow_reader/row_filter/async.rs index 059755cb7bb3..cbcf404256b2 100644 --- a/parquet/tests/arrow_reader/row_filter/async.rs +++ b/parquet/tests/arrow_reader/row_filter/async.rs @@ -315,6 +315,8 @@ async fn test_mask_nested_projection_with_different_page_boundaries() { .iter() .map(|column| { column + .as_ref() + .unwrap() .page_locations() .iter() .map(|page| page.first_row_index) diff --git a/parquet/tests/arrow_writer/layout.rs b/parquet/tests/arrow_writer/layout.rs index 1c63a3144391..9bde70bfa7f3 100644 --- a/parquet/tests/arrow_writer/layout.rs +++ b/parquet/tests/arrow_writer/layout.rs @@ -94,11 +94,13 @@ fn assert_layout(file_reader: &Bytes, meta: &ParquetMetaData, layout: &Layout) { for (column_index, column_layout) in offset_index.iter().zip(&row_group_layout.columns) { assert_eq!( - column_index.page_locations.len(), + column_index.as_ref().unwrap().page_locations.len(), column_layout.pages.len(), "index page count mismatch" ); for (idx, (page, page_layout)) in column_index + .as_ref() + .unwrap() .page_locations .iter() .zip(&column_layout.pages) @@ -110,6 +112,8 @@ fn assert_layout(file_reader: &Bytes, meta: &ParquetMetaData, layout: &Layout) { "index page {idx} size mismatch" ); let next_first_row_index = column_index + .as_ref() + .unwrap() .page_locations .get(idx + 1) .map(|x| x.first_row_index) @@ -598,8 +602,8 @@ fn test_per_column_data_page_size_limit() { // Get page counts from offset index let offset_index = metadata.offset_index().unwrap(); - let col_a_page_count = offset_index[0][0].page_locations.len(); - let col_b_page_count = offset_index[0][1].page_locations.len(); + let col_a_page_count = offset_index[0][0].as_ref().unwrap().page_locations.len(); + let col_b_page_count = offset_index[0][1].as_ref().unwrap().page_locations.len(); // col_a should have many more pages than col_b due to smaller page size limit // col_a: 500 byte limit for 8000 bytes of data -> 16 pages diff --git a/parquet/tests/encryption/encryption_util.rs b/parquet/tests/encryption/encryption_util.rs index daf7e07b7bc2..758b852856f1 100644 --- a/parquet/tests/encryption/encryption_util.rs +++ b/parquet/tests/encryption/encryption_util.rs @@ -188,8 +188,8 @@ pub(crate) fn verify_column_indexes(metadata: &ParquetMetaData) { // Check float column, which is encrypted in the non-uniform test file let float_col_idx = 4; let offset_index = &offset_index[0][float_col_idx]; - assert_eq!(offset_index.page_locations.len(), 1); - assert!(offset_index.page_locations[0].offset > 0); + assert_eq!(offset_index.as_ref().unwrap().page_locations.len(), 1); + assert!(offset_index.as_ref().unwrap().page_locations[0].offset > 0); let column_index = metadata.column_index().unwrap(); assert_eq!(column_index.len(), 1); @@ -197,7 +197,7 @@ pub(crate) fn verify_column_indexes(metadata: &ParquetMetaData) { let column_index = &column_index[0][float_col_idx]; match column_index { - parquet::file::page_index::column_index::ColumnIndexMetaData::FLOAT(float_index) => { + Some(parquet::file::page_index::column_index::ColumnIndexMetaData::FLOAT(float_index)) => { assert_eq!(float_index.num_pages(), 1); assert_eq!(float_index.min_value(0), Some(&0.0f32)); assert!( From 250644d04e05ad64e595510208a9682a8329cc91 Mon Sep 17 00:00:00 2001 From: seidl Date: Tue, 11 Aug 2026 14:42:58 -0700 Subject: [PATCH 2/3] remove the NONE variant from ColumnIndexMetaData --- parquet/src/bin/parquet-index.rs | 2 +- parquet/src/file/metadata/memory.rs | 1 - parquet/src/file/metadata/writer.rs | 40 ++++++++------------- parquet/src/file/page_index/column_index.rs | 18 ---------- 4 files changed, 15 insertions(+), 46 deletions(-) diff --git a/parquet/src/bin/parquet-index.rs b/parquet/src/bin/parquet-index.rs index ae99d9841c2b..7b6026c24e6e 100644 --- a/parquet/src/bin/parquet-index.rs +++ b/parquet/src/bin/parquet-index.rs @@ -108,7 +108,7 @@ impl Args { let row_counts = compute_row_counts(offset_index.page_locations(), row_group.num_rows()); match column_indices[column_idx].as_ref() { - None | Some(ColumnIndexMetaData::NONE) => println!("NO INDEX"), + None => println!("NO INDEX"), Some(ColumnIndexMetaData::BOOLEAN(v)) => { print_index::(v, offset_index, &row_counts)? } diff --git a/parquet/src/file/metadata/memory.rs b/parquet/src/file/metadata/memory.rs index 8a5937c43c65..b554dd48793b 100644 --- a/parquet/src/file/metadata/memory.rs +++ b/parquet/src/file/metadata/memory.rs @@ -242,7 +242,6 @@ impl HeapSize for OffsetIndexMetaData { impl HeapSize for ColumnIndexMetaData { fn heap_size(&self) -> usize { match self { - Self::NONE => 0, Self::BOOLEAN(native_index) => native_index.heap_size(), Self::INT32(native_index) => native_index.heap_size(), Self::INT64(native_index) => native_index.heap_size(), diff --git a/parquet/src/file/metadata/writer.rs b/parquet/src/file/metadata/writer.rs index 574522c1fdb1..54b443ac23f8 100644 --- a/parquet/src/file/metadata/writer.rs +++ b/parquet/src/file/metadata/writer.rs @@ -527,14 +527,8 @@ impl MetadataObjectWriter { _column_idx: usize, sink: impl Write, ) -> Result { - match column_index { - // Missing indexes may also have the placeholder ColumnIndexMetaData::NONE - ColumnIndexMetaData::NONE => Ok(false), - _ => { - Self::write_thrift_object(column_index, sink)?; - Ok(true) - } - } + Self::write_thrift_object(column_index, sink)?; + Ok(true) } /// No-op implementation of row-group metadata encryption @@ -624,25 +618,19 @@ impl MetadataObjectWriter { column_idx: usize, sink: impl Write, ) -> Result { - match column_index { - // Missing indexes may also have the placeholder ColumnIndexMetaData::NONE - ColumnIndexMetaData::NONE => Ok(false), - _ => { - match &self.file_encryptor { - Some(file_encryptor) => Self::write_thrift_object_with_encryption( - column_index, - sink, - file_encryptor, - column_chunk, - ModuleType::ColumnIndex, - row_group_idx, - column_idx, - )?, - None => Self::write_thrift_object(column_index, sink)?, - } - Ok(true) - } + match &self.file_encryptor { + Some(file_encryptor) => Self::write_thrift_object_with_encryption( + column_index, + sink, + file_encryptor, + column_chunk, + ModuleType::ColumnIndex, + row_group_idx, + column_idx, + )?, + None => Self::write_thrift_object(column_index, sink)?, } + Ok(true) } /// If encryption is enabled and configured, encrypt row group metadata. diff --git a/parquet/src/file/page_index/column_index.rs b/parquet/src/file/page_index/column_index.rs index 665b4bc454b0..7889ce48aa8e 100644 --- a/parquet/src/file/page_index/column_index.rs +++ b/parquet/src/file/page_index/column_index.rs @@ -531,11 +531,6 @@ macro_rules! colidx_enum_func { Self::DOUBLE(ref typed) => typed.$func($arg), Self::BYTE_ARRAY(ref typed) => typed.$func($arg), Self::FIXED_LEN_BYTE_ARRAY(ref typed) => typed.$func($arg), - _ => panic!(concat!( - "Cannot call ", - stringify!($func), - " on ColumnIndexMetaData::NONE" - )), } }}; ($self:ident, $func:ident) => {{ @@ -548,11 +543,6 @@ macro_rules! colidx_enum_func { Self::DOUBLE(ref typed) => typed.$func(), Self::BYTE_ARRAY(ref typed) => typed.$func(), Self::FIXED_LEN_BYTE_ARRAY(ref typed) => typed.$func(), - _ => panic!(concat!( - "Cannot call ", - stringify!($func), - " on ColumnIndexMetaData::NONE" - )), } }}; } @@ -566,10 +556,6 @@ macro_rules! colidx_enum_func { #[derive(Debug, Clone, PartialEq)] #[allow(non_camel_case_types)] pub enum ColumnIndexMetaData { - /// Sometimes reading page index from parquet file - /// will only return pageLocations without min_max index, - /// `NONE` represents this lack of index information - NONE, /// Boolean type index BOOLEAN(PrimitiveColumnIndex), /// 32-bit integer type index @@ -602,7 +588,6 @@ impl ColumnIndexMetaData { /// Get boundary_order of this page index. pub fn get_boundary_order(&self) -> Option { match self { - Self::NONE => None, Self::BOOLEAN(index) => Some(index.boundary_order), Self::INT32(index) => Some(index.boundary_order), Self::INT64(index) => Some(index.boundary_order), @@ -619,7 +604,6 @@ impl ColumnIndexMetaData { /// Returns `None` if no null counts have been set in the index pub fn null_counts(&self) -> Option<&Vec> { match self { - Self::NONE => None, Self::BOOLEAN(index) => index.null_counts.as_ref(), Self::INT32(index) => index.null_counts.as_ref(), Self::INT64(index) => index.null_counts.as_ref(), @@ -636,7 +620,6 @@ impl ColumnIndexMetaData { /// Returns `None` if no NaN counts have been set in the index pub fn nan_counts(&self) -> Option<&Vec> { match self { - Self::NONE => None, Self::BOOLEAN(index) => index.nan_counts.as_ref(), Self::INT32(index) => index.nan_counts.as_ref(), Self::INT64(index) => index.nan_counts.as_ref(), @@ -752,7 +735,6 @@ impl WriteThrift for ColumnIndexMetaData { ColumnIndexMetaData::DOUBLE(index) => index.write_thrift(writer), ColumnIndexMetaData::BYTE_ARRAY(index) => index.write_thrift(writer), ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => index.write_thrift(writer), - _ => Err(general_err!("Cannot serialize NONE index")), } } } From 511ffbab07821eb25a712532ab6dab6bce51fb81 Mon Sep 17 00:00:00 2001 From: seidl Date: Tue, 11 Aug 2026 15:00:24 -0700 Subject: [PATCH 3/3] doc fix --- parquet/src/file/metadata/writer.rs | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/parquet/src/file/metadata/writer.rs b/parquet/src/file/metadata/writer.rs index 54b443ac23f8..f536abc47c51 100644 --- a/parquet/src/file/metadata/writer.rs +++ b/parquet/src/file/metadata/writer.rs @@ -515,8 +515,7 @@ impl MetadataObjectWriter { /// Write a column [`ColumnIndex`] in Thrift format /// - /// If `column_index` is [`ColumnIndexMetaData::NONE`] the index will not be written and - /// this will return `false`. Returns `true` otherwise. + /// Returns `true` unless there is an error. /// /// [`ColumnIndex`]: https://github.com/apache/parquet-format/blob/master/PageIndex.md fn write_column_index( @@ -606,8 +605,7 @@ impl MetadataObjectWriter { /// Write a column [`ColumnIndex`] in Thrift format, possibly encrypting it if required /// - /// If `column_index` is [`ColumnIndexMetaData::NONE`] the index will not be written and - /// this will return `false`. Returns `true` otherwise. + /// Returns `true` unless there is an error. /// /// [`ColumnIndex`]: https://github.com/apache/parquet-format/blob/master/PageIndex.md fn write_column_index(