Conversation
| /// | ||
| /// [Page Index]: https://parquet.apache.org/docs/file-format/pageindex/ | ||
| #[derive(Debug, Clone, PartialEq, Eq, Default)] | ||
| pub struct ColumnChunkMask { |
There was a problem hiding this comment.
A note on naming: this was originally PageIndexSelection, but it occurred to me that we might want to use this for the file metadata as well. I changed it to RowGroupColumnSelection, but that was a bit unwieldy. I asked Google for some naming suggestions, and this name made the most sense to me since the intersection of row group and column is a column chunk.
| pub(crate) column_index: PageIndexPolicy, | ||
| pub(crate) offset_index: PageIndexPolicy, |
There was a problem hiding this comment.
Drive by change. Accessors have been added for these, so they no longer need to be visible.
|
run benchmark metadata env:
BENCH_FILTER: page index |
|
🤖 Arrow criterion benchmark running (GKE) | trigger CPU Details (lscpu)Comparing page_index_selection (69ab61a) to e3a21bf (merge-base) diff Run configurationrun benchmark metadata
env:
BENCH_FILTER: "page index"BENCH_COMMAND=cargo bench --features=arrow,async,test_common,experimental,object_store --bench metadata File an issue against this benchmark runner |
|
🤖 Arrow criterion benchmark completed (GKE) | trigger Instance: Comparing page_index_selection (69ab61a) to e3a21bf (merge-base) diff Run configurationrun benchmark metadata
env:
BENCH_FILTER: "page index"CPU Details (lscpu)Details
Resource Usagebase (merge-base)
branch
File an issue against this benchmark runner |
There was a problem hiding this comment.
@etseidl similar disclaimer here: this review text is largely AI generated, although I worked with my agent extensively to guide the style of review, things to look for and how to present them as well as discussing some of the individual findings (but not all). If this is not helpful please let me know and I will change strategies.
Masks work for the sync reader and for ParquetMetaDataReader::parse_and_finish. I found one bug that silently returns wrong page index data (C1), one reader error that masks make common (C2), and misaligned page statistics (C3). Some API choices are hard to change after a release (C4, C5, C13).
Scope: the GitHub diff also contains older copies of #11031 and #11132. I review those in their own PRs. This review covers commits 3a8ee6a to b682ebe. All outputs below come from scratch tests on this head (not committed).
| id | severity | summary | where |
|---|---|---|---|
| C1 | Bug | Suffix load + prefetch hint + mask returns the page index of another column chunk, with no error. Without a mask, the same call fails. | body, inline reader.rs L622 |
| C2 | Bug (on main) |
Async reader and push decoder fail with column was not fetched when the offset index is partial. Masks make this the normal path. |
body |
| C3 | Bug (on main) |
StatisticsConverter::data_page_* arrays have different lengths with a partial index. The options docs suggest the mask setup that causes it. |
body, inline arrow_reader/mod.rs L779 |
| C4 | API | An empty set selects everything. columns([-1]) selects everything, columns([100]) selects nothing. |
inline reader.rs L212 |
| C5 | API | i32 in public constructors. The neighbour APIs use usize. |
inline reader.rs L157 |
| C6 | Docs | The column indexes are leaf indexes. The docs do not say so. | inline reader.rs L112 |
| C13 | API | No none(), is_all(), union, getters, Hash, or ArrowReaderOptions::with_page_index_mask. |
body |
| C7 | Footgun | Masks are silently ignored in 4 paths. | inline arrow_reader/mod.rs L1012 |
| C8 | Perf | ParquetMetaDataReader fetches one covering range, so a column mask saves almost no I/O. |
inline reader.rs L251 |
| C9 | Perf | Push decoder cost is O(ranges x buffers): about 1 s for 30,000 ranges. | inline push_decoder.rs L451 |
| C17 | Perf | Ranges merge only when adjacent in emission order. Duplicates and overlaps stay separate. | inline push_decoder.rs L560 |
| C15 | Perf (nit) | The parse loops visit every row group x column. | inline parser.rs L306 |
| C10, C18 | API | Push decoder clients can loop forever in two ways. | body |
| C11 | Footgun | A second read_page_indexes replaces the old index, unless the new one is empty. |
body |
| C12 | API | Required checks only the selected chunks and depends on the column index mask. |
body, docs in C8 suggestion |
| C14 | Test gap | No test for the ArrowReaderOptions masks, ParquetMetaDataReader::with_page_index_mask, or a data read with a mask. |
inline page_index.rs L236 |
| N1 | Test gap | test_chunk_mask runs only with the async feature. |
inline reader.rs L1530 |
| C16 | Docs | Typo, stale start_offset doc. |
inline |
C1: suffix load + prefetch hint + mask returns another chunk's index, with no error
File: 4 Int32 columns c0..c3, 3 row groups of 10 rows, c1 = 1..=10, c2 = 2..=11, 3183 bytes. All calls use ParquetMetaDataReader::with_prefetch_hint(..).load_via_suffix_and_finish(..).
| mask | prefetch hint | result |
|---|---|---|
none, OI Required |
3172 | Err("Corrupted parquet file: index data range (1372..1503) exceeds remainder length (1492)") |
CI row_groups_and_columns([0], [1]) |
none | c1 page min / max = [1] / [10] (correct) |
CI row_groups_and_columns([0], [1]) |
3156 | Ok, c1 page min / max = [2] / [11], which are the stats of c2 |
OI row_groups_and_columns([0], [1]), Required |
3172 | Ok, c1 page offset = 232, which is the data_page_offset of c2 (correct: 145) |
A page pruner that evaluates c1 = 1 with these stats skips the page that holds the row. ParquetObjectReader::with_footer_size_hint + ArrowReaderOptions::with_offset_index_mask gives the same wrong offset index. An Arrow read of c1 (no RowSelection) with that metadata then panics: range end out of bounds: 141 <= 87 (in_memory_row_group.rs L303). A custom AsyncFileReader that follows the documented get_metadata example and adds a prefetch hint is affected too.
Cause (older than this PR, reader.rs L826 and L637-648):
flowchart LR
A["fetch_suffix(S)<br/>suffix = file[len-S..len]"] --> B["remainder = (0, suffix[..metadata_start])<br/>real start is len-S, not 0"]
B --> C["load_page_index_with_remainder<br/>slices at range.start - 0"]
C --> D["reads file[range.start + len-S ..]<br/>= index of a later chunk"]
Without a mask, the covering range ends at the footer. So the wrong slice is out of bounds, and the call fails. A mask makes the range shorter, so the wrong slice is in bounds.
Fix: the file offset of the suffix is not known, so do not reuse it as a remainder.
Ok((
self.decode_footer_metadata(slice, file_size, footer)?,
- Some((0, suffix.slice(..metadata_start))),
+ // the file offset of `suffix` is not known, so it cannot hold the page index
+ None,
))With this change, all repros above return the correct index. cargo test -p parquet --features arrow,async,object_store --lib passes for file::metadata and arrow::async_reader. The change is one line, so it can go in this PR. Separately, ColumnChunkData::get_bytes (in_memory_row_group.rs L303) should return an error, not panic, for a page outside the chunk.
C2: async reader and push decoder fail with a partial offset index (bug on main)
File: 2 Int32 columns, 2 row groups of 500 rows, 100 rows per page. Options: with_page_index_policy(Optional).with_offset_index_mask(..). Err(n) = Invalid column index n, column was not fetched.
| reader | all() |
columns([0]) |
columns([1]) |
row_groups([1]) |
|---|---|---|---|---|
sync, RowSelection (40 rows) |
40 | 40 | 40 | 40 |
push, RowSelection |
40 | Err(1) | Err(1) | Err(0) |
async, RowSelection |
40 | Err(1) | Err(1) | Err(0) |
async, RowFilter only (500 rows) |
500 | Err(1) | 500 | Err(1) |
The row_groups([1]) column is the "load the offset index only for row groups that survive pruning" use case. The columns([0]) column is the use case of the with_offset_index_mask docs. The cause and a tested fix are in my review of #11132: InMemoryRowGroup aligns page_start_offsets to the columns by position. The bug is on main since #10719 (5ec9eaf).
This PR does not cause the bug. I recommend (not require) that the fix lands on main first, in a separate PR: #11182. Its regression tests then also cover the new mask paths of this PR. After that, add one push decoder or async test with a mask here.
C3: StatisticsConverter page arrays lose alignment with a partial index (bug on main)
Columns a, b; 3 row groups of 10 pages; converter for a; row groups [0, 1, 2]:
| CI mask | OI mask | data_page_mins len |
data_page_row_counts len |
|---|---|---|---|
all() |
all() |
30 | 30 |
columns([0]) (predicate) |
columns([1]) (projection) |
30 | 0 |
columns([0]) |
row_groups_and_columns([1], [0]) |
30 | 10 |
row_groups_and_columns([1, 2], [0]) |
same | 20 | 20 |
has_offset_indexes() is true in all rows, so a caller cannot detect the problem. Row 2 is the setup that the ArrowReaderOptions docs suggest. Cause: data_page_row_counts skips a row group with no offset index (statistics.rs L2044-2046), and data_page_* use num_data_pages(..).unwrap_or(0) (L1911-1914 and 3 copies). Fix: return nulls (or None) for such row groups, not fewer entries. The inline doc change on with_offset_index_mask covers the user side.
API shape (C4, C5, C6, C13)
Two choices are one-way doors. If "empty = all" ships, a later fix silently changes what every caller loads. If i32 ships, a later change to usize breaks every caller. Everything else below is additive. A proposal that keeps the current names:
#[derive(Debug, Clone, PartialEq, Eq, Hash, Default)] // Default = all()
pub struct ColumnChunkMask { /* per axis: Option<Arc<[u32]>> (sorted), None = all */ }
impl ColumnChunkMask {
pub fn all() -> Self;
pub fn none() -> Self;
/// Leaf (Parquet) column indexes, as in `SchemaDescriptor::column(i)`.
/// An empty iterator selects no columns, as in `ProjectionMask::leaves`.
pub fn columns(leaves: impl IntoIterator<Item = usize>) -> Self;
/// An empty iterator selects no row groups, as in `with_row_groups(vec![])`.
pub fn row_groups(row_groups: impl IntoIterator<Item = usize>) -> Self;
pub fn row_groups_and_columns(/* same rules */) -> Self;
pub fn from_projection(projection: &ProjectionMask, schema: &SchemaDescriptor) -> Self;
pub fn union(&self, other: &Self) -> Self;
pub fn is_all(&self) -> bool;
pub fn selected_columns(&self) -> Option<&[u32]>; // None = all
pub fn selected_row_groups(&self) -> Option<&[u32]>;
// includes_row_group / includes_column as today
}
impl ArrowReaderOptions {
pub fn with_page_index_mask(self, mask: ColumnChunkMask) -> Self;
}Hash and the getters let a metadata cache use the mask in its key and see what it loaded.
Other items
- C10 (API): the push decoder waits for one pushed buffer per merged range. A client that pushes the per-chunk ranges (for example from
column_index_range()) pushes every byte, but the decoder asks again:asks again for [526..580, 607..661, 688..709, 720..742], forever. Fix: gate on the per-chunk ranges that the parser reads, and merge only for theNeedsDataoutput. - C18 (API): a client that calls
clear_all_ranges()and then pushes exactly the requested ranges each round loops forever after a partial prefetch: rounds alternate between[11140..11191, 11242..11344, 11395..11468]and[11490..11534, 11558..11582](not finished after 101 rounds). Onmainthe same client finishes in 1 round, because there is one covering range. Document thatNeedsDatais relative to the data that is buffered now. - C11 (Footgun): a second
read_page_indexes(ortry_new_with_metadataon metadata with a page index) replaces the old index. If the new index is empty, it keeps the old one (parser.rsL288-291). Output: first read mask[0]gives cells[(0, 0)]; second read mask[1]gives[(0, 1)]; second read mask[7]gives[(0, 0)].is_complete()istruefor the partial index. Choose "replace" or "merge", and document it onread_page_indexes. - C12 (API):
Required+columns([42])returnsOkwith no page index. On a file with no offset index,OI Requiredfails with CI maskall(), but returnsOk(no page index)with CI maskcolumns([5])or CISkip. The check inparse_offset_indexruns only if some range was requested. At least document it (see the C8 suggestion). - P4:
ParquetMetaDataWriteron a masked page index writes stale index offsets. The output then fails to read back:EOF("Parquet file too small. Page index range 0..253 overlaps with file metadata 37..525"). Details are in my review of #11031. - C16 (Nit): the comment at
reader.rsL488 still namesrange_for_page_index().
| /// Sets the [`ColumnChunkMask`] for the Parquet [ColumnIndex] structure. | ||
| /// | ||
| /// The column index can be costly to decode and store, especially when it is needed | ||
| /// only for a subset of row groups or columns (such as when filtering by a predicate | ||
| /// on a single column). Providing a [`ColumnChunkMask`] can greatly decrease | ||
| /// the time needed to decode this metadata. | ||
| /// | ||
| /// [ColumnIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | ||
| pub fn with_column_index_mask(mut self, mask: ColumnChunkMask) -> Self { | ||
| self.column_index_mask = mask; | ||
| self | ||
| } | ||
|
|
||
| /// Sets the [`ColumnChunkMask`] for the Parquet [OffsetIndex] structure. | ||
| /// | ||
| /// The offset index can be costly to decode and store, especially when it is needed | ||
| /// only for a subset of row groups or columns (such as when projecting a small subset | ||
| /// of columns). Providing a [`ColumnChunkMask`] can greatly decrease | ||
| /// the time needed to decode this metadata. | ||
| /// | ||
| /// [OffsetIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | ||
| pub fn with_offset_index_mask(mut self, mask: ColumnChunkMask) -> Self { | ||
| self.offset_index_mask = mask; | ||
| self | ||
| } |
There was a problem hiding this comment.
Docs / Bug (C3). These docs suggest a column index mask for the predicate columns and an offset index mask for the projected columns. With that setup, StatisticsConverter returns arrays of different lengths (see the review body):
CI columns([0]) (predicate), OI columns([1]) (projection): data_page_mins len 30, data_page_row_counts len 0
Page pruning needs the offset index of the predicate columns too. The suggestion also adds two other traps (C7): a mask without a policy does nothing, and ArrowReaderMetadata::try_new ignores the masks.
| /// Sets the [`ColumnChunkMask`] for the Parquet [ColumnIndex] structure. | |
| /// | |
| /// The column index can be costly to decode and store, especially when it is needed | |
| /// only for a subset of row groups or columns (such as when filtering by a predicate | |
| /// on a single column). Providing a [`ColumnChunkMask`] can greatly decrease | |
| /// the time needed to decode this metadata. | |
| /// | |
| /// [ColumnIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | |
| pub fn with_column_index_mask(mut self, mask: ColumnChunkMask) -> Self { | |
| self.column_index_mask = mask; | |
| self | |
| } | |
| /// Sets the [`ColumnChunkMask`] for the Parquet [OffsetIndex] structure. | |
| /// | |
| /// The offset index can be costly to decode and store, especially when it is needed | |
| /// only for a subset of row groups or columns (such as when projecting a small subset | |
| /// of columns). Providing a [`ColumnChunkMask`] can greatly decrease | |
| /// the time needed to decode this metadata. | |
| /// | |
| /// [OffsetIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | |
| pub fn with_offset_index_mask(mut self, mask: ColumnChunkMask) -> Self { | |
| self.offset_index_mask = mask; | |
| self | |
| } | |
| /// Sets the [`ColumnChunkMask`] for the Parquet [ColumnIndex] structure. | |
| /// | |
| /// The column index can be costly to decode and store, especially when it is needed | |
| /// only for a subset of row groups or columns (such as when filtering by a predicate | |
| /// on a single column). Providing a [`ColumnChunkMask`] can greatly decrease | |
| /// the time needed to decode this metadata. | |
| /// | |
| /// The mask applies only if the column index policy is not [`PageIndexPolicy::Skip`] | |
| /// (the default), and only when the page index is loaded with these options (for | |
| /// example by [`ArrowReaderMetadata::load`]). [`ArrowReaderMetadata::try_new`] ignores it. | |
| /// | |
| /// [ColumnIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | |
| pub fn with_column_index_mask(mut self, mask: ColumnChunkMask) -> Self { | |
| self.column_index_mask = mask; | |
| self | |
| } | |
| /// Sets the [`ColumnChunkMask`] for the Parquet [OffsetIndex] structure. | |
| /// | |
| /// The offset index can be costly to decode and store, especially when it is needed | |
| /// only for a subset of row groups or columns (such as when projecting a small subset | |
| /// of columns). Providing a [`ColumnChunkMask`] can greatly decrease | |
| /// the time needed to decode this metadata. | |
| /// | |
| /// Page pruning with the column index also needs the offset index of the same | |
| /// columns. So include the predicate columns in this mask, not only the projected | |
| /// columns. The same notes as for [`Self::with_column_index_mask`] apply. | |
| /// | |
| /// [OffsetIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | |
| pub fn with_offset_index_mask(mut self, mask: ColumnChunkMask) -> Self { | |
| self.offset_index_mask = mask; | |
| self | |
| } |
There was a problem hiding this comment.
I think this one's getting a little pedantic 😅. I mean, I think users can figure out that a function that explicitly states that it won't load the page index won't be affected by options controlling the loading of the page index.
Fair point re noting that you'll want the offset index for the predicate columns as well as the projected columns.
| fn iter_to_set(indices: impl IntoIterator<Item = i32>) -> Option<Arc<BTreeSet<i32>>> { | ||
| let set: BTreeSet<i32> = indices.into_iter().filter(|&i| i >= 0).collect(); | ||
| (!set.is_empty()).then_some(Arc::new(set)) | ||
| } |
There was a problem hiding this comment.
API, one-way door (C4). An empty set selects everything. The common callers build the set from a list that can be empty:
columns([]) -> loads 12 of 12 offset indexes (SELECT count(*): no projected leaves)
row_groups(<no surviving rgs>) -> loads 12 of 12 (stats pruned every row group)
columns([-1]) -> loads 12 of 12
columns([100]) -> loads 0 of 12, page_index() = None
columns([-1, 100]) -> loads 0 of 12
The neighbour APIs do the opposite: ProjectionMask::leaves(schema, []) selects no leaves, with_row_groups(vec![]) reads nothing, and ParquetStatisticsPolicy::skip_except(&[]) is SkipAll. After a release, a change here silently changes what callers load. Suggestion: an empty set selects nothing, and all() stays the only way to select everything. test_chunk_mask (L1545-1549) and the constructor docs (L155-156, L166-167, L177-179) need the same change.
| fn iter_to_set(indices: impl IntoIterator<Item = i32>) -> Option<Arc<BTreeSet<i32>>> { | |
| let set: BTreeSet<i32> = indices.into_iter().filter(|&i| i >= 0).collect(); | |
| (!set.is_empty()).then_some(Arc::new(set)) | |
| } | |
| fn iter_to_set(indices: impl IntoIterator<Item = i32>) -> Option<Arc<BTreeSet<i32>>> { | |
| // An empty set selects nothing. `Self::all()` selects everything. | |
| Some(Arc::new(indices.into_iter().filter(|&i| i >= 0).collect())) | |
| } |
There was a problem hiding this comment.
I'm fine with revisiting this. I think I was just trying to avoid users inadvertently specifying empty sets, but I agree that empty (as opposed to None) should mean just that.
| /// | ||
| /// Any indices in `columns` that are less than zero will be ignored. Passing an empty | ||
| /// set is treated the same as selecting all columns. | ||
| pub fn columns(columns: impl IntoIterator<Item = i32>) -> Self { |
There was a problem hiding this comment.
API, one-way door (C5). The constructors take i32. The other APIs around this one use usize: includes_* here, ProjectionMask::leaves, with_row_groups, PageIndexProvider::column_index, ParquetStatisticsPolicy::skip_except. So every caller writes .map(|i| i as i32). That cast wraps silently, and a negative result then means "ignore" (and, with C4, "all"). A related detail: ColumnChunkMask::all().includes_column(u32::MAX as usize) is false.
Suggestion: take impl IntoIterator<Item = usize> in the public API. Keep i32/u32 storage internal if you want it. A change after release breaks every caller, so this is the time to decide.
There was a problem hiding this comment.
Well, the issue is parquet doesn't allow arrays to be sized beyond i32::MAX. It was probably a mistake to use usize elsewhere, but I guess that horse is out of the barn. I'll rethink this.
| for col_idx in 0..rg.num_columns() { | ||
| if !mask.includes_column(col_idx) { | ||
| continue; | ||
| } |
There was a problem hiding this comment.
Perf nit (C15). For each selected row group, this loop visits every column and does a BTreeSet lookup, also when the mask selects 1 of 10,000 columns. add_ranges in push_decoder.rs does the same, on every try_decode. A pub(crate) helper can iterate the selected set directly (compiled and tested on this head):
// reader.rs
impl ColumnChunkMask {
/// Selected column indexes in `0..num_columns`, in order
pub(crate) fn column_indices(
&self,
num_columns: usize,
) -> Box<dyn Iterator<Item = usize> + '_> {
match &self.columns {
None => Box::new(0..num_columns),
Some(set) => {
let end = i32::try_from(num_columns).unwrap_or(i32::MAX);
Box::new(set.range(0..end).map(|&i| i as usize))
}
}
}
}
// here, and in parse_offset_index / add_ranges
for col_idx in mask.column_indices(rg.num_columns()) {| assert_eq!(metadata.memory_size(), 11326); | ||
| #[cfg(feature = "encryption")] | ||
| assert_eq!(metadata.memory_size(), 11750); | ||
| } |
There was a problem hiding this comment.
Test gap (C14). No test sets a mask through ArrowReaderOptions or ParquetMetaDataReader::with_page_index_mask. Mutation testing (on the #11159 head, same code) shows it: ArrowReaderOptions::with_column_index_mask, with_offset_index_mask, both getters, and ParquetMetaDataReader::with_page_index_mask can each return Default::default(), and all tests still pass. If a refactor drops the mask plumbing, readers decode the full index again and no test fails. There, these two tests kill all 5 mutants. They pass on this head.
Also missing: a data read (push decoder or async) with a mask. That test would have found C2. Add it after the C2 fix.
| } | |
| } | |
| /// Asserts which cells of the page index are populated | |
| fn assert_page_index_cells( | |
| metadata: &parquet::file::metadata::ParquetMetaData, | |
| expect_ci: impl Fn(usize, usize) -> bool, | |
| expect_oi: impl Fn(usize, usize) -> bool, | |
| ) { | |
| let page_index = metadata.page_index().expect("page index should be loaded"); | |
| let num_cols = metadata.file_metadata().schema_descr().num_columns(); | |
| for rg in 0..metadata.num_row_groups() { | |
| for col in 0..num_cols { | |
| let ci = page_index.column_index(rg, col).is_some(); | |
| let oi = page_index.offset_index(rg, col).is_some(); | |
| assert_eq!(ci, expect_ci(rg, col), "column index rg={rg} col={col}"); | |
| assert_eq!(oi, expect_oi(rg, col), "offset index rg={rg} col={col}"); | |
| } | |
| } | |
| } | |
| #[test] | |
| fn test_arrow_reader_options_page_index_masks() { | |
| use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions}; | |
| let file = create_test_file(); | |
| let column_mask = ColumnChunkMask::row_groups_and_columns([1], [0]); | |
| let offset_mask = ColumnChunkMask::columns([2]); | |
| let options = ArrowReaderOptions::new() | |
| .with_page_index_policy(PageIndexPolicy::Required) | |
| .with_column_index_mask(column_mask.clone()) | |
| .with_offset_index_mask(offset_mask.clone()); | |
| assert_eq!(options.column_index_mask(), &column_mask); | |
| assert_eq!(options.offset_index_mask(), &offset_mask); | |
| let expect_ci = |rg: usize, col: usize| rg == 1 && col == 0; | |
| let expect_oi = |_rg: usize, col: usize| col == 2; | |
| // sync reader builders | |
| let arrow_metadata = ArrowReaderMetadata::load(&file, options.clone()).unwrap(); | |
| assert_page_index_cells(arrow_metadata.metadata(), expect_ci, expect_oi); | |
| // `AsyncFileReader::get_metadata` implementations | |
| let metadata = ParquetMetaDataReader::new() | |
| .with_arrow_reader_options(Some(&options)) | |
| .parse_and_finish(&file) | |
| .unwrap(); | |
| assert_page_index_cells(&metadata, expect_ci, expect_oi); | |
| } | |
| #[test] | |
| fn test_parse_with_page_index_mask() { | |
| let file = create_test_file(); | |
| let metadata = ParquetMetaDataReader::new() | |
| .with_page_index_policy(PageIndexPolicy::Required) | |
| .with_page_index_mask(ColumnChunkMask::row_groups_and_columns([2], [1, 3])) | |
| .parse_and_finish(&file) | |
| .unwrap(); | |
| let expect = |rg: usize, col: usize| rg == 2 && (col == 1 || col == 3); | |
| assert_page_index_cells(&metadata, expect, expect); | |
| } |
| #[test] | ||
| fn test_chunk_mask() { |
There was a problem hiding this comment.
Test gap (N1). test_chunk_mask is in mod async_tests, which has #[cfg(all(feature = "async", feature = "arrow", test))]. The default features of parquet do not include async, so cargo test -p parquet does not run it. The test uses no async code. Move it to mod tests (L977).
| /// | ||
| /// Any indices in `row_groups` or `columns` that are less than zero will be ignored. | ||
| /// Passing an empty set for `row_groups` is treated as selecting all row groups, and | ||
| /// an empty set for `columns` as selectiong all columns. |
There was a problem hiding this comment.
Nit (C16). Typo.
| /// an empty set for `columns` as selectiong all columns. | |
| /// an empty set for `columns` as selecting all columns. |
| /// * `bytes` - The byte slice containing the page index data. | ||
| /// * `bytes` - [`PushBuffers`] that should have already been populated with the bytes containing | ||
| /// the page indexes. | ||
| /// * `start_offset` - The offset where `bytes` begin in the file. |
There was a problem hiding this comment.
Nit (C16). start_offset is no longer a parameter.
| /// * `start_offset` - The offset where `bytes` begin in the file. |
…to page_index_selection
This reverts commit dc23d99.
|
Ok, I reverted the page index merging, and instead opted for the "replace" path for dealing with repeated calls to |
| /// | ||
| /// Returns None if no page indexes are needed | ||
| pub fn range_for_page_index( | ||
| pub(crate) fn range_for_page_index( |
There was a problem hiding this comment.
another pub function that isn't really
|
I will try and review this later today or tomorrow |
…ex (apache#11182) # Which issue does this PR close? - None. There is no separate issue. This PR is related to apache#11132 and apache#11157. I found the bug while I reviewed them. # Rationale for this change The async reader (`ParquetRecordBatchStreamBuilder`) and the push decoder (`ParquetPushDecoderBuilder`) fail on a valid read. The sync reader (`ParquetRecordBatchReaderBuilder`) reads the same data correctly. Example: `file` is a valid Parquet file with the columns `a`, `b` and `c` (400 rows, 2 row groups, 50 rows per page). Columns `a` and `c` have an offset index. Column `b` does not have an offset index. ```rust let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional); let selection = RowSelection::from(vec![RowSelector::skip(150), RowSelector::select(100)]); // Sync reader: OK, 100 rows let batches = ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options.clone())? .with_row_selection(selection.clone()) .build()? .collect::<Result<Vec<_>, _>>()?; // Async reader: Err(General("Invalid column index 2, column was not fetched")) let batches: Vec<RecordBatch> = ParquetRecordBatchStreamBuilder::new_with_options(std::io::Cursor::new(file), options) .await? .with_row_selection(selection) .build()? .try_collect() .await?; ``` <details><summary>How to make <code>file</code></summary> ```rust use arrow_array::{ArrayRef, Int32Array, RecordBatch}; use bytes::Bytes; use parquet::arrow::ArrowWriter; use parquet::column::writer::ColumnCloseResult; use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader}; use parquet::file::properties::WriterProperties; use parquet::file::writer::SerializedFileWriter; use std::sync::Arc; fn make_file() -> Bytes { // A file with a full page index let column = |start| Arc::new(Int32Array::from_iter_values(start..start + 400)) as ArrayRef; let batch = RecordBatch::try_from_iter([("a", column(0)), ("b", column(400)), ("c", column(800))]) .unwrap(); let props = WriterProperties::builder() .set_max_row_group_row_count(Some(200)) .set_data_page_row_count_limit(50) .set_write_batch_size(50) .build(); let mut buf = Vec::new(); let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props)).unwrap(); writer.write(&batch).unwrap(); writer.close().unwrap(); let source = Bytes::from(buf); // Copy the column chunks to a new file, without the page index of column `b` let metadata = ParquetMetaDataReader::new() .with_page_index_policy(PageIndexPolicy::Required) .parse_and_finish(&source) .unwrap(); let page_index = metadata.page_index().unwrap(); let schema = metadata.file_metadata().schema_descr().root_schema_ptr(); let mut buf = Vec::new(); let mut writer = SerializedFileWriter::new(&mut buf, schema, Default::default()).unwrap(); for (rg, rg_meta) in metadata.row_groups().iter().enumerate() { let mut rg_writer = writer.next_row_group().unwrap(); for (col, col_meta) in rg_meta.columns().iter().enumerate() { let keep = col != 1; let close = ColumnCloseResult { bytes_written: col_meta.compressed_size() as u64, rows_written: rg_meta.num_rows() as u64, metadata: col_meta.clone(), bloom_filter: None, column_index: page_index.column_index(rg, col).filter(|_| keep).cloned(), offset_index: page_index.offset_index(rg, col).filter(|_| keep).cloned(), }; rg_writer.append_column(&source, close).unwrap(); } rg_writer.close().unwrap(); } writer.close().unwrap(); Bytes::from(buf) } ``` </details> | Reader | Before (`main`, 60.0.0) | After (this PR) | |---|---|---| | Sync (`ParquetRecordBatchReaderBuilder`) | 100 rows | 100 rows | | Async (`ParquetRecordBatchStreamBuilder`) | Error: `Invalid column index 2, column was not fetched` | 100 rows | | Push decoder (`ParquetPushDecoderBuilder`) | Error: `Invalid column index 2, column was not fetched` | 100 rows | The error occurs when all of these conditions are true: 1. The reader loads the page index (`PageIndexPolicy::Optional` or `Required`). 2. In a row group, some of the column chunks to fetch have an offset index, and some do not. 3. The read has a `RowSelection` or a `RowFilter`. After a predicate, the reader always fetches the remaining columns with a selection. Condition 2 occurs with a valid file in which only some column chunks have an offset index, or with a custom `PageIndexProvider`. The offset index mask in apache#11157 makes condition 2 a usual case. Release 60.0.0 is affected. The bug started with apache#10719 (commit apache@5ec9eaf033), which made the offset index an `Option` for each column chunk. Release 59.3.0 is not affected. ## Cause `InMemoryRowGroup::fetch_ranges` fetches the full column chunk when the chunk has no offset index, but it does not push an entry to `page_start_offsets`. `fill_column_chunks` reads `page_start_offsets` by position. Thus the entries go to the wrong columns, and the last column gets no data. | Column | Offset index | `page_start_offsets` entry | Before this PR | After this PR | |---|---|---|---|---| | `a` | Yes | `offsets_a` | `a` gets `offsets_a` | `Sparse` with `offsets_a` | | `b` | No | Missing (now `None`) | `b` gets `offsets_c` (wrong) | `Dense` (full chunk) | | `c` | Yes | `offsets_c` | `c` gets no data: error | `Sparse` with `offsets_c` | # What changes are included in this PR? - `page_start_offsets` is now `Option<Vec<Option<Vec<u64>>>>` in `FetchRanges`, `fill_column_chunks` and the push decoder `DataRequest`. - `fetch_ranges` pushes `None` for a column chunk without an offset index. - `fill_column_chunks` stores a `None` entry as `ColumnChunkData::Dense` (the full chunk). It stores a `Some(offsets)` entry as `ColumnChunkData::Sparse`, as before. A `Sparse` chunk with one range at the chunk start does not work. Without page locations, `SerializedPageReader` reads each page at its own offset, and `Sparse` accepts only an exact page start. That change gives the error `Invalid offset in sparse column chunk data: ..., no matching page found`. # Are these changes tested? Yes. The new file `parquet/tests/arrow_reader/partial_offset_index.rs` makes a valid file in memory in which only some column chunks have an offset index. It uses `SerializedRowGroupWriter::append_column` for this. Each test reads rows 150..250 with the sync reader (control), the async reader and the push decoder. Each test uses a `RowSelection`, then a `RowFilter`, and compares the output with the expected rows. | Test | Column chunks without an offset index | Async reader and push decoder on `main` | |---|---|---| | `test_no_offset_index_first_column` | `a` (before the indexed columns) | Fail with `RowSelection` | | `test_no_offset_index_middle_column` | `b` (between the indexed columns) | Fail with `RowSelection` and with `RowFilter` | | `test_no_offset_index_last_column` | `c` (after the indexed columns) | Fail with `RowSelection` and with `RowFilter` | | `test_no_offset_index_row_group` | All columns in row group 1 | Fail with `RowSelection` and with `RowFilter` | All 4 tests pass with this PR. With only a `RowFilter` on `a`, the first-column case passes on `main`, because the predicate step already fetched column `a`. # Are there any user-facing changes? No API changes. Reads that failed with `Invalid column index N, column was not fetched` now return the correct rows. Note on AI use: Claude Code wrote the fix and the tests. The author reviewed them. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
sunchao
left a comment
There was a problem hiding this comment.
Reviewed c47bb75 with Codex assistance. No blocking correctness issues found. Mask semantics, partial-index consumption in sync/async/push readers, statistics alignment, preload plumbing, and repeated-load replacement look consistent. The explicit replacement behavior keeps the design straightforward.
Validation: 68 focused tests passed with async/object_store enabled (46 metadata unit tests, 16 page-index-related integration tests, and 6 partial-offset-index tests). One temporary test exercised 27 combinations of disjoint row selections, index masks, and sync/async/push readers, including explicit Mask and Selectors policies. All six reported CI workflows passed at this head. Full-workspace tests and performance benchmarks were not rerun.
Left two optional performance/simplification suggestions below; neither blocks approval. The single covering read and dense row-group × column storage remain limitations for sparse selections, with the latter inherited from the existing representation.
| rg_idx, | ||
| col_idx, | ||
| )?; | ||
| let idx_bytes = bytes.get_bytes(r.start, (r.end - r.start) as usize)?; |
There was a problem hiding this 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.
There was a problem hiding this comment.
Yeah, this can be reverted. I was hoping to at some point reintroduce more targeted ranges for the page index, but that may be over-engineering. I could see at least wanting one range for the column index and a second for the offset index if the column selection is small and the table very wide. Later work can handle that case.
| Self::axis_indices(self.columns.as_deref(), num_columns) | ||
| } | ||
|
|
||
| fn axis_indices(axis: Option<&[u32]>, len: usize) -> Box<dyn Iterator<Item = usize> + '_> { |
There was a problem hiding this 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.
There was a problem hiding this comment.
I had Codex whip up some tests to show the difference. Added a concrete iterator in ad86518
alamb
left a comment
There was a problem hiding this comment.
Thanks @etseidl -- this is looking great.
I didn't quite finish my review today, but I am getting close . I am also running some performance tests with some vibe code in alamb/parquet_footer_parsing#2 to see if this PR actually delivers the proposed API (don't read indexes). I suspect it will
| .with_offset_index_mask(mask) | ||
| } | ||
|
|
||
| /// Sets the [`ColumnChunkMask`] for the Parquet [ColumnIndex] structure. |
| (offset_index_policy, offset_index_mask, false), | ||
| ] { | ||
| if policy != PageIndexPolicy::Skip { | ||
| for row_group in mask.row_group_indices(metadata.num_row_groups()) { |
There was a problem hiding this comment.
nit is you can could reduce the indent with something like
if policy == PageIndexPolicy::Skip {
continue;
}
alamb
left a comment
There was a problem hiding this comment.
All in all, this is great @etseidl -- thank you. ALso thanks to @adriangb and @sunchao for the reviews
Benchmark results 📈
I updated my parsing time benchmark :
My interpretation of the results is that this PR realizes the theoretical savings of not parsing the entire page index (this PR is the green bar):

Concern
Clarifies and standardizes repeated page index reads: each successful call now replaces the existing page index state rather than implicitly preserving previously loaded indexes.
I think some users may be bitten by this implicit change in behavior.
it sounds like you already tried to preserve the existing index and merge in the newly requested indexes and that got costly. However, it seems like If users want the "clear existing index" behavior they could explicitly clear the indexes before using the MetadataDecoder.
I think this is fine, but maybe it is worth seeing if we can preserve whatever the existing behavior is (I think "do nothing if the index doesn't exist 🤔 )
| /// honored by loading APIs such as [`ArrowReaderMetadata::load`]; | ||
| /// [`ArrowReaderMetadata::try_new`] does not load or filter page indexes. | ||
| /// | ||
| /// [ColumnIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md |
There was a problem hiding this comment.
The interaction between PageIndexPolicy and the ColumnChunkMask is complicated, especially when there is an existing PageIndex in ParquetMetadata
One potential API we could consider is a new variant in PageIndexPolicy::Mask to avoid having to thread through another set of variables.
pub enum PageIndexPolicy {
/// Do not read the page index.
#[default]
Skip,
...
/// Ensure the specified indexes exist, and error if not
Required(ColumnChunkMask),
/// Load the specified indexes if they exist
Optional(ColumnChunkMask),
}However, that would be an API change as PageIndexPolicy isn't marked as non_exhastive so we would have to either
- Wait for the next major release
- Go with this API, and then unify PageIndexPolicy, and deprecate
with_column_index_masket al methods in next major version 🤔
There was a problem hiding this comment.
My original design was something similar, but instead added three new variants to the policy which allowed for column only selection, row only selection, or both for a grid. The ColumnChunkMask grew out of that as a way to get the behavior without a breaking API change. I didn't want to code that up, but LLMs make the plumbing changes too easy 😅.
Anyway, I agree that's a cleaner solution that requires less added baggage. I guess the question is how soon do we want this? Given the work @mkleen is doing with direct parsing of the column index in #11285, maybe there's less urgency here and we can wait for the breaking window to open for 61.0.0.
I'm not in love with option 2.
Tangent: since the change to optional cells, I've wondered if we still need the Required vs Optional distinction any longer. The policy enum was introduced when offset indexes were all or nothing, and we needed to decide what to do if a file lacked some, but not all, cells since there was no way to have a placeholder for a missing cell. Now that each cell is optional, the offset index behavior can mirror the column index. If we're waiting for a breaking change, perhaps we can deprecate Required and only support Optional 🤷
There was a problem hiding this comment.
Tangent: since the change to optional cells, I've wondered if we still need the Required vs Optional distinction any longer. The policy enum was introduced when offset indexes were all or nothing, and we needed to decide what to do if a file lacked some, but not all, cells since there was no way to have a placeholder for a missing cell. Now that each cell is optional, the offset index behavior can mirror the column index. If we're waiting for a breaking change, perhaps we can deprecate Required and only support Optional 🤷
I think at least at some point in the past the Required variant was to fail fast when a user didn't correctly configure / provide the PageIndex to the parquet reader and the caller wanted to ensure they were properly using the index. Maybe the APIs are sophisticated enough these days that we could just add a flag to ArrowReaderOptions 🤔 . I agree the idea of Required for the metadata reader doesn't really make a lot of sense.
There was a problem hiding this comment.
The origin story is that back in the day, the column index had a NONE placeholder variant, so a missing column index was replaced with that. The offset index had no such placeholder, as the belief was all offset indexes would be populated (which makes sense, you can't do page pruning without them). The offset index parser simply panic'd if an individual cell was missing. To get rid of the panic, the policy was introduced; if the policy was Required, then an error was returned if an index was missing, if Optional then the entire offset index would be thrown away and decoding stopped. Now that individual cells of the offset index can be None, there's less need for Required. I think we should at least change the policy to Optional when the old with_page_index(true) is used.
| 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]>>, |
There was a problem hiding this comment.
I recommend documenting what the [u32] means and the invariants. I think it is something like "sorted indexes"
I was somewhat confused on first read if it was using bitmaps; It is fine that it isn't, but if we want to get fancy / need more space efficiency here in the future we could switch to bitmaps
There was a problem hiding this comment.
I'll look at the ColumnChunkMask docs and see what needs better clarification. I don't know that usize vs u32 matters so much here, but it does matter more when we use a similar construct in the Grid storage (#11159). No need to store the wider type in the metadata when parquet itself doesn't support array indices greater than 2^31 - 1
| /// 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. |
There was a problem hiding this comment.
It may also be worth mentioning this is cheap to clone (some Arcs)
There was a problem hiding this comment.
While we're on cloning, if you have a moment to opine on #11297 I'd appreciate your thoughts.
| /// | ||
| /// [ColumnIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md | ||
| pub fn with_column_index_mask(mut self, mask: ColumnChunkMask) -> Self { | ||
| self.column_index_mask = mask; |
There was a problem hiding this comment.
is it worth making this API return Result and checking that the column index mask doesn't refer to out of bounds columns (e.g. some column index greater than the number of columns?)
There was a problem hiding this comment.
Sure, that's a good point. I think the PageIndexProvider interface suffers from being infallible (which I'm trying to address in #11159).
There was a problem hiding this comment.
Looked at this, but I think I'll leave it as-is. No other setters here return a Result, so it would be kind of weird. Also, we won't have the footer yet, so we can't really know the proper bounds. Also, it's a mask...nothing wrong with mask bits set that will always be 0 after an AND. It's probably more important to do bounds checking where the mask is used.
|
Thanks for the detailed review @alamb 🙏 I do want to get this right before merging, so I think it's worth discussing the replacement behavior.
So one issue is we unwittingly introduced a change in behavior in 60.0.0. Prior to that, the column and offset indexes were separate monolithic entities. The push decoder separately parsed them, so they were preserved on subsequent calls where the policy was 60.0.0 merged both indexes into a single So we've already lost some of the preservation that existed prior. @adriangb's review pointed this out and how an early version of this PR exacerbated it. This is finding C11 above.
I did, and the problem wasn't in this PR, but when changing the storage format. The current nested Vec storage allows for easy index merging since there is a slot already allocated for each cell. Subsequent calls could update the cells they target and leave the others alone. The problem really arises if we try to get fancy with storage and only allocate enough space for the requested subgrid; later trying to append columns or row groups requires a reallocation and move of the existing cells. But maybe that's putting the cart before the horse. The issue is with sparse indexes, why waste storage on things we don't want. The current nested vec is awful. We need to at least change to a single allocation and calculate positions manually. But an So perhaps the preserve path isn't so bad.* As C11 points out, we need to do something and document it. The previous behavior was not a contract, just a consequence of how things were implemented. 60.0.0 introduced a behavior change that we should either revert or document anyway. I guess at this point I'm still open to either path, either what I have now (always replace), or go back to preserve across multiple calls, and take that into account as we try to make the back-end storage more efficient. Sounds like @alamb is a vote for the latter; I abstain. Other votes? @adriangb, @zhuqi-lucas, @sunchao? * I just remembered another issue with preservation. The old page index is behind an |
|
Another thought is we could divorce the page index parsing from the |
Another thing to think about is how to (eventually) unify BloomFilter into this whole thing -- specifically if the BloomFilter should be part of ParquetMetadata or not. It is kind of nice to be able to pass around one ParquetMetadata thing that has pointers to the various metadata structures, but it is weird that it is not a 1:1 correspondence of wht is in the footer |
Maybe we are worrying about old behavior that no one really cares about 🤔 I think this will be more important once people starting to use the incremental page index parsing. I think the usecase would go something like:
Ideally it will be possible / easy to just parse the PageIndex for |
Ugh, I keep diving into the weeds on this response. TL;DR is I agree, but I no longer think the current metadata parser is the place to handle this. The parser should parse components, and something else should cache and synthesize those parts into a form usable by the rest of the crate to return record batches.
Again, I agree. ParquetMetaData in my mind would eventually wrap a Bloom filter provider, which would follow the pattern we settle on for the page indexes. |
add is_none function enforce indices <= i32::MAX, not u32::MAX
Which issue does this PR close?
Rationale for this change
The Parquet Page Index can be quite large, but in many cases (column projection, row group pruning) only a small subset of the structure is required.
What changes are included in this PR?
This PR introduces a new struct
ColumnChunkMaskwhich can be used to identify the subset of the column and offset indexes that need to be retrieved and parsed.This also modifies the
ParquetMetaDataPushDecoderto request ranges based on the offset and column index masks to reduce unnecessary I/O.Clarifies and standardizes repeated page index reads: each successful call now replaces the existing page index state rather than implicitly preserving previously loaded indexes.
Are these changes tested?
Yes
Are there any user-facing changes?
There is a behavior change to repeated page index reads as mentioned above, but the prior behavior was not documented.
Adds some new public APIs, but no breaking changes to existing public APIs
AI assistance
Claude Code and Codex were used in the development of this PR, but I take responsibility for the end result.