mkleen opened a new pull request, #11285:
URL: https://github.com/apache/arrow-rs/pull/11285
# Which issue does this PR close?
- Closes #9296.
# Rationale for this change
Today, getting page-level statistics as Arrow arrays copies the data twice:
1. `decode_column_index` reads the Thrift bytes into `ThriftColumnIndex`.
`PrimitiveColumnIndex::try_new` then copies the mins and maxes into typed
`Vec<T>`s, putting default values in place for null pages.
2. `data_page_mins`, `data_page_maxes`, `data_page_null_counts` and
`data_page_nan_counts` go through that `ColumnIndexMetaData` again and push
each value into an Arrow builder one at a time. Each call walks the index
separately.
# What changes are included in this PR?
This PR adds a way to skip loading the column index, then decode only the
columns you need straight into Arrow arrays. The decoder lives in the arrow
reader, so `ParquetMetaData`, `ColumnIndexMetaData` and `PageIndexProvider` are
unchanged and contain no Arrow types.
### New API
```rust
pub struct DataPageStatistics {
pub mins: ArrayRef,
pub maxes: ArrayRef,
pub null_counts: UInt64Array,
pub nan_counts: UInt64Array,
}
impl StatisticsConverter<'_> {
pub fn data_page_statistics_from_bytes<'b, I>(&self, column_indexes: I)
-> Result<DataPageStatistics>
where I: IntoIterator<Item = (usize, Option<&'b [u8]>)>;
}
```
- **Input:** one `(num_pages, bytes)` pair per row group. The bytes are the
column chunk's serialized `ColumnIndex`, found with
`ColumnChunkMetaData::column_index_range()`. To avoid parsing the index when
the metadata is loaded, load with `PageIndexPolicy::Skip`.
- **Output:** one call returns all four arrays, matching `data_page_mins`,
`data_page_maxes`, `data_page_null_counts` and `data_page_nan_counts`.
- **Missing data:** a row group with `None` bytes, or a column that isn't in
the Parquet file, gives `num_pages` nulls in all four arrays.
The decoder is `parquet/src/arrow/arrow_reader/statistics/page_index.rs`. It
reads the Thrift compact encoding of `ColumnIndex` in one pass, on top of
`ThriftSliceInputProtocol`.
| Field | Thrift type | Handling |
|---|---|---|
| 1 `null_pages` | `list<bool>` | Stored one byte per element (`0x01` =
true; `0x00` and `0x02` = false; any other byte is an error). Inverted to "page
has a min/max" and packed into the `NullBuffer` shared by mins and maxes. |
| 2 `min_values`, 3 `max_values` | `list<binary>` | Each element is decoded
straight into a buffer of the column's physical type (see below). The bytes of
null pages are skipped without being checked. |
| 4 `boundary_order` | `i32` | Read and checked, then discarded. |
| 5 `null_counts`, 8 `nan_counts` | `list<i64>` | Zigzag varints decoded
into `Vec<u64>`. A negative count is an error. A row group without the list
gives null entries. |
| 6, 7 level histograms, and unknown fields | any | Skipped. For integer
lists, skipping just counts the bytes whose top bit is clear, which is much
faster than decoding each varint. |
**Physical type to buffer:**
| Physical type | Buffer | How each value is decoded |
|---|---|---|
| `BOOLEAN` | `Vec<bool>` | First byte ≠ 0 |
| `INT32` | `Vec<i32>` | First 4 bytes, little-endian |
| `INT64` | `Vec<i64>` | First 8 bytes, little-endian |
| `FLOAT` | `Vec<f32>` | First 4 bytes, little-endian |
| `DOUBLE` | `Vec<f64>` | First 8 bytes, little-endian |
| `BYTE_ARRAY` / `FIXED_LEN_BYTE_ARRAY` read as `Decimal32`/`64`/`128`/`256`
| `Vec<i32>` … `Vec<i256>` | Big-endian two's complement, sign-extended while
reading (`from_bytes_to_i*`). There is no intermediate `BinaryArray`. A value
must be 1 to 4/8/16/32 bytes long. |
| Other `BYTE_ARRAY` / `FIXED_LEN_BYTE_ARRAY` | `i32` offsets + `Vec<u8>` →
`BinaryArray` | Bytes copied as they are. Fixed-length values are kept as
variable-length because a stored min or max may have been truncated. |
| `INT96` | Count only | Each value is checked to be at least 12 bytes. The
result is always nulls, as before. |
As in the old decoder, a fixed-width value with extra trailing bytes is
accepted, and one with too few bytes is an error with the same message: `error
converting value, expected N bytes got M`.
### Benchmark
`parquet/benches/arrow_statistics.rs` compares
`data_page_statistics_from_bytes` with `decode_column_index` plus all four
`data_page_*` calls, starting from the same raw bytes in both cases:
| Type | Pages | Existing route | `data_page_statistics_from_bytes` |
Speedup |
|---|---:|---:|---:|---:|
| Int64 | 2,000 | 29.7 µs | 7.0 µs | 4.22× |
| Int64 | 10,000 | 121.3 µs | 29.4 µs | 4.12× |
| Utf8 | 2,000 | 109.7 µs | 32.6 µs | 3.36× |
| Utf8 | 10,000 | 500.5 µs | 159.1 µs | 3.15× |
| Utf8View | 2,000 | 109.9 µs | 41.7 µs | 2.64× |
| Utf8View | 10,000 | 515.9 µs | 203.7 µs | 2.53× |
| Decimal128(20, 2) | 2,000 | 49.2 µs | 15.9 µs | 3.10× |
| Decimal128(20, 2) | 10,000 | 220.2 µs | 71.8 µs | 3.07× |
Each file has 20 row groups with 10 rows per page, and every 7th value is
null. The Utf8 values are 14 bytes, so as views they point into the data buffer
rather than being stored inline.
# Are these changes tested?
- **Unit tests** in `page_index.rs`:
- Every physical type is compared against every Arrow type with the
existing route
- Specific tests cover INT96, all-null pages, zero pages, no row groups, a
column missing from the file, unknown fields and histograms, and mins stored
before `null_pages`.
- Every error case is checked, including empty and over-wide decimal
values for both byte column types, and decoding is tried on every truncation of
a valid index.
- **Integration tests:** `parquet/tests/arrow_reader/statistics.rs` now also
checks every existing page-level case through
`data_page_statistics_from_bytes`, using the column index bytes read from the
file.
# Are there any user-facing changes?
Yes. This PR adds a new public API but nothing existing changes.
- New method StatisticsConverter::data_page_statistics_from_bytes.
- New struct parquet::arrow::arrow_reader::statistics::DataPageStatistics
## LLM-generated code disclosure
This PR includes LLM-generated code and comments. All LLM-generated content
has been manually reviewed.
--
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]