sunchao commented on code in PR #11157:
URL: https://github.com/apache/arrow-rs/pull/11157#discussion_r4138198174
##########
parquet/src/file/metadata/parser.rs:
##########
@@ -236,84 +237,84 @@ pub(crate) fn decode_metadata(
parquet_metadata_from_bytes(buf, options)
}
-/// Parses page index from the provided bytes and adds it to the metadata.
+/// Parses page indexes from the provided bytes and replaces those in the
metadata.
///
/// Arguments
-/// * `metadata` - The ParquetMetaData to which the parsed column index will
be added.
+/// * `metadata` - The ParquetMetaData whose page index will be replaced.
/// * `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 page index data.
-/// * `start_offset` - The offset where `bytes` begin in the file.
+/// * `column_index_mask` - The row groups and leaf columns whose column
indexes are parsed.
+/// Ignored when `column_index_policy` is [`PageIndexPolicy::Skip`].
+/// * `offset_index_mask` - The row groups and leaf columns whose offset
indexes are parsed.
+/// Ignored when `offset_index_policy` is [`PageIndexPolicy::Skip`].
+/// * `bytes` - [`PushBuffers`] that should have already been populated with
the bytes containing
+/// the page indexes.
pub(crate) fn parse_page_index(
metadata: &mut ParquetMetaData,
column_index_policy: PageIndexPolicy,
offset_index_policy: PageIndexPolicy,
- bytes: &Bytes,
- start_offset: u64,
+ column_index_mask: &ColumnChunkMask,
+ offset_index_mask: &ColumnChunkMask,
+ bytes: &PushBuffers,
) -> crate::errors::Result<()> {
- if column_index_policy == PageIndexPolicy::Skip && offset_index_policy ==
PageIndexPolicy::Skip
- {
- return Ok(());
- }
let num_row_groups = metadata.num_row_groups();
let num_columns = metadata.file_metadata().schema_descr().num_columns();
let mut builder = PageIndexBuilder::default();
+
if column_index_policy != PageIndexPolicy::Skip {
builder.allocate_column_indexes(num_row_groups, num_columns);
parse_column_index(
metadata,
column_index_policy,
+ column_index_mask,
&mut builder,
bytes,
- start_offset,
)?;
}
if offset_index_policy != PageIndexPolicy::Skip {
builder.allocate_offset_indexes(num_row_groups, num_columns);
parse_offset_index(
metadata,
offset_index_policy,
+ offset_index_mask,
&mut builder,
bytes,
- start_offset,
)?;
}
+ // Always replace the page index, even if both policies are Skip.
+ // This ensures that repeated calls to read_page_indexes replace the
existing index
+ // rather than preserving it.
let page_index = builder.build();
- // if both indexes are missing from the file, return without modifying
`metadata`
- if !page_index.has_column_indexes() && !page_index.has_offset_indexes() {
- return Ok(());
+ if page_index.has_column_indexes() || page_index.has_offset_indexes() {
+ metadata.set_page_index(Some(Arc::new(page_index)));
+ } else {
+ // If no indexes were read, clear the page index
+ metadata.set_page_index(None);
}
- metadata.set_page_index(Some(Arc::new(page_index)));
Ok(())
}
fn parse_column_index(
metadata: &ParquetMetaData,
column_index_policy: PageIndexPolicy,
+ mask: &ColumnChunkMask,
page_index_builder: &mut PageIndexBuilder,
- bytes: &Bytes,
- start_offset: u64,
+ bytes: &PushBuffers,
) -> crate::errors::Result<()> {
if column_index_policy == PageIndexPolicy::Skip {
return Ok(());
}
- for rg_idx in 0..metadata.num_row_groups() {
+ for rg_idx in mask.row_group_indices(metadata.num_row_groups()) {
let rg = metadata.row_group(rg_idx);
- for col_idx in 0..rg.num_columns() {
+ for col_idx in mask.column_indices(rg.num_columns()) {
let col = rg.column(col_idx);
if let Some(r) = col.column_index_range() {
- let r_start = usize::try_from(r.start - start_offset)?;
- let r_end = usize::try_from(r.end - start_offset)?;
- let idx = inner::parse_single_column_index(
- &bytes[r_start..r_end],
- metadata,
- col,
- rg_idx,
- col_idx,
- )?;
+ let idx_bytes = bytes.get_bytes(r.start, (r.end - r.start) as
usize)?;
Review Comment:
Non-blocking simplification: now that the decoder again requires one buffer
covering the entire selected index range, could this path resolve that buffer
once and borrow subslices, as before? Each get_bytes here (and in
parse_offset_index) scans PushBuffers and creates/drops a Bytes slice. Reusing
the covering buffer would avoid repeated searches and reference-count
operations while keeping the same I/O behavior. I have not measured a material
regression, so this can be a follow-up.
##########
parquet/src/file/metadata/reader.rs:
##########
@@ -108,6 +115,166 @@ impl From<bool> for PageIndexPolicy {
}
}
+/// Struct to specify column chunks for which metadata is required.
+///
+/// Column chunks are identified by row group index and leaf column index (the
index of the
+/// column in [`SchemaDescriptor::columns`], not the index of a root or Arrow
field). This struct
+/// allows for specifying vertical slices of column chunk data (via
[`Self::columns`]),
+/// horizontal slices (via [`Self::row_groups`]), or the intersection of the
two
+/// (via [`Self::row_groups_and_columns`]).
+///
+/// At present this is only used to select elements of the [Page Index] for
decoding.
+///
+/// # Examples
+///
+/// To select columns 0 and 1 from all row groups:
+/// ```rust
+/// # use parquet::file::metadata::ColumnChunkMask;
+/// let mask = ColumnChunkMask::columns([0, 1]);
+/// ```
+///
+/// To select all columns from row group 2:
+/// ```rust
+/// # use parquet::file::metadata::ColumnChunkMask;
+/// let mask = ColumnChunkMask::row_groups([2]);
+/// ```
+///
+/// To select columns 1 and 3 from row group 0:
+/// ```rust
+/// # use parquet::file::metadata::ColumnChunkMask;
+/// let mask = ColumnChunkMask::row_groups_and_columns([0], [1, 3]);
+/// ```
+///
+/// [Page Index]: https://parquet.apache.org/docs/file-format/pageindex/
+#[derive(Debug, Clone, PartialEq, Eq, Hash, Default)]
+pub struct ColumnChunkMask {
+ // `None` means all, while `Some(empty)` means none. Store u32 because
+ // Parquet/Thrift collections cannot contain more than i32::MAX entries.
+ row_groups: Option<Arc<[u32]>>,
+ columns: Option<Arc<[u32]>>,
+}
+
+impl ColumnChunkMask {
+ /// Select all row groups and columns.
+ pub fn all() -> Self {
+ Self::default()
+ }
+
+ /// Select no row groups or columns.
+ pub fn none() -> Self {
+ Self {
+ row_groups: Some(Arc::from([])),
+ columns: Some(Arc::from([])),
+ }
+ }
+
+ /// Select only the listed columns.
+ ///
+ /// Passing an empty iterator selects no columns.
+ pub fn columns(columns: impl IntoIterator<Item = usize>) -> Self {
+ Self {
+ row_groups: None,
+ columns: Self::iter_to_set(columns),
+ }
+ }
+
+ /// Select only the listed row groups.
+ ///
+ /// Passing an empty iterator selects no row groups.
+ pub fn row_groups(row_groups: impl IntoIterator<Item = usize>) -> Self {
+ Self {
+ row_groups: Self::iter_to_set(row_groups),
+ columns: None,
+ }
+ }
+
+ /// Select only the listed row groups and columns.
+ ///
+ /// An empty iterator for either dimension selects no column chunks.
+ pub fn row_groups_and_columns(
+ row_groups: impl IntoIterator<Item = usize>,
+ columns: impl IntoIterator<Item = usize>,
+ ) -> Self {
+ Self {
+ row_groups: Self::iter_to_set(row_groups),
+ columns: Self::iter_to_set(columns),
+ }
+ }
+
+ /// Test if `idx` is in the row group set.
+ pub fn includes_row_group(&self, idx: usize) -> bool {
+ Self::includes_index(self.row_groups.as_ref(), idx)
+ }
+
+ /// Test if `idx` is in the column set.
+ pub fn includes_column(&self, idx: usize) -> bool {
+ Self::includes_index(self.columns.as_ref(), idx)
+ }
+
+ fn includes_index(keep: Option<&Arc<[u32]>>, idx: usize) -> bool {
+ // return false for out-of-bounds index
+ let Ok(idx) = u32::try_from(idx) else {
+ return false;
+ };
+ keep.is_none_or(|keep| keep.binary_search(&idx).is_ok())
+ }
+
+ /// Returns `true` when this mask selects every column chunk.
+ pub fn is_all(&self) -> bool {
+ self.row_groups.is_none() && self.columns.is_none()
+ }
+
+ /// Returns selected row groups, or `None` when all row groups are
selected.
+ pub fn selected_row_groups(&self) -> Option<&[u32]> {
+ self.row_groups.as_deref()
+ }
+
+ /// Returns selected leaf columns, or `None` when all columns are selected.
+ pub fn selected_columns(&self) -> Option<&[u32]> {
+ self.columns.as_deref()
+ }
+
+ /// Creates a mask selecting the leaf columns in an Arrow projection.
+ #[cfg(feature = "arrow")]
+ pub fn from_projection(projection: &ProjectionMask, schema:
&SchemaDescriptor) -> Self {
+ Self::columns((0..schema.num_columns()).filter(|&i|
projection.leaf_included(i)))
+ }
+
+ /// Returns an iterator over the row group indices selected by this mask
+ pub fn row_group_indices(&self, num_row_groups: usize) -> Box<dyn
Iterator<Item = usize> + '_> {
+ Self::axis_indices(self.row_groups.as_deref(), num_row_groups)
+ }
+
+ /// Returns an iterator over the column indices selected by this mask
+ pub fn column_indices(&self, num_columns: usize) -> Box<dyn Iterator<Item
= usize> + '_> {
+ Self::axis_indices(self.columns.as_deref(), num_columns)
+ }
+
+ fn axis_indices(axis: Option<&[u32]>, len: usize) -> Box<dyn Iterator<Item
= usize> + '_> {
Review Comment:
Optional performance follow-up: these boxed iterators add an allocation per
selected row group when iterating columns, in both range discovery and parsing,
plus dynamic dispatch during iteration. A small concrete iterator covering the
all-indices and selected-indices cases could avoid that overhead. It would be
useful to measure a many-row-group sparse-mask case before adding
implementation complexity; this does not block approval.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]