diff --git a/parquet/benches/arrow_statistics.rs b/parquet/benches/arrow_statistics.rs index 389e6f1f32e4..3c78dee0e575 100644 --- a/parquet/benches/arrow_statistics.rs +++ b/parquet/benches/arrow_statistics.rs @@ -226,33 +226,16 @@ fn criterion_benchmark(c: &mut Criterion) { .unwrap(); if data_page_row_count_limit.is_some() { - let column_page_index = reader + let page_index = reader .metadata() - .column_index() - .expect("File should have column page indices"); + .page_index() + .expect("File should have page indices"); - let column_offset_index = reader - .metadata() - .offset_index() - .expect("File should have column offset indices"); - - let _ = converter.data_page_mins( - column_page_index, - column_offset_index, - &row_group_indices, - ); - let _ = converter.data_page_maxes( - column_page_index, - column_offset_index, - &row_group_indices, - ); - let _ = converter.data_page_null_counts( - column_page_index, - column_offset_index, - &row_group_indices, - ); + let _ = converter.data_page_mins(page_index, &row_group_indices); + let _ = converter.data_page_maxes(page_index, &row_group_indices); + let _ = converter.data_page_null_counts(page_index, &row_group_indices); let _ = converter.data_page_row_counts( - column_offset_index, + page_index, row_groups, &row_group_indices, ); diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index 52e7461835d9..6374b46155de 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -1313,12 +1313,11 @@ impl ReaderPageIterator { fn next_page_reader(&self, rg_idx: usize) -> Result> { let rg = self.metadata.row_group(rg_idx); let column_chunk_metadata = rg.column(self.column_idx); - let offset_index = self.metadata.offset_index(); - // `offset_index` may not exist and `i[rg_idx]` will be empty. - // 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()); + let page_locations = self + .metadata + .page_index() + .map(|i| i.page_locations(rg_idx, self.column_idx).cloned()) + .unwrap_or(None); let total_rows = rg.num_rows() as usize; let reader = self.reader.clone(); @@ -4686,7 +4685,19 @@ pub(crate) mod tests { ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required), ) .unwrap(); - assert!(!builder.metadata().offset_index().unwrap()[0].is_empty()); + let page_index = builder + .metadata() + .page_index() + .expect("page index should be present"); + let num_columns = builder.metadata().row_group(0).num_columns(); + let offset_indexes = page_index.offset_indexes_for_rowgroup(0); + assert!(offset_indexes.is_some_and(|ois| ois.len() == num_columns)); + let column_indexes = page_index.offset_indexes_for_rowgroup(0); + assert!(column_indexes.is_some_and(|cis| cis.len() == num_columns)); + assert!(page_index.offset_index(0, 0).is_some()); + assert!(page_index.column_index(0, 0).is_some()); + assert!(page_index.page_locations(0, 0).is_some()); + assert_eq!(page_index.num_data_pages(0, 0), Some(325)); let reader = builder.build().unwrap(); let batches = reader.collect::, _>>().unwrap(); assert_eq!(batches.len(), 8); @@ -4703,7 +4714,7 @@ pub(crate) mod tests { .unwrap(); // Although `Vec>` of each row group is empty, // we should read the file successfully. - assert!(builder.metadata().offset_index().is_none()); + assert!(builder.metadata().page_index().is_none()); let reader = builder.build().unwrap(); let batches = reader.collect::, _>>().unwrap(); assert_eq!(batches.len(), 1); diff --git a/parquet/src/arrow/arrow_reader/statistics.rs b/parquet/src/arrow/arrow_reader/statistics.rs index dd6ef3607989..c08806d48cf8 100644 --- a/parquet/src/arrow/arrow_reader/statistics.rs +++ b/parquet/src/arrow/arrow_reader/statistics.rs @@ -23,7 +23,7 @@ use crate::arrow::buffer::bit_util::sign_extend_be; use crate::arrow::parquet_column; use crate::basic::Type as PhysicalType; use crate::errors::{ParquetError, Result}; -use crate::file::metadata::{ParquetColumnIndex, ParquetOffsetIndex, RowGroupMetaData}; +use crate::file::metadata::{PageIndex, RowGroupMetaData}; use crate::file::page_index::column_index::ColumnIndexMetaData; use crate::file::statistics::Statistics as ParquetStatistics; use crate::schema::types::SchemaDescriptor; @@ -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); } @@ -1847,7 +1847,7 @@ impl<'a> StatisticsConverter<'a> { /// In Parquet files, in addition to the Column Chunk level statistics /// (stored for each column for each row group) there are also /// optional statistics stored for each data page, as part of - /// the [`ParquetColumnIndex`]. + /// the [`PageIndex`]. /// /// Since a single Column Chunk is stored as one or more pages, /// page level statistics can prune at a finer granularity. @@ -1858,11 +1858,7 @@ impl<'a> StatisticsConverter<'a> { /// /// # Parameters: /// - /// * `column_page_index`: The parquet column page indices, read from - /// `ParquetMetaData` column_index - /// - /// * `column_offset_index`: The parquet column offset indices, read from - /// `ParquetMetaData` offset_index + /// * `page_index`: The parquet page indices, read from `ParquetMetaData` /// /// * `row_group_indices`: The indices of the row groups, that are used to /// extract the column page index and offset index on a per row group @@ -1895,8 +1891,7 @@ impl<'a> StatisticsConverter<'a> { /// * the stored statistic value can not be converted to the requested type pub fn data_page_mins( &self, - column_page_index: &ParquetColumnIndex, - column_offset_index: &ParquetOffsetIndex, + page_index: &PageIndex, row_group_indices: I, ) -> Result where @@ -1910,12 +1905,12 @@ 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(); + page_index.column_index(*rg_index, parquet_index); + let num_data_pages = page_index + .num_data_pages(*rg_index, parquet_index) + .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) @@ -1926,8 +1921,7 @@ impl<'a> StatisticsConverter<'a> { /// See docs on [`Self::data_page_mins`] for details. pub fn data_page_maxes( &self, - column_page_index: &ParquetColumnIndex, - column_offset_index: &ParquetOffsetIndex, + page_index: &PageIndex, row_group_indices: I, ) -> Result where @@ -1941,12 +1935,12 @@ 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(); + page_index.column_index(*rg_index, parquet_index); + let num_data_pages = page_index + .num_data_pages(*rg_index, parquet_index) + .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) @@ -1957,8 +1951,7 @@ impl<'a> StatisticsConverter<'a> { /// See docs on [`Self::data_page_mins`] for details. pub fn data_page_null_counts( &self, - column_page_index: &ParquetColumnIndex, - column_offset_index: &ParquetOffsetIndex, + page_index: &PageIndex, row_group_indices: I, ) -> Result where @@ -1971,12 +1964,12 @@ 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(); + page_index.column_index(*rg_index, parquet_index); + let num_data_pages = page_index + .num_data_pages(*rg_index, parquet_index) + .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) } @@ -1986,8 +1979,7 @@ impl<'a> StatisticsConverter<'a> { /// See docs on [`Self::data_page_mins`] for details. pub fn data_page_nan_counts( &self, - column_page_index: &ParquetColumnIndex, - column_offset_index: &ParquetOffsetIndex, + page_index: &PageIndex, row_group_indices: I, ) -> Result where @@ -2000,12 +1992,12 @@ 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(); + page_index.column_index(*rg_index, parquet_index); + let num_data_pages = page_index + .num_data_pages(*rg_index, parquet_index) + .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) } @@ -2029,7 +2021,7 @@ impl<'a> StatisticsConverter<'a> { /// See docs on [`Self::data_page_mins`] for details. pub fn data_page_row_counts( &self, - column_offset_index: &ParquetOffsetIndex, + page_index: &PageIndex, row_group_metadatas: &'a [RowGroupMetaData], row_group_indices: I, ) -> Result> @@ -2046,7 +2038,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) = page_index.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 947f07a6270e..c49c919fa0d4 100644 --- a/parquet/src/arrow/arrow_writer/mod.rs +++ b/parquet/src/arrow/arrow_writer/mod.rs @@ -3119,10 +3119,13 @@ mod tests { "Expected a dictionary page" ); - 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_index = reader + .metadata() + .page_index() + .expect("page index should be present"); + let page_locations = page_index + .page_locations(0, 0) + .expect("page locations should exist"); // 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. @@ -3526,10 +3529,12 @@ mod tests { assert!(column.column_index_length().is_some()); } } - assert!(file_meta_data.column_index().is_some()); - if let Some(col_indexes) = file_meta_data.column_index() { - for rg_idx in col_indexes { - for idx in rg_idx { + if let Some(page_index) = file_meta_data.page_index() { + for rg in 0..file_meta_data.num_row_groups() { + for col in 0..file_meta_data.row_group(rg).num_columns() { + let idx = page_index + .column_index(rg, col) + .expect("column index should exist"); assert!(idx.nan_counts().is_some()); let ColumnIndexMetaData::DOUBLE(float_idx) = idx else { panic!("expected double statistics") @@ -3547,6 +3552,8 @@ mod tests { } } } + } else { + panic!("page index should be present"); } } @@ -3604,12 +3611,12 @@ mod tests { assert_eq!(col_stats.min_bytes_opt(), Some((-1.0f64).as_bytes())); assert_eq!(col_stats.max_bytes_opt(), Some(1.0f64.as_bytes())); - 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!(file_meta_data.page_index().is_some()); + let col_idx = &file_meta_data.page_index().unwrap().column_index(0, 0); + assert_eq!(col_idx.as_ref().unwrap().num_pages(), 4); // test each page - let ColumnIndexMetaData::DOUBLE(float_idx) = col_idx else { + let Some(ColumnIndexMetaData::DOUBLE(float_idx)) = col_idx else { panic!("expected double statistics") }; @@ -5099,12 +5106,10 @@ mod tests { let bytes = Bytes::from(buf); let options = ReadOptionsBuilder::new().with_page_index().build(); let reader = SerializedFileReader::new_with_options(bytes, options).unwrap(); - let index = reader.metadata().offset_index().unwrap(); + let index = reader.metadata().page_index().unwrap(); - 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.num_data_pages(0, 0), Some(1)); // 1 page + assert_eq!(index.num_data_pages(0, 1), Some(1)); // 1 page } #[test] @@ -5165,21 +5170,15 @@ mod tests { // The column chunk for column "b" shouldn't have statistics assert!(b_col.statistics().is_none()); - let offset_index = reader.metadata().offset_index().unwrap(); - assert_eq!(offset_index.len(), 1); // 1 row group - assert_eq!(offset_index[0].len(), 2); // 2 columns - - let column_index = reader.metadata().column_index().unwrap(); - assert_eq!(column_index.len(), 1); // 1 row group - assert_eq!(column_index[0].len(), 2); // 2 columns + let page_index = reader.metadata().page_index().unwrap(); - let a_idx = &column_index[0][0]; + let a_idx = page_index.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:?}"); + let b_idx = page_index.column_index(0, 1); + assert!(b_idx.is_none(), "{b_idx:?}"); } #[test] @@ -5240,14 +5239,12 @@ mod tests { // The column chunk for column "b" shouldn't have statistics assert!(b_col.statistics().is_none()); - let column_index = reader.metadata().column_index().unwrap(); - assert_eq!(column_index.len(), 1); // 1 row group - assert_eq!(column_index[0].len(), 2); // 2 columns + let page_index = reader.metadata().page_index().unwrap(); - let a_idx = &column_index[0][0]; - assert!(matches!(a_idx, ColumnIndexMetaData::NONE), "{a_idx:?}"); - let b_idx = &column_index[0][1]; - assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}"); + let a_idx = page_index.column_index(0, 0); + assert!(a_idx.is_none(), "{a_idx:?}"); + let b_idx = page_index.column_index(0, 1); + assert!(b_idx.is_none(), "{b_idx:?}"); } #[test] diff --git a/parquet/src/arrow/async_reader/mod.rs b/parquet/src/arrow/async_reader/mod.rs index 0bff84b3d836..85c8fa463a25 100644 --- a/parquet/src/arrow/async_reader/mod.rs +++ b/parquet/src/arrow/async_reader/mod.rs @@ -942,8 +942,8 @@ mod tests { use crate::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions}; use crate::arrow::schema::virtual_type::RowNumber; use crate::arrow::{ArrowWriter, AsyncArrowWriter, ProjectionMask}; - use crate::file::metadata::PageIndexPolicy; use crate::file::metadata::ParquetMetaDataReader; + use crate::file::metadata::{PageIndex, PageIndexPolicy}; use crate::file::properties::WriterProperties; use arrow::compute::kernels::cmp::eq; use arrow::error::Result as ArrowResult; @@ -1126,25 +1126,23 @@ mod tests { let metadata_with_index = builder.metadata(); assert_eq!(metadata_with_index.num_row_groups(), 1); - // Check offset indexes are present for all columns - let offset_index = metadata_with_index.offset_index().unwrap(); - let column_index = metadata_with_index.column_index().unwrap(); - - assert_eq!(offset_index.len(), metadata_with_index.num_row_groups()); - assert_eq!(column_index.len(), metadata_with_index.num_row_groups()); - + // Check offset indexes are present for all columns of all row groups + let page_index = metadata_with_index.page_index().unwrap(); + let num_rowgroups = metadata_with_index.num_row_groups(); let num_columns = metadata_with_index .file_metadata() .schema_descr() .num_columns(); - - // Check page indexes are present for all columns - offset_index - .iter() - .for_each(|x| assert_eq!(x.len(), num_columns)); - column_index - .iter() - .for_each(|x| assert_eq!(x.len(), num_columns)); + for rgidx in 0..num_rowgroups { + let column_index = page_index.column_indexes_for_rowgroup(rgidx); + let offset_index = page_index.offset_indexes_for_rowgroup(rgidx); + assert!(column_index.is_some_and(|ci| ci.len() == num_columns)); + assert!(offset_index.is_some_and(|oi| oi.len() == num_columns)); + // some column indexes are not defined, but all offset indexes should be + for colidx in 0..num_columns { + assert!(page_index.offset_index(rgidx, colidx).is_some()); + } + } let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]); let stream = builder @@ -1714,7 +1712,8 @@ mod tests { .await .unwrap(); - metadata.set_offset_index(Some(vec![])); + let page_index = PageIndex::new(None, Some(vec![])); + metadata.set_page_index(Some(page_index)); let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required); let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap(); let reader = diff --git a/parquet/src/arrow/async_reader/store.rs b/parquet/src/arrow/async_reader/store.rs index d4f5beb817f9..b7fec0a8d31c 100644 --- a/parquet/src/arrow/async_reader/store.rs +++ b/parquet/src/arrow/async_reader/store.rs @@ -259,8 +259,8 @@ impl AsyncFileReader for ParquetObjectReader { #[cfg(test)] #[expect(deprecated)] mod tests { - use crate::arrow::async_reader::ArrowReaderOptions; use crate::file::metadata::PageIndexPolicy; + use crate::{arrow::async_reader::ArrowReaderOptions, file::metadata::PageIndex}; use std::sync::{ Arc, atomic::{AtomicUsize, Ordering}, @@ -447,7 +447,7 @@ mod tests { let metadata = reader.get_metadata(Some(&options)).await.unwrap(); // With preload=true, indexes should be loaded since the test file has them - assert!(metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } #[tokio::test] @@ -468,7 +468,7 @@ mod tests { // With Optional policy, it will TRY to load indexes but won't fail if they don't exist // The test file has page indexes, so they will be some - assert!(metadata.column_index().is_some()); + assert!(metadata.page_index().is_some()); } #[tokio::test] @@ -496,8 +496,8 @@ mod tests { // Both should succeed (no panic/error) // metadata1 (Skip) uses preload=false -> Skip policy // metadata2 (Optional) overrides preload=false -> Optional policy - assert!(metadata1.column_index().is_none()); - assert!(metadata2.column_index().is_some()); + assert!(metadata1.page_index().is_none()); + assert!(metadata2.page_index().is_some()); } #[tokio::test] @@ -516,6 +516,6 @@ mod tests { // With no options provided, preload flags (true) should be respected // and converted to Optional policy internally (preload=true -> Optional) // The test file has page indexes, so they will be some - assert!(metadata.column_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } } 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..8d80e6cf18b1 100644 --- a/parquet/src/arrow/push_decoder/reader_builder/data.rs +++ b/parquet/src/arrow/push_decoder/reader_builder/data.rs @@ -224,10 +224,9 @@ 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) - .filter(|index| !index.is_empty()) - .map(|x| x[row_group_idx].as_slice()) + .page_index() + .map(|pi| pi.offset_indexes_for_rowgroup(row_group_idx)) + .unwrap_or(None) } diff --git a/parquet/src/arrow/push_decoder/reader_builder/mod.rs b/parquet/src/arrow/push_decoder/reader_builder/mod.rs index 89b6d58fbb3b..92ebdeee3cc9 100644 --- a/parquet/src/arrow/push_decoder/reader_builder/mod.rs +++ b/parquet/src/arrow/push_decoder/reader_builder/mod.rs @@ -851,12 +851,14 @@ 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()) - .and_then(|index| index.get(row_group_idx)) - .map(|columns| columns.as_slice()) + .page_index() + .map(|pi| pi.offset_indexes_for_rowgroup(row_group_idx)) + .unwrap_or(None) } } @@ -883,7 +885,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() { @@ -908,7 +910,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?; @@ -918,7 +920,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 @@ -947,17 +950,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![ @@ -980,7 +985,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() @@ -991,7 +996,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..3949be2d597b 100644 --- a/parquet/src/bin/parquet-concat.rs +++ b/parquet/src/bin/parquet-concat.rs @@ -100,18 +100,18 @@ impl Args { let mut writer = SerializedFileWriter::new(output, schema, props)?; for (input, metadata) in inputs { - let column_indexes = metadata.column_index(); - let offset_indexes = metadata.offset_index(); + let page_index = metadata.page_index(); for (rg_idx, rg) in metadata.row_groups().iter().enumerate() { - let rg_column_indexes = column_indexes.and_then(|ci| ci.get(rg_idx)); - let rg_offset_indexes = offset_indexes.and_then(|oi| oi.get(rg_idx)); 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 offset_index = rg_offset_indexes.and_then(|row| row.get(col_idx)).cloned(); + let column_index = page_index + .and_then(|pi| pi.column_index(rg_idx, col_idx)) + .cloned(); + let offset_index = page_index + .and_then(|pi| pi.offset_index(rg_idx, col_idx)) + .cloned(); 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..72401aa11932 100644 --- a/parquet/src/bin/parquet-index.rs +++ b/parquet/src/bin/parquet-index.rs @@ -75,48 +75,39 @@ impl Args { ParquetError::General(format!("Failed to find column {}", self.column)) })?; - // Column index data for all row groups and columns - let column_index = reader + let page_index = reader .metadata() - .column_index() - .ok_or_else(|| ParquetError::General("Column index not found".to_string()))?; - - // Offset index data for all row groups and columns - let offset_index = reader - .metadata() - .offset_index() - .ok_or_else(|| ParquetError::General("Offset index not found".to_string()))?; + .page_index() + .ok_or_else(|| ParquetError::General("Page index not found".to_string()))?; // Iterate through each row group - for (row_group_idx, ((column_indices, offset_indices), row_group)) in column_index - .iter() - .zip(offset_index) - .zip(reader.metadata().row_groups()) - .enumerate() - { + for row_group_idx in 0..reader.metadata().num_row_groups() { + let row_group = reader.metadata().row_group(row_group_idx); println!("Row Group: {row_group_idx}"); - let offset_index = offset_indices.get(column_idx).ok_or_else(|| { - ParquetError::General(format!( - "No offset index for row group {row_group_idx} column chunk {column_idx}" - )) - })?; + let offset_index = page_index + .offset_index(row_group_idx, column_idx) + .ok_or_else(|| { + ParquetError::General(format!( + "No offset index 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 page_index.column_index(row_group_idx, column_idx) { + 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 +138,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 +177,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/memory.rs b/parquet/src/file/metadata/memory.rs index 8a5937c43c65..1325a39b6658 100644 --- a/parquet/src/file/metadata/memory.rs +++ b/parquet/src/file/metadata/memory.rs @@ -21,8 +21,8 @@ use crate::basic::{BoundaryOrder, ColumnOrder, CompressionCodec, Encoding, PageType}; use crate::data_type::private::ParquetValueType; use crate::file::metadata::{ - ColumnChunkMetaData, FileMetaData, KeyValue, PageEncodingStats, ParquetPageEncodingStats, - RowGroupMetaData, SortingColumn, + ColumnChunkMetaData, FileMetaData, KeyValue, PageEncodingStats, PageIndex, + ParquetPageEncodingStats, RowGroupMetaData, SortingColumn, }; use crate::file::page_index::column_index::{ ByteArrayColumnIndex, ColumnIndex, ColumnIndexMetaData, PrimitiveColumnIndex, @@ -233,6 +233,12 @@ impl HeapSize for Statistics { } } +impl HeapSize for PageIndex { + fn heap_size(&self) -> usize { + self.column_indexes.heap_size() + self.offset_indexes.heap_size() + } +} + impl HeapSize for OffsetIndexMetaData { fn heap_size(&self) -> usize { self.page_locations.heap_size() + self.unencoded_byte_array_data_bytes.heap_size() @@ -242,7 +248,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/mod.rs b/parquet/src/file/metadata/mod.rs index b19cf4fc8820..15a16829badd 100644 --- a/parquet/src/file/metadata/mod.rs +++ b/parquet/src/file/metadata/mod.rs @@ -134,36 +134,285 @@ use std::sync::Arc; pub use writer::ParquetMetaDataWriter; pub(crate) use writer::ThriftMetadataWriter; -/// Page level statistics for each column chunk of each row group. +/// Encapsulates the Parquet [Page Index] for efficient page-level data skipping /// -/// This structure is an in-memory representation of multiple [`ColumnIndex`] -/// structures in a parquet file footer, as described in the Parquet [PageIndex -/// documentation]. Each [`ColumnIndex`] holds statistics about all the pages in a -/// particular column chunk. +/// The Page Index is optional metadata that enables query engines to skip irrelevant +/// data pages during scans, significantly improving I/O efficiency. It consists of two +/// complementary structures: /// -/// `column_index[row_group_number][column_number]` holds the -/// [`ColumnIndex`] corresponding to column `column_number` of row group -/// `row_group_number`. +/// * **[`ColumnIndex`]**: Per-page min/max value boundaries that enable predicate-based +/// page filtering. Allows determining which pages might contain rows matching a query +/// predicate without reading the actual data pages. /// -/// For example `column_index[2][3]` holds the [`ColumnIndex`] for the fourth -/// column in the third row group of the parquet file. +/// * **[`OffsetIndex`]**: Physical locations and sizes of data pages, plus the first row +/// index of each page. Used to locate and read only the pages identified as relevant +/// by the ColumnIndex. /// -/// [PageIndex documentation]: https://github.com/apache/parquet-format/blob/master/PageIndex.md -/// [`ColumnIndex`]: crate::file::page_index::column_index::ColumnIndexMetaData -pub type ParquetColumnIndex = Vec>; - -/// [`OffsetIndexMetaData`] for each data page of each row group of each column +/// Together, these indexes enable: +/// - Single-row lookups reading only one data page per column (on sorted columns) +/// - Range scans reading only pages containing values in the query range +/// - Efficient cross-column filtering by skipping corresponding row ranges +/// +/// # Structure +/// +/// Within a Parquet file, both indexes are organized as a two-level structure, with +/// indexes arranged first by row group, and then column. The [`ColumnChunkMetaData`] +/// contains pointers to the indexes for a given column chunk, so they may be +/// populated piecemeal. This struct allows access either by row group index +/// ([Self::column_indexes_for_rowgroup], [Self::offset_indexes_for_rowgroup]) or +/// individual access by row group index and column number ([Self::column_index], +/// [Self::offset_index]). +/// +/// Each entry is `Option` because: +/// - The entire page index might be absent (old files, disabled during write) +/// - Individual columns might lack indexes (unsupported types, statistics disabled) +/// +/// # Example: Checking if Page Index is Available /// -/// This structure is the parsed representation of the [`OffsetIndex`] from the -/// Parquet file footer, as described in the Parquet [PageIndex documentation]. +/// ``` +/// use parquet::file::metadata::ParquetMetaData; +/// # use parquet::errors::Result; +/// +/// fn check_page_index_availability(metadata: &ParquetMetaData) -> Result<()> { +/// if let Some(page_index) = metadata.page_index() { +/// println!("Page index present:"); +/// println!(" Has offset indexes: {}", page_index.has_offset_indexes()); +/// println!(" Has column indexes: {}", page_index.has_column_indexes()); +/// +/// // Check availability for first row group, first column +/// if let Some(col_idx) = page_index.column_index(0, 0) { +/// println!(" Column index found for row group 0, column 0"); +/// println!(" Number of pages: {}", col_idx.num_pages()); +/// } /// -/// `offset_index[row_group_number][column_number]` holds -/// the [`OffsetIndexMetaData`] corresponding to column -/// `column_number`of row group `row_group_number`. +/// if let Some(offset_idx) = page_index.offset_index(0, 0) { +/// println!(" Offset index found for row group 0, column 0"); +/// println!(" Number of pages: {}", offset_idx.page_locations().len()); +/// } +/// } else { +/// println!("No page index available"); +/// } +/// Ok(()) +/// } +/// ``` +/// +/// # Example: Using Page Index for Predicate Pushdown /// -/// [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>; +/// ``` +/// use parquet::file::metadata::ParquetMetaData; +/// use parquet::file::page_index::column_index::ColumnIndexMetaData; +/// # use parquet::errors::Result; +/// +/// /// Identifies which pages in a column might contain values >= min_value +/// fn find_relevant_pages( +/// metadata: &ParquetMetaData, +/// row_group_idx: usize, +/// column_idx: usize, +/// min_value: i32, +/// ) -> Vec { +/// let mut relevant_pages = Vec::new(); +/// +/// let Some(page_index) = metadata.page_index() else { +/// // No page index - must read all pages +/// return relevant_pages; +/// }; +/// +/// let Some(column_index) = page_index.column_index(row_group_idx, column_idx) else { +/// // No column index - must read all pages +/// return relevant_pages; +/// }; +/// +/// // Check each page's statistics +/// match column_index { +/// ColumnIndexMetaData::INT32(index) => { +/// for (page_num, max_value) in index.max_values_iter().enumerate() { +/// // Page might contain matching rows if its max >= our min +/// if let Some(max) = max_value { +/// if *max >= min_value { +/// relevant_pages.push(page_num); +/// } +/// } +/// } +/// } +/// _ => { +/// // Wrong column type - read all pages +/// } +/// } +/// +/// relevant_pages +/// } +/// ``` +/// +/// [Page Index]: https://github.com/apache/parquet-format/blob/master/PageIndex.md +/// [`ColumnIndex`]: crate::file::page_index::column_index::ColumnIndexMetaData +/// [`OffsetIndex`]: crate::file::page_index::offset_index::OffsetIndexMetaData +#[derive(Debug, Clone, PartialEq)] +pub struct PageIndex { + column_indexes: Option>>>, + offset_indexes: Option>>>, +} + +impl PageIndex { + pub(crate) fn new( + column_indexes: Option>>>, + offset_indexes: Option>>>, + ) -> Self { + Self { + column_indexes, + offset_indexes, + } + } + + /// Returns `true` if offset index structures are present + /// + /// This indicates whether [`OffsetIndexMetaData`] structures were loaded or created. + /// Returns `true` even if some individual columns lack offset indexes. + /// + /// To check if a specific column has an offset index, use [`Self::offset_index`]. + pub fn has_offset_indexes(&self) -> bool { + self.offset_indexes.is_some() + } + + /// Returns `true` if column index structures are present + /// + /// This indicates whether [`ColumnIndexMetaData`] structures were loaded or created. + /// Returns `true` even if some individual columns lack column indexes. + /// + /// To check if a specific column has a column index, use [`Self::column_index`]. + pub fn has_column_indexes(&self) -> bool { + self.column_indexes.is_some() + } + + /// Returns `true` if both the offset and column index structures are present + /// + /// This is equivalent to both [`Self::has_offset_indexes`] and [`Self::has_column_indexes`] + /// returning `true`. + pub fn is_complete(&self) -> bool { + self.has_column_indexes() && self.has_offset_indexes() + } + + /// Returns column indexes for all columns in the specified row group + /// + /// Returns `None` if: + /// - Column indexes were not loaded or are not available + /// - The row group index is out of bounds + /// + /// Returns `Some(&[Option])` where: + /// - The slice length equals the number of columns in the row group + /// - Each element is `Some` if that column has statistics, `None` otherwise + pub fn column_indexes_for_rowgroup( + &self, + row_group_idx: usize, + ) -> Option<&[Option]> { + match self.column_indexes.as_ref() { + None => None, + Some(indexes) => indexes.get(row_group_idx).map(|ci| ci.as_slice()), + } + } + + /// Returns the column index for a specific row group and column + /// + /// This is the primary method for accessing page-level min/max statistics + /// used in predicate pushdown and page skipping optimizations. + /// + /// Returns: + /// * `Some(&ColumnIndexMetaData)` - Column index is available with statistics + /// * `None` - Index unavailable (not loaded, row group/column out of bounds, or no statistics) + pub fn column_index( + &self, + row_group_idx: usize, + column_idx: usize, + ) -> Option<&ColumnIndexMetaData> { + if let Some(column_indexes) = self.column_indexes.as_ref() { + let rg = column_indexes.get(row_group_idx)?; + rg.get(column_idx)?.as_ref() + } else { + None + } + } + + /// Returns offset indexes for all columns in the specified row group + /// + /// Returns `None` if: + /// - Offset indexes were not loaded or are not available + /// - The row group index is out of bounds + /// + /// Returns `Some(&[Option])` where: + /// - The slice length equals the number of columns in the row group + /// - Each element is `Some` if that column has location metadata, `None` otherwise + pub fn offset_indexes_for_rowgroup( + &self, + row_group_idx: usize, + ) -> Option<&[Option]> { + match self.offset_indexes.as_ref() { + None => None, + Some(indexes) => indexes.get(row_group_idx).map(|oi| oi.as_slice()), + } + } + + /// Returns the offset index for a specific row group and column + /// + /// This provides physical locations and sizes of data pages, enabling: + /// - Direct seeking to specific pages identified by column index filtering + /// - Reading only relevant pages without scanning entire column chunks + /// - Efficient cross-column row-based filtering + /// + /// Returns: + /// * `Some(&OffsetIndexMetaData)` - Offset index is available + /// * `None` - Index unavailable (not loaded, row group/column out of bounds) + pub fn offset_index( + &self, + row_group_idx: usize, + column_idx: usize, + ) -> Option<&OffsetIndexMetaData> { + if let Some(offset_indexes) = self.offset_indexes.as_ref() { + let rg = offset_indexes.get(row_group_idx)?; + rg.get(column_idx)?.as_ref() + } else { + None + } + } + + /// Returns the expected number of data pages for a specific column chunk + /// + /// This count includes only data pages, not dictionary pages or other metadata pages. + /// + /// Returns: + /// * `Some(usize)` - Number of data pages if any index is available + /// * `None` - No index information available for this column + pub fn num_data_pages(&self, row_group_idx: usize, column_idx: usize) -> Option { + match self.offset_index(row_group_idx, column_idx) { + Some(offset_index) => Some(offset_index.page_locations.len()), + None => Some(self.column_index(row_group_idx, column_idx)?.num_pages() as usize), + } + } + + /// Returns the physical locations of all data pages in a column chunk + /// + /// Each [`PageLocation`] contains: + /// - File offset where the page begins + /// - Compressed size of the page + /// - First row index within the row group + /// + /// This enables direct I/O to specific pages without reading the entire column chunk. + /// + /// Returns: + /// * `Some(&Vec)` - Vector of page locations if offset index exists + /// * `None` - Offset index not available + pub fn page_locations( + &self, + row_group_idx: usize, + column_idx: usize, + ) -> Option<&Vec> { + if let Some(offset_indexes) = self.offset_indexes.as_ref() { + let rg = offset_indexes.get(row_group_idx)?; + let off_idx = rg.get(column_idx)?.as_ref()?; + Some(off_idx.page_locations()) + } else { + None + } + } +} /// Parsed metadata for a single Parquet file /// @@ -174,7 +423,7 @@ pub type ParquetOffsetIndex = Vec>; /// The fields of this structure are: /// * [`FileMetaData`]: Information about the overall file (such as the schema) (See [`Self::file_metadata`]) /// * [`RowGroupMetaData`]: Information about each Row Group (see [`Self::row_groups`]) -/// * [`ParquetColumnIndex`] and [`ParquetOffsetIndex`]: Optional "Page Index" structures (see [`Self::column_index`] and [`Self::offset_index`]) +/// * [`PageIndex`]: Optional "Page Index" structures (see [`Self::page_index`]) /// /// This structure is read by the various readers in this crate or can be read /// directly from a file using the [`ParquetMetaDataReader`] struct. @@ -189,9 +438,7 @@ pub struct ParquetMetaData { /// Row group metadata row_groups: Vec, /// Page level index for each page in each column chunk - column_index: Option, - /// Offset index for each page in each column chunk - offset_index: Option, + page_index: Option, /// Optional file decryptor #[cfg(feature = "encryption")] file_decryptor: Option>, @@ -204,8 +451,7 @@ impl ParquetMetaData { ParquetMetaData { file_metadata, row_groups, - column_index: None, - offset_index: None, + page_index: None, #[cfg(feature = "encryption")] file_decryptor: None, } @@ -250,24 +496,14 @@ impl ParquetMetaData { &self.row_groups } - /// Returns the column index for this file if loaded - /// - /// Returns `None` if the parquet file does not have a `ColumnIndex` or - /// [ArrowReaderOptions::with_page_index] was set to false. - /// - /// [ArrowReaderOptions::with_page_index]: https://docs.rs/parquet/latest/parquet/arrow/arrow_reader/struct.ArrowReaderOptions.html#method.with_page_index - pub fn column_index(&self) -> Option<&ParquetColumnIndex> { - self.column_index.as_ref() - } - - /// Returns offset indexes in this file, if loaded + /// Returns the page index for this file if loaded /// - /// Returns `None` if the parquet file does not have a `OffsetIndex` or + /// Returns `None` if the parquet file lacks page indexes or /// [ArrowReaderOptions::with_page_index] was set to false. /// /// [ArrowReaderOptions::with_page_index]: https://docs.rs/parquet/latest/parquet/arrow/arrow_reader/struct.ArrowReaderOptions.html#method.with_page_index - pub fn offset_index(&self) -> Option<&ParquetOffsetIndex> { - self.offset_index.as_ref() + pub fn page_index(&self) -> Option<&PageIndex> { + self.page_index.as_ref() } /// Estimate of the bytes allocated to store `ParquetMetadata` @@ -293,19 +529,13 @@ impl ParquetMetaData { std::mem::size_of::() + self.file_metadata.heap_size() + self.row_groups.heap_size() - + self.column_index.heap_size() - + self.offset_index.heap_size() + + self.page_index.heap_size() + encryption_size } - /// Override the column index - pub(crate) fn set_column_index(&mut self, index: Option) { - self.column_index = index; - } - - /// Override the offset index - pub(crate) fn set_offset_index(&mut self, index: Option) { - self.offset_index = index; + /// Override the page index + pub(crate) fn set_page_index(&mut self, index: Option) { + self.page_index = index; } } @@ -386,35 +616,19 @@ impl ParquetMetaDataBuilder { } /// Sets the column index - pub fn set_column_index(mut self, column_index: Option) -> Self { - self.0.column_index = column_index; + pub fn set_page_index(mut self, page_index: Option) -> Self { + self.0.page_index = page_index; self } /// Returns the current column index from the builder, replacing it with `None` - pub fn take_column_index(&mut self) -> Option { - std::mem::take(&mut self.0.column_index) + pub fn take_page_index(&mut self) -> Option { + std::mem::take(&mut self.0.page_index) } /// Return a reference to the current column index, if any - pub fn column_index(&self) -> Option<&ParquetColumnIndex> { - self.0.column_index.as_ref() - } - - /// Sets the offset index - pub fn set_offset_index(mut self, offset_index: Option) -> Self { - self.0.offset_index = offset_index; - self - } - - /// Returns the current offset index from the builder, replacing it with `None` - pub fn take_offset_index(&mut self) -> Option { - std::mem::take(&mut self.0.offset_index) - } - - /// Return a reference to the current offset index, if any - pub fn offset_index(&self) -> Option<&ParquetOffsetIndex> { - self.0.offset_index.as_ref() + pub fn page_index(&self) -> Option<&PageIndex> { + self.0.page_index.as_ref() } /// Sets the file decryptor needed to decrypt this metadata. @@ -2120,12 +2334,16 @@ 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 page_index = PageIndex::new( + Some(vec![vec![Some(ColumnIndexMetaData::BOOLEAN(native_index))]]), + Some(vec![vec![offset_index]]), + ); let parquet_meta = ParquetMetaDataBuilder::new(file_metadata) .set_row_groups(row_group_meta) - .set_column_index(Some(vec![vec![ColumnIndexMetaData::BOOLEAN(native_index)]])) - .set_offset_index(Some(vec![vec![offset_index]])) + .set_page_index(Some(page_index)) .build(); #[cfg(not(feature = "encryption"))] diff --git a/parquet/src/file/metadata/parser.rs b/parquet/src/file/metadata/parser.rs index 9df6bcdd7185..b804d4f5948a 100644 --- a/parquet/src/file/metadata/parser.rs +++ b/parquet/src/file/metadata/parser.rs @@ -23,7 +23,7 @@ use crate::errors::ParquetError; use crate::file::metadata::thrift::parquet_metadata_from_bytes; use crate::file::metadata::{ - ColumnChunkMetaData, PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions, + ColumnChunkMetaData, PageIndex, PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions, }; use crate::file::page_index::column_index::ColumnIndexMetaData; @@ -233,23 +233,47 @@ pub(crate) fn decode_metadata( parquet_metadata_from_bytes(buf, options) } -/// Parses column index from the provided bytes and adds it to the metadata. +/// Parses page index from the provided bytes and adds it to the metadata. /// /// Arguments /// * `metadata` - The ParquetMetaData to which the parsed column index will be added. /// * `column_index_policy` - The policy for handling column index parsing (e.g., /// Required, Optional, Skip). +/// * `offset_index_policy` - The policy for handling offset index parsing (e.g., +/// Required, Optional, Skip). /// * `bytes` - The byte slice containing the column index data. /// * `start_offset` - The offset where `bytes` begin in the file. -pub(crate) fn parse_column_index( +pub(crate) fn parse_page_index( metadata: &mut ParquetMetaData, column_index_policy: PageIndexPolicy, + offset_index_policy: PageIndexPolicy, bytes: &Bytes, start_offset: u64, ) -> crate::errors::Result<()> { - if column_index_policy == PageIndexPolicy::Skip { + if column_index_policy == PageIndexPolicy::Skip && offset_index_policy == PageIndexPolicy::Skip + { return Ok(()); } + let column_indexes = parse_column_index(metadata, column_index_policy, bytes, start_offset)?; + let offset_indexes = parse_offset_index(metadata, offset_index_policy, bytes, start_offset)?; + // this likely shouldn't happen, but check just in case + if column_indexes.is_none() && offset_indexes.is_none() { + return Ok(()); + } + let page_index = PageIndex::new(column_indexes, offset_indexes); + metadata.set_page_index(Some(page_index)); + Ok(()) +} + +fn parse_column_index( + metadata: &ParquetMetaData, + column_index_policy: PageIndexPolicy, + bytes: &Bytes, + start_offset: u64, +) -> crate::errors::Result>>>> { + if column_index_policy == PageIndexPolicy::Skip { + return Ok(None); + } let index = metadata .row_groups() .iter() @@ -269,25 +293,25 @@ pub(crate) fn parse_column_index( rg_idx, col_idx, ) + .map(Some) } - None => Ok(ColumnIndexMetaData::NONE), + None => Ok(None), }) .collect::>>() }) .collect::>>()?; - metadata.set_column_index(Some(index)); - Ok(()) + Ok(Some(index)) } -pub(crate) fn parse_offset_index( - metadata: &mut ParquetMetaData, +fn parse_offset_index( + metadata: &ParquetMetaData, offset_index_policy: PageIndexPolicy, bytes: &Bytes, start_offset: u64, -) -> crate::errors::Result<()> { +) -> crate::errors::Result>>>> { if offset_index_policy == PageIndexPolicy::Skip { - return Ok(()); + return Ok(None); } let row_groups = metadata.row_groups(); let mut all_indexes = Vec::with_capacity(row_groups.len()); @@ -305,26 +329,20 @@ 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); } - metadata.set_offset_index(Some(all_indexes)); - Ok(()) + Ok(Some(all_indexes)) } diff --git a/parquet/src/file/metadata/push_decoder.rs b/parquet/src/file/metadata/push_decoder.rs index c9307340ccf9..c77aa2e34954 100644 --- a/parquet/src/file/metadata/push_decoder.rs +++ b/parquet/src/file/metadata/push_decoder.rs @@ -20,7 +20,7 @@ use crate::DecodeResult; use crate::encryption::decrypt::FileDecryptionProperties; use crate::errors::{ParquetError, Result}; use crate::file::FOOTER_SIZE; -use crate::file::metadata::parser::{MetadataParser, parse_column_index, parse_offset_index}; +use crate::file::metadata::parser::{MetadataParser, parse_page_index}; use crate::file::metadata::{FooterTail, PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions}; use crate::file::page_index::index_reader::acc_range; use crate::file::reader::ChunkReader; @@ -426,8 +426,13 @@ impl ParquetMetaDataPushDecoder { let buffer = self.get_bytes(&page_index_range)?; let offset = page_index_range.start; - parse_column_index(&mut metadata, self.column_index_policy, &buffer, offset)?; - parse_offset_index(&mut metadata, self.offset_index_policy, &buffer, offset)?; + parse_page_index( + &mut metadata, + self.column_index_policy, + self.offset_index_policy, + &buffer, + offset, + )?; self.state = DecodeState::Finished; return Ok(DecodeResult::Data(*metadata)); } @@ -503,6 +508,7 @@ pub fn range_for_page_index( mod tests { use super::*; use crate::arrow::ArrowWriter; + use crate::file::metadata::PageIndex; use crate::file::properties::WriterProperties; use arrow_array::{ArrayRef, Int64Array, RecordBatch, StringViewArray}; use bytes::Bytes; @@ -524,8 +530,7 @@ mod tests { assert_eq!(metadata.num_row_groups(), 2); assert_eq!(metadata.row_group(0).num_rows(), 200); assert_eq!(metadata.row_group(1).num_rows(), 200); - assert!(metadata.column_index().is_some()); - assert!(metadata.offset_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } /// It is possible to feed some, but not all, of the footer into the metadata decoder @@ -544,8 +549,7 @@ mod tests { assert_eq!(metadata.num_row_groups(), 2); assert_eq!(metadata.row_group(0).num_rows(), 200); assert_eq!(metadata.row_group(1).num_rows(), 200); - assert!(metadata.column_index().is_some()); - assert!(metadata.offset_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } /// It is possible to pre-fetch some, but not all, of the necessary data @@ -572,8 +576,7 @@ mod tests { assert_eq!(metadata.num_row_groups(), 2); assert_eq!(metadata.row_group(0).num_rows(), 200); assert_eq!(metadata.row_group(1).num_rows(), 200); - assert!(metadata.column_index().is_some()); - assert!(metadata.offset_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } #[test] @@ -619,8 +622,7 @@ mod tests { assert_eq!(metadata.num_row_groups(), 2); assert_eq!(metadata.row_group(0).num_rows(), 200); assert_eq!(metadata.row_group(1).num_rows(), 200); - assert!(metadata.column_index().is_some()); - assert!(metadata.offset_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } /// Decode the metadata incrementally, but without reading the page indexes @@ -647,8 +649,7 @@ mod tests { assert_eq!(metadata.num_row_groups(), 2); assert_eq!(metadata.row_group(0).num_rows(), 200); assert_eq!(metadata.row_group(1).num_rows(), 200); - assert!(metadata.column_index().is_none()); // of course, we did not read the column index - assert!(metadata.offset_index().is_none()); // or the offset index + assert!(metadata.page_index().is_none()); // of course, we did not read the page index } static TEST_BATCH: LazyLock = LazyLock::new(|| { diff --git a/parquet/src/file/metadata/reader.rs b/parquet/src/file/metadata/reader.rs index f700bba9e947..73c3be2237cf 100644 --- a/parquet/src/file/metadata/reader.rs +++ b/parquet/src/file/metadata/reader.rs @@ -64,8 +64,7 @@ use crate::arrow::async_reader::{MetadataFetch, MetadataSuffixFetch}; /// .with_page_index_policy(PageIndexPolicy::Required); /// reader.try_parse(&file).unwrap(); /// let metadata = reader.finish().unwrap(); -/// assert!(metadata.column_index().is_some()); -/// assert!(metadata.offset_index().is_some()); +/// assert!(metadata.page_index().is_some()); /// ``` #[derive(Default, Debug)] pub struct ParquetMetaDataReader { @@ -842,6 +841,7 @@ fn parse_index_data(push_decoder: &mut ParquetMetaDataPushDecoder) -> Result panic!("unexpected error"), } @@ -941,8 +937,7 @@ mod tests { } } let metadata = reader.finish().unwrap(); - assert!(metadata.column_index.is_some()); - assert!(metadata.offset_index.is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // not enough for page index but lie about file size let bytes = bytes_for_range(323584..len); @@ -1010,6 +1005,7 @@ mod async_tests { use tempfile::NamedTempFile; use crate::arrow::ArrowWriter; + use crate::file::metadata::PageIndex; use crate::file::properties::WriterProperties; use crate::file::reader::Length; use crate::util::test_common::file_util::get_test_file; @@ -1276,7 +1272,7 @@ mod async_tests { loader.try_load(f, len).await.unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 3); let metadata = loader.finish().unwrap(); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // Prefetch just footer exactly fetch_count.store(0, Ordering::SeqCst); @@ -1287,7 +1283,7 @@ mod async_tests { loader.try_load(f, len).await.unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 2); let metadata = loader.finish().unwrap(); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // Prefetch more than footer but not enough fetch_count.store(0, Ordering::SeqCst); @@ -1298,7 +1294,7 @@ mod async_tests { loader.try_load(f, len).await.unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 2); let metadata = loader.finish().unwrap(); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // Prefetch exactly enough fetch_count.store(0, Ordering::SeqCst); @@ -1310,7 +1306,7 @@ mod async_tests { .await .unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 1); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // Prefetch more than enough but less than the entire file fetch_count.store(0, Ordering::SeqCst); @@ -1322,7 +1318,7 @@ mod async_tests { .await .unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 1); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // Prefetch the entire file fetch_count.store(0, Ordering::SeqCst); @@ -1334,7 +1330,7 @@ mod async_tests { .await .unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 1); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); // Prefetch more than the entire file fetch_count.store(0, Ordering::SeqCst); @@ -1346,7 +1342,7 @@ mod async_tests { .await .unwrap(); assert_eq!(fetch_count.load(Ordering::SeqCst), 1); - assert!(metadata.offset_index().is_some() && metadata.column_index().is_some()); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); } fn write_parquet_file(offset_index_disabled: bool) -> Result { diff --git a/parquet/src/file/metadata/thrift/encryption.rs b/parquet/src/file/metadata/thrift/encryption.rs index 4e40acbf8186..b7f3f74c527f 100644 --- a/parquet/src/file/metadata/thrift/encryption.rs +++ b/parquet/src/file/metadata/thrift/encryption.rs @@ -298,8 +298,7 @@ pub(crate) fn parquet_metadata_with_encryption( let ParquetMetaData { mut file_metadata, row_groups, - column_index: _, - offset_index: _, + page_index: _, file_decryptor: _, } = parquet_meta; diff --git a/parquet/src/file/metadata/writer.rs b/parquet/src/file/metadata/writer.rs index 4b88077d5857..e8d0c77cd139 100644 --- a/parquet/src/file/metadata/writer.rs +++ b/parquet/src/file/metadata/writer.rs @@ -16,9 +16,7 @@ // under the License. use crate::file::metadata::thrift::FileMeta; -use crate::file::metadata::{ - ColumnChunkMetaData, ParquetColumnIndex, ParquetOffsetIndex, RowGroupMetaData, -}; +use crate::file::metadata::{ColumnChunkMetaData, PageIndex, RowGroupMetaData}; use crate::schema::types::{SchemaDescPtr, SchemaDescriptor}; use crate::{ basic::ColumnOrder, @@ -136,7 +134,7 @@ impl<'a, W: Write> ThriftMetadataWriter<'a, W> { } /// Serialize the column indexes and transform to `Option` - fn finalize_column_indexes(&mut self) -> Result> { + fn finalize_column_indexes(&mut self) -> Result>>>> { let column_indexes = std::mem::take(&mut self.column_indexes); // Write column indexes to file @@ -149,27 +147,15 @@ 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` - fn finalize_offset_indexes(&mut self) -> Result> { + fn finalize_offset_indexes(&mut self) -> Result>>>> { let offset_indexes = std::mem::take(&mut self.offset_indexes); // Write offset indexes to file @@ -182,18 +168,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 @@ -275,8 +254,7 @@ impl<'a, W: Write> ThriftMetadataWriter<'a, W> { // to be usable for retrieving the row group statistics for example, without users // needing to decrypt the metadata. let builder = ParquetMetaDataBuilder::new(file_metadata) - .set_column_index(column_indexes) - .set_offset_index(offset_indexes); + .set_page_index(Some(PageIndex::new(column_indexes, offset_indexes))); Ok(match unencrypted_row_groups { Some(rg) => builder.set_row_groups(rg).build(), @@ -464,8 +442,7 @@ 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 page_index = self.metadata.page_index().cloned(); let mut encoder = ThriftMetadataWriter::new( &mut self.buf, @@ -476,12 +453,18 @@ impl<'a, W: Write> ParquetMetaDataWriter<'a, W> { self.write_path_in_schema, ); - if let Some(column_indexes) = column_indexes { - encoder = encoder.with_column_indexes(column_indexes); - } + if let Some(PageIndex { + column_indexes, + offset_indexes, + }) = page_index + { + if let Some(column_indexes) = column_indexes { + encoder = encoder.with_column_indexes(column_indexes); + } - if let Some(offset_indexes) = offset_indexes { - encoder = encoder.with_offset_indexes(offset_indexes); + if let Some(offset_indexes) = offset_indexes { + encoder = encoder.with_offset_indexes(offset_indexes); + } } if let Some(key_value_metadata) = key_value_metadata { @@ -491,40 +474,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)] @@ -568,8 +517,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( @@ -580,14 +528,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 @@ -665,8 +607,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( @@ -677,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 b68e8e811d74..3c4f59d8d813 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,28 +543,19 @@ 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" - )), } }}; } /// Parsed [`ColumnIndex`] information for a Parquet file. /// -/// See [`ParquetColumnIndex`] for more information. +/// See [`PageIndex`] for more information. /// -/// [`ParquetColumnIndex`]: crate::file::metadata::ParquetColumnIndex +/// [`PageIndex`]: crate::file::metadata::PageIndex /// [`ColumnIndex`]: https://github.com/apache/parquet-format/blob/master/PageIndex.md #[derive(Debug, Clone, PartialEq)] #[expect(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), - ColumnIndexMetaData::NONE => Err(general_err!("Cannot serialize NONE index")), } } } diff --git a/parquet/src/file/page_index/offset_index.rs b/parquet/src/file/page_index/offset_index.rs index 06d21efb68c8..e4caac849189 100644 --- a/parquet/src/file/page_index/offset_index.rs +++ b/parquet/src/file/page_index/offset_index.rs @@ -48,9 +48,9 @@ thrift_struct!( /// [`OffsetIndex`] information for a column chunk. Contains offsets and sizes for each page /// in the chunk. Optionally stores fully decoded page sizes for BYTE_ARRAY columns. /// -/// See [`ParquetOffsetIndex`] for more information. +/// See [`PageIndex`] for more information. /// -/// [`ParquetOffsetIndex`]: crate::file::metadata::ParquetOffsetIndex +/// [`PageIndex`]: crate::file::metadata::PageIndex /// [`OffsetIndex`]: https://github.com/apache/parquet-format/blob/master/PageIndex.md pub struct OffsetIndexMetaData { /// Vector of [`PageLocation`] objects, one per page in the chunk. diff --git a/parquet/src/file/properties.rs b/parquet/src/file/properties.rs index c7a4f550b88f..b6369f1ace62 100644 --- a/parquet/src/file/properties.rs +++ b/parquet/src/file/properties.rs @@ -1131,7 +1131,7 @@ impl WriterPropertiesBuilder { /// /// Setting this value to `true` can greatly increase the size of the resulting Parquet /// file while yielding very little added benefit. Most modern Parquet implementations - /// will use the min/max values stored in the [`ParquetColumnIndex`] rather than + /// will use the min/max values stored in the [`PageIndex`] rather than /// those in the page header. /// /// # Note @@ -1142,7 +1142,7 @@ impl WriterPropertiesBuilder { /// specification. See [issue #7580] for more details. /// /// [`Statistics`]: crate::file::statistics::Statistics - /// [`ParquetColumnIndex`]: crate::file::metadata::ParquetColumnIndex + /// [`PageIndex`]: crate::file::metadata::PageIndex /// [`Page`]: EnabledStatistics::Page /// [issue #7580]: https://github.com/apache/arrow-rs/issues/7580 pub fn set_write_page_header_statistics(mut self, value: bool) -> Self { @@ -1417,10 +1417,10 @@ pub enum EnabledStatistics { /// Setting this option will store one set of statistics for each relevant /// column for each row group. In addition, this will enable the writing /// of the column index (the offset index is always written regardless of - /// this setting). See [`ParquetColumnIndex`] for + /// this setting). See [`PageIndex`] for /// more information. /// - /// [`ParquetColumnIndex`]: crate::file::metadata::ParquetColumnIndex + /// [`PageIndex`]: crate::file::metadata::PageIndex Page, } diff --git a/parquet/src/file/serialized_reader.rs b/parquet/src/file/serialized_reader.rs index f563fe2423a6..ed035d66ad89 100644 --- a/parquet/src/file/serialized_reader.rs +++ b/parquet/src/file/serialized_reader.rs @@ -312,7 +312,10 @@ impl FileReader for SerializedFileReader { Ok(Box::new(SerializedRowGroupReader::new( f, row_group_metadata, - self.metadata.offset_index().map(|x| x[i].as_slice()), + self.metadata + .page_index() + .map(|pi| pi.offset_indexes_for_rowgroup(i)) + .unwrap_or(None), props, )?)) } @@ -326,7 +329,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>, } @@ -336,7 +339,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() { @@ -371,7 +374,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( @@ -1748,17 +1755,22 @@ mod tests { row_group_metadata, file_reader .metadata - .offset_index() - .map(|x| x[row_group].as_slice()), + .page_index() + .map(|pi| pi.offset_indexes_for_rowgroup(row_group)) + .unwrap_or(None), props, )? }; 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( @@ -2059,8 +2071,7 @@ mod tests { let reader = SerializedFileReader::new_with_options(test_file, read_options)?; let metadata = reader.metadata(); assert_eq!(metadata.num_row_groups(), 1); - assert_eq!(metadata.column_index().unwrap().len(), 1); - assert_eq!(metadata.offset_index().unwrap().len(), 1); + assert!(metadata.page_index().is_some()); // true, false predicate let test_file = get_test_file("alltypes_tiny_pages.parquet"); @@ -2072,8 +2083,7 @@ mod tests { let reader = SerializedFileReader::new_with_options(test_file, read_options)?; let metadata = reader.metadata(); assert_eq!(metadata.num_row_groups(), 0); - assert!(metadata.column_index().is_none()); - assert!(metadata.offset_index().is_none()); + assert!(metadata.page_index().is_none()); // false, true predicate let test_file = get_test_file("alltypes_tiny_pages.parquet"); @@ -2085,8 +2095,7 @@ mod tests { let reader = SerializedFileReader::new_with_options(test_file, read_options)?; let metadata = reader.metadata(); assert_eq!(metadata.num_row_groups(), 0); - assert!(metadata.column_index().is_none()); - assert!(metadata.offset_index().is_none()); + assert!(metadata.page_index().is_none()); // false, false predicate let test_file = get_test_file("alltypes_tiny_pages.parquet"); @@ -2098,8 +2107,7 @@ mod tests { let reader = SerializedFileReader::new_with_options(test_file, read_options)?; let metadata = reader.metadata(); assert_eq!(metadata.num_row_groups(), 0); - assert!(metadata.column_index().is_none()); - assert!(metadata.offset_index().is_none()); + assert!(metadata.page_index().is_none()); Ok(()) } @@ -2145,11 +2153,10 @@ mod tests { let metadata = reader.metadata(); assert_eq!(metadata.num_row_groups(), 1); - let column_index = metadata.column_index().unwrap(); + let page_index = metadata.page_index().expect("page index should be present"); // only one row group - assert_eq!(column_index.len(), 1); - let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][0] else { + let Some(ColumnIndexMetaData::BYTE_ARRAY(index)) = page_index.column_index(0, 0) else { unreachable!() }; @@ -2163,11 +2170,14 @@ mod tests { assert_eq!(b"Hello", min.as_bytes()); assert_eq!(b"today", max.as_bytes()); - let offset_indexes = metadata.offset_index().unwrap(); // 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 offset_index = page_index + .offset_index(0, 0) + .expect("offset index should be present"); + let page_offset = offset_index + .page_locations() + .first() + .expect("offset index too small"); assert_eq!(4, page_offset.offset); assert_eq!(152, page_offset.compressed_page_size); @@ -2187,172 +2197,263 @@ mod tests { let metadata = reader.metadata(); assert_eq!(metadata.num_row_groups(), 1); - let column_index = metadata.column_index().unwrap(); - let row_group_offset_indexes = &metadata.offset_index().unwrap()[0]; + let page_index = metadata.page_index().unwrap(); + let row_group_offset_indexes = page_index.offset_indexes_for_rowgroup(0).unwrap(); // only one row group - assert_eq!(column_index.len(), 1); 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 = page_index.column_index(0, 0).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 = page_index.column_index(0, 1).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 = page_index.column_index(0, 2).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 = page_index.column_index(0, 3).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 = page_index.column_index(0, 4).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 = page_index.column_index(0, 5).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 = page_index.column_index(0, 6).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 = page_index.column_index(0, 7).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 = page_index.column_index(0, 8).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 = page_index.column_index(0, 9).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!(page_index.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 = page_index.column_index(0, 11).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 = page_index.column_index(0, 12).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!() } @@ -2615,16 +2716,10 @@ mod tests { let b = Bytes::from(out); let options = ReadOptionsBuilder::new().with_page_index().build(); let reader = SerializedFileReader::new_with_options(b, options).unwrap(); - let index = reader.metadata().column_index().unwrap(); + let page_index = reader.metadata().page_index().unwrap(); - // 1 row group - assert_eq!(index.len(), 1); - let c = &index[0]; - // 1 column - assert_eq!(c.len(), 1); - - match &c[0] { - ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v) => { + match page_index.column_index(0, 0) { + 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]); @@ -2747,19 +2842,16 @@ mod tests { assert_eq!(metadata.row_group(0).ordinal(), Some(2)); // check we only got the relevant page indexes - assert!(metadata.column_index().is_some()); - assert!(metadata.offset_index().is_some()); - assert_eq!(metadata.column_index().unwrap().len(), 1); - assert_eq!(metadata.offset_index().unwrap().len(), 1); - let col_idx = metadata.column_index().unwrap(); - let off_idx = metadata.offset_index().unwrap(); + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); + let page_index = metadata.page_index().unwrap(); + let col_stats = metadata.row_group(0).column(0).statistics().unwrap(); - let pg_idx = &col_idx[0][0]; - let off_idx_i = &off_idx[0][0]; + let pg_idx = page_index.column_index(0, 0); + let off_idx_i = page_index.offset_index(0, 0); // 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()); @@ -2770,7 +2862,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() ); @@ -2790,21 +2882,18 @@ mod tests { assert_eq!(metadata.row_group(1).ordinal(), Some(3)); // check we only got the relevant page indexes - assert!(metadata.column_index().is_some()); - assert!(metadata.offset_index().is_some()); - assert_eq!(metadata.column_index().unwrap().len(), 2); - assert_eq!(metadata.offset_index().unwrap().len(), 2); - let col_idx = metadata.column_index().unwrap(); - let off_idx = metadata.offset_index().unwrap(); - - for (i, col_idx_i) in col_idx.iter().enumerate().take(metadata.num_row_groups()) { - let col_stats = metadata.row_group(i).column(0).statistics().unwrap(); - let pg_idx = &col_idx_i[0]; - let off_idx_i = &off_idx[i][0]; + assert!(metadata.page_index().is_some_and(PageIndex::is_complete)); + + let page_index = metadata.page_index().unwrap(); + + for rg_idx in 0..metadata.num_row_groups() { + let col_stats = metadata.row_group(rg_idx).column(0).statistics().unwrap(); + let pg_idx = page_index.column_index(rg_idx, 0); + let off_idx_i = page_index.offset_index(rg_idx, 0); // 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()); @@ -2815,8 +2904,8 @@ mod tests { // check offset index matches too assert_eq!( - off_idx_i.page_locations[0].offset, - metadata.row_group(i).column(0).data_page_offset() + off_idx_i.as_ref().unwrap().page_locations[0].offset, + metadata.row_group(rg_idx).column(0).data_page_offset() ); } } diff --git a/parquet/src/file/writer.rs b/parquet/src/file/writer.rs index 1f66b5d4eb90..6bdc2f885f0b 100644 --- a/parquet/src/file/writer.rs +++ b/parquet/src/file/writer.rs @@ -2192,18 +2192,13 @@ mod tests { let options = ReadOptionsBuilder::new().with_page_index().build(); let reader = SerializedFileReader::new_with_options(Bytes::from(file), options).unwrap(); - let offset_index = reader.metadata().offset_index().unwrap(); - assert_eq!(offset_index.len(), 1); // 1 row group - assert_eq!(offset_index[0].len(), 2); // 2 columns - - let column_index = reader.metadata().column_index().unwrap(); - assert_eq!(column_index.len(), 1); // 1 row group - assert_eq!(column_index[0].len(), 2); // 2 column - - let a_idx = &column_index[0][0]; - assert!(matches!(a_idx, ColumnIndexMetaData::INT32(_)), "{a_idx:?}"); - let b_idx = &column_index[0][1]; - assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}"); + let a_idx = reader.metadata().page_index().unwrap().column_index(0, 0); + assert!( + matches!(a_idx, Some(ColumnIndexMetaData::INT32(_))), + "{a_idx:?}" + ); + let b_idx = reader.metadata().page_index().unwrap().column_index(0, 1); + assert!(b_idx.is_none(), "{b_idx:?}"); } #[test] @@ -2279,27 +2274,31 @@ mod tests { ); // check histogram in column index as well - assert!(reader.metadata().column_index().is_some()); - 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] { - assert_eq!(index.num_pages(), 1); - index - } else { - unreachable!() - }; + assert!(reader.metadata().page_index().is_some()); + let page_index = reader.metadata().page_index().unwrap(); + let col_idx = + if let Some(ColumnIndexMetaData::BYTE_ARRAY(index)) = page_index.column_index(0, 0) { + assert_eq!(index.num_pages(), 1); + index + } else { + unreachable!() + }; assert!(col_idx.repetition_level_histogram(0).is_none()); assert!(col_idx.definition_level_histogram(0).is_some()); check_def_hist(col_idx.definition_level_histogram(0).unwrap()); - assert!(reader.metadata().offset_index().is_some()); - 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()); - let page_sizes = offset_index[0][0] + assert!(page_index.offset_index(0, 0).is_some()); + assert!( + page_index + .offset_index(0, 0) + .unwrap() + .unencoded_byte_array_data_bytes + .is_some() + ); + let page_sizes = page_index + .offset_index(0, 0) + .unwrap() .unencoded_byte_array_data_bytes .as_ref() .unwrap(); @@ -2473,12 +2472,17 @@ mod tests { check_def_hist(column.definition_level_histogram().unwrap().values()); check_rep_hist(column.repetition_level_histogram().unwrap().values()); + assert!( + reader + .metadata() + .page_index() + .is_some_and(PageIndex::is_complete) + ); + let page_index = reader.metadata().page_index().unwrap(); + // check histogram in column index as well - assert!(reader.metadata().column_index().is_some()); - 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)) = page_index.column_index(0, 0) + { assert_eq!(index.num_pages(), 1); index } else { @@ -2488,11 +2492,13 @@ mod tests { check_def_hist(col_idx.definition_level_histogram(0).unwrap()); check_rep_hist(col_idx.repetition_level_histogram(0).unwrap()); - assert!(reader.metadata().offset_index().is_some()); - 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!( + page_index + .offset_index(0, 0) + .unwrap() + .unencoded_byte_array_data_bytes + .is_none() + ); } #[test] @@ -2654,16 +2660,24 @@ mod tests { let output = Vec::::new(); let mut writer = SerializedFileWriter::new(output, schema, props).unwrap(); - let column_indexes = metadata.column_index(); - let offset_indexes = metadata.offset_index(); + let page_index = metadata.page_index(); for (rg_idx, rg) in metadata.row_groups().iter().enumerate() { - let rg_column_indexes = column_indexes.and_then(|ci| ci.get(rg_idx)); - let rg_offset_indexes = offset_indexes.and_then(|oi| oi.get(rg_idx)); + let rg_column_indexes = + page_index.and_then(|pi| pi.column_indexes_for_rowgroup(rg_idx)); + let rg_offset_indexes = + page_index.and_then(|pi| pi.offset_indexes_for_rowgroup(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| { + let c = row.get(col_idx)?; + c.clone() + }); + let offset_index = rg_offset_indexes.and_then(|row| { + let o = row.get(col_idx)?; + o.clone() + }); + let result = ColumnCloseResult { bytes_written: column.compressed_size() as _, rows_written: rg.num_rows() as _, @@ -2706,8 +2720,14 @@ mod tests { let min = stats.min_bytes_opt().expect("min stats missing"); let max = stats.max_bytes_opt().expect("max stats missing"); - let col_idx = metadata.column_index().expect("column index not present"); - let ColumnIndexMetaData::INT96(col0) = &col_idx[0][0] else { + assert!( + metadata + .page_index() + .is_some_and(|pi| pi.has_column_indexes()) + ); + let Some(ColumnIndexMetaData::INT96(col0)) = + metadata.page_index().unwrap().column_index(0, 0) + else { panic!("expected INT96 stats") }; let col_min = col0.min_value(0).expect("ColumnIndex min not present"); diff --git a/parquet/tests/arrow_reader/io/async_reader.rs b/parquet/tests/arrow_reader/io/async_reader.rs index db06dda8ee89..1435ca549aad 100644 --- a/parquet/tests/arrow_reader/io/async_reader.rs +++ b/parquet/tests/arrow_reader/io/async_reader.rs @@ -328,8 +328,7 @@ async fn async_builder( .as_ref() .clone() .into_builder() - .set_column_index(None) - .set_offset_index(None) + .set_page_index(None) .build(); Arc::new(metadata) }; diff --git a/parquet/tests/arrow_reader/io/mod.rs b/parquet/tests/arrow_reader/io/mod.rs index cab3c24e7aa2..ef4ba1d53846 100644 --- a/parquet/tests/arrow_reader/io/mod.rs +++ b/parquet/tests/arrow_reader/io/mod.rs @@ -53,10 +53,10 @@ use parquet::arrow::async_reader::AsyncFileReader; use parquet::arrow::{ArrowWriter, ProjectionMask}; use parquet::data_type::AsBytes; use parquet::file::FOOTER_SIZE; -use parquet::file::metadata::PageIndexPolicy; #[cfg(feature = "async")] use parquet::file::metadata::ParquetMetaDataReader; -use parquet::file::metadata::{FooterTail, ParquetMetaData, ParquetOffsetIndex}; +use parquet::file::metadata::{FooterTail, ParquetMetaData}; +use parquet::file::metadata::{PageIndex, PageIndexPolicy}; use parquet::file::page_index::offset_index::PageLocation; use parquet::file::properties::WriterProperties; use parquet::schema::types::SchemaDescriptor; @@ -261,11 +261,11 @@ impl TestParquetFile { let parquet_metadata = Arc::clone(builder.metadata()); - let offset_index = parquet_metadata - .offset_index() + let page_index = parquet_metadata + .page_index() .expect("Parquet metadata should have a page index"); - let row_groups = TestRowGroups::new(&parquet_metadata, offset_index); + let row_groups = TestRowGroups::new(&parquet_metadata, page_index); // figure out the footer location in the file let footer_location = bytes.len() - FOOTER_SIZE..bytes.len(); @@ -342,7 +342,7 @@ struct TestRowGroups { } impl TestRowGroups { - fn new(parquet_metadata: &ParquetMetaData, offset_index: &ParquetOffsetIndex) -> Self { + fn new(parquet_metadata: &ParquetMetaData, page_index: &PageIndex) -> Self { let row_groups = parquet_metadata .row_groups() .iter() @@ -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 = page_index + .offset_index(rg_index, col_idx) + .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 2e2c0b6ea46f..f5bf2f364c6f 100644 --- a/parquet/tests/arrow_reader/row_filter/async.rs +++ b/parquet/tests/arrow_reader/row_filter/async.rs @@ -187,8 +187,12 @@ async fn test_cached_mask_reads_sparse_pages_without_error() { .unwrap(); let schema = builder.parquet_schema().clone(); let projection = ProjectionMask::leaves(&schema, [0]); - let page_first_rows = builder.metadata().offset_index().unwrap()[0][0] - .page_locations() + let page_first_rows = builder + .metadata() + .page_index() + .unwrap() + .page_locations(0, 0) + .unwrap() .iter() .map(|page| page.first_row_index) .collect::>(); @@ -379,10 +383,17 @@ async fn test_mask_nested_projection_with_different_page_boundaries() { ) .await .unwrap(); - let page_first_rows = builder.metadata().offset_index().unwrap()[0] + let page_first_rows = builder + .metadata() + .page_index() + .unwrap() + .offset_indexes_for_rowgroup(0) + .unwrap() .iter() .map(|column| { column + .as_ref() + .unwrap() .page_locations() .iter() .map(|page| page.first_row_index) diff --git a/parquet/tests/arrow_reader/statistics.rs b/parquet/tests/arrow_reader/statistics.rs index e0740d38e93e..11d6fac48e39 100644 --- a/parquet/tests/arrow_reader/statistics.rs +++ b/parquet/tests/arrow_reader/statistics.rs @@ -257,20 +257,15 @@ impl Test<'_> { let row_groups = reader.metadata().row_groups(); if check.data_page() { - let column_page_index = reader + let page_index = reader .metadata() - .column_index() - .expect("File should have column page indices"); - - let column_offset_index = reader - .metadata() - .offset_index() - .expect("File should have column offset indices"); + .page_index() + .expect("File should have page indices"); let row_group_indices: Vec<_> = (0..row_groups.len()).collect(); let min = converter - .data_page_mins(column_page_index, column_offset_index, &row_group_indices) + .data_page_mins(page_index, &row_group_indices) .unwrap(); assert_eq!( &min, &expected_min, @@ -278,7 +273,7 @@ impl Test<'_> { ); let max = converter - .data_page_maxes(column_page_index, column_offset_index, &row_group_indices) + .data_page_maxes(page_index, &row_group_indices) .unwrap(); assert_eq!( &max, &expected_max, @@ -286,7 +281,7 @@ impl Test<'_> { ); let null_counts = converter - .data_page_null_counts(column_page_index, column_offset_index, &row_group_indices) + .data_page_null_counts(page_index, &row_group_indices) .unwrap(); assert_eq!( @@ -296,7 +291,7 @@ impl Test<'_> { ); let row_counts = converter - .data_page_row_counts(column_offset_index, row_groups, &row_group_indices) + .data_page_row_counts(page_index, row_groups, &row_group_indices) .unwrap(); assert_eq!( row_counts, expected_row_counts, @@ -2946,12 +2941,9 @@ mod test { let parquet_schema = reader.parquet_schema(); let row_groups = metadata.row_groups(); let row_group_indices = [0]; - let column_page_index = metadata - .column_index() - .expect("file should have column page indices"); - let column_offset_index = metadata - .offset_index() - .expect("file should have column offset indices"); + let page_index = metadata + .page_index() + .expect("file should have page indices"); let DataType::Struct(fields) = schema.field_with_name("c1").unwrap().data_type() else { unreachable!("c1 must be a struct field") @@ -2984,34 +2976,22 @@ mod test { assert_eq!(leaf_row_counts, Some(UInt64Array::from(vec![6]))); let leaf_page_mins = leaf_converter - .data_page_mins( - column_page_index, - column_offset_index, - row_group_indices.iter(), - ) + .data_page_mins(page_index, row_group_indices.iter()) .unwrap(); assert_eq!(&leaf_page_mins, &i32_array([Some(1), Some(4)])); let leaf_page_maxes = leaf_converter - .data_page_maxes( - column_page_index, - column_offset_index, - row_group_indices.iter(), - ) + .data_page_maxes(page_index, row_group_indices.iter()) .unwrap(); assert_eq!(&leaf_page_maxes, &i32_array([Some(3), Some(9)])); let leaf_page_null_counts = leaf_converter - .data_page_null_counts( - column_page_index, - column_offset_index, - row_group_indices.iter(), - ) + .data_page_null_counts(page_index, row_group_indices.iter()) .unwrap(); assert_eq!(leaf_page_null_counts, UInt64Array::from(vec![1, 0])); let leaf_page_row_counts = leaf_converter - .data_page_row_counts(column_offset_index, row_groups, row_group_indices.iter()) + .data_page_row_counts(page_index, row_groups, row_group_indices.iter()) .unwrap(); assert_eq!(leaf_page_row_counts, Some(UInt64Array::from(vec![3, 3]))); @@ -3031,11 +3011,7 @@ mod test { ); let amount_page_mins = amount_converter - .data_page_mins( - column_page_index, - column_offset_index, - row_group_indices.iter(), - ) + .data_page_mins(page_index, row_group_indices.iter()) .unwrap(); assert_eq!( &amount_page_mins, @@ -3043,11 +3019,7 @@ mod test { ); let amount_page_maxes = amount_converter - .data_page_maxes( - column_page_index, - column_offset_index, - row_group_indices.iter(), - ) + .data_page_maxes(page_index, row_group_indices.iter()) .unwrap(); assert_eq!( &amount_page_maxes, diff --git a/parquet/tests/arrow_writer/layout.rs b/parquet/tests/arrow_writer/layout.rs index 1c63a3144391..55a489272f88 100644 --- a/parquet/tests/arrow_writer/layout.rs +++ b/parquet/tests/arrow_writer/layout.rs @@ -82,23 +82,28 @@ fn do_test(test: LayoutTest) { fn assert_layout(file_reader: &Bytes, meta: &ParquetMetaData, layout: &Layout) { assert_eq!(meta.row_groups().len(), layout.row_groups.len()); - let iter = meta - .row_groups() - .iter() - .zip(&layout.row_groups) - .zip(meta.offset_index().unwrap()); - for ((row_group, row_group_layout), offset_index) in iter { + for rg_idx in 0..meta.num_row_groups() { + let row_group = meta.row_group(rg_idx); + let row_group_layout = &layout.row_groups[rg_idx]; + let offset_index = meta + .page_index() + .unwrap() + .offset_indexes_for_rowgroup(rg_idx) + .unwrap(); + // Check against offset index assert_eq!(offset_index.len(), row_group_layout.columns.len()); 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 +115,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) @@ -597,9 +604,9 @@ fn test_per_column_data_page_size_limit() { assert_eq!(row_group.columns().len(), 2); // 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 page_index = metadata.page_index().unwrap(); + let col_a_page_count = page_index.offset_index(0, 0).unwrap().page_locations.len(); + let col_b_page_count = page_index.offset_index(0, 1).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 9031813134a5..d3e29ef003f5 100644 --- a/parquet/tests/encryption/encryption_util.rs +++ b/parquet/tests/encryption/encryption_util.rs @@ -181,23 +181,23 @@ pub(crate) fn verify_encryption_test_data( /// Verifies that the column and offset indexes were successfully read from an /// encrypted test file. pub(crate) fn verify_column_indexes(metadata: &ParquetMetaData) { - let offset_index = metadata.offset_index().unwrap(); + assert!(metadata.page_index().is_some()); + let page_index = metadata.page_index().unwrap(); + let offset_index = page_index.offset_indexes_for_rowgroup(0).unwrap(); // 1 row group, 8 columns - assert_eq!(offset_index.len(), 1); - assert_eq!(offset_index[0].len(), 8); + assert_eq!(offset_index.len(), 8); // 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); + let offset_index = &offset_index[float_col_idx]; + 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); - assert_eq!(column_index[0].len(), 8); - let column_index = &column_index[0][float_col_idx]; + let column_index = page_index.column_indexes_for_rowgroup(0).unwrap(); + assert_eq!(column_index.len(), 8); + let column_index = &column_index[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!( diff --git a/parquet/tests/ieee754_nan_interop.rs b/parquet/tests/ieee754_nan_interop.rs index cc3326dae555..ec8e6548f8f0 100644 --- a/parquet/tests/ieee754_nan_interop.rs +++ b/parquet/tests/ieee754_nan_interop.rs @@ -61,11 +61,7 @@ fn validate_float_metadata( assert_eq!(&mins, &exp); // verify page mins (should be 1 page per row group, so should be same) - let page_mins = converter.data_page_mins( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_mins = converter.data_page_mins(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_mins, &exp); let exp: Arc = Arc::new(Float32Array::from(FLOAT_MAXS.to_vec())); @@ -73,22 +69,16 @@ fn validate_float_metadata( assert_eq!(&maxs, &exp); // verify page maxs (should be 1 page per row group, so should be same) - let page_maxs = converter.data_page_maxes( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_maxs = + converter.data_page_maxes(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_maxs, &exp); let exp = UInt64Array::from(NAN_COUNTS.to_vec()); let nans = converter.row_group_nan_counts(metadata.row_groups())?; assert_eq!(&nans, &exp); - let page_nans = converter.data_page_nan_counts( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_nans = + converter.data_page_nan_counts(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_nans, &exp); Ok(()) @@ -116,11 +106,7 @@ fn validate_double_metadata( assert_eq!(&mins, &exp); // verify page mins (should be 1 page per row group, so should be same) - let page_mins = converter.data_page_mins( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_mins = converter.data_page_mins(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_mins, &exp); let exp: Arc = Arc::new(Float64Array::from(DOUBLE_MAXS.to_vec())); @@ -128,22 +114,16 @@ fn validate_double_metadata( assert_eq!(&maxs, &exp); // verify page maxs (should be 1 page per row group, so should be same) - let page_maxs = converter.data_page_maxes( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_maxs = + converter.data_page_maxes(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_maxs, &exp); let exp = UInt64Array::from(NAN_COUNTS.to_vec()); let nans = converter.row_group_nan_counts(metadata.row_groups())?; assert_eq!(&nans, &exp); - let page_nans = converter.data_page_nan_counts( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_nans = + converter.data_page_nan_counts(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_nans, &exp); Ok(()) @@ -183,11 +163,7 @@ fn validate_float16_metadata( assert_eq!(&mins, &exp); // verify page mins (should be 1 page per row group, so should be same) - let page_mins = converter.data_page_mins( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_mins = converter.data_page_mins(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_mins, &exp); let exp: Arc = Arc::new(Float16Array::from(FLOAT16_MAXS.to_vec())); @@ -195,22 +171,16 @@ fn validate_float16_metadata( assert_eq!(&maxs, &exp); // verify page maxs (should be 1 page per row group, so should be same) - let page_maxs = converter.data_page_maxes( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_maxs = + converter.data_page_maxes(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_maxs, &exp); let exp = UInt64Array::from(NAN_COUNTS.to_vec()); let nans = converter.row_group_nan_counts(metadata.row_groups())?; assert_eq!(&nans, &exp); - let page_nans = converter.data_page_nan_counts( - metadata.column_index().unwrap(), - metadata.offset_index().unwrap(), - &row_group_indices, - )?; + let page_nans = + converter.data_page_nan_counts(metadata.page_index().unwrap(), &row_group_indices)?; assert_eq!(&page_nans, &exp); Ok(())