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]

Reply via email to