zhuxiangyi opened a new pull request, #10394:
URL: https://github.com/apache/paimon/pull/10394

   ### Purpose
   
   **Problem.** A partial-column update of a data evolution table (for example
   `MERGE INTO ... UPDATE SET score = ...`) is stored as a separate file that 
covers the same row id
   range as the base file. Reading such a range stitches its files together by 
row position, so
   `DataEvolutionSplitRead` pushes no filter into the format readers of a 
merged group: dropping rows
   in only one of the readers would misalign them (see the class javadoc, and 
`createUnionReader`,
   which builds the format reader with no filters). The file index, and since 
#10289 a bitmap
   selection, can skip rows of a merged group, but the row group statistics and 
the page indexes of
   the Parquet files are not used at all.
   
   As a result, once a query reads an updated column together with any other 
column of a range, every
   filter of the query, including filters on columns that were never updated, 
is evaluated only after
   all rows of the range have been read. The same range stored in one file 
reads just the pages which
   can match. In the benchmark below a 1% range filter reads 29.2 MB when the 
ranges are merged and
   1.8 MB when they are single files.
   
   **Why not the existing options.**
   - Compacting a range into one file brings the pushdown back, but rewrites 
every row of every
     updated range, and has to be repeated after the next update. Producing the 
single file layout of
     the benchmark rewrites 408 MB.
   - #10289 needs a bitmap file index on each filtered column, and the bitmap 
index only answers `=`,
     `<>`, `IN`, `NOT IN` and null checks, not the range predicates on time or 
id columns which are
     the common selective filters.
   
   **What this PR does.** The rows which may match are computed for the whole 
group from the
   metadata of the files holding the latest values of the filtered columns 
(Parquet row group
   statistics, dictionaries, bloom filters and page indexes), and every file of 
the group then reads
   the same rows through the existing row range path. Every file skips the same 
pages and cuts the
   same rows, so the readers stay aligned.
   
   - `FormatReaderFactory#candidateRowRanges` returns the rows which may match 
from file metadata
     only, or null for formats which can not tell. `ParquetReaderFactory` 
implements it on top of the
     row group and column index filtering of `ParquetFileReader`.
   - Only the latest copy of a field may prune: an older copy is stale and 
would exclude rows which
     do match. The winning files are resolved by the logic of the file index 
pushdown, extracted into
     `DataEvolutionSplitRead#winningFieldIds`. A filter is evaluated on one 
file, so a filter whose
     fields are won by different files, such as an OR across them, is not used.
   - The candidate ranges are intersected with the row ranges of the split (for 
example from a global
     index) and with the bitmap selection of #10289, which is applied to the 
same readers. When they
     cover the whole group nothing changes, when they are empty the group is 
skipped.
   - The footer read for the candidate rows is reused when the readers of the 
group open the same
     file, through `FormatReaderFactory.Context#metadataCache`.
   - Failing to compute the candidate rows reads the whole group, so a lost or 
corrupt file is still
     reported by opening it, respecting `scan.ignore-lost-file` and 
`scan.ignore-corrupt-file`.
   
   New option `data-evolution.merged-read.stats-pushdown.enabled`, default true.
   
   **Scope and limitations.**
   - Parquet only, and groups of normal data files sharing one row id range. 
Groups with blob or
     vector-store files, and tables with nested field evolution, keep the 
current behavior.
   - It pays off only when the filtered column is clustered in write order, as 
time or id columns
     usually are, because pruning is by statistics. Otherwise nothing is 
pruned, see the benchmark.
   - Computing the candidate rows opens each file which wins a filtered column 
once more and reads
     the column index of the filtered columns. See the I/O pattern below.
   
   ### Tests
   
   - `ParquetCandidateRowRangesTest`: the ranges cover all matches and prune, 
file positions across
     row groups, no match and no filter, reading the selected rows, the footer 
read once with a cache.
   - `DataEvolutionStatsPushDownTest`, each case compared with the pushdown 
disabled and with an
     in-memory model:
     - a stale copy of the filtered column does not exclude matching rows
     - the newest of several updates of a column wins
     - filters whose fields are won by different files are intersected
     - an OR across files
     - deletion vectors
     - intersection with the row ranges of the split
     - a renamed column
     - row sidecars enabled
     - a lost or a corrupt update file behaves as without the pushdown, with 
and without
       `scan.ignore-lost-file` / `scan.ignore-corrupt-file`
     - the pushdown disabled
     - 30 rounds of random updates, deletions and filters
   - The existing data evolution and Parquet tests pass with the option enabled 
by default, including
     the bitmap pushdown tests of #10289.
   
   **Benchmark.** A table of 4,000,000 rows in 8 row id ranges of 500,000 rows 
(11 columns, Parquet
   with default settings). Every range has one base file and one update file 
rewriting `ts2`, `rnd2`
   and `d0`: `ts2` and `ts` increase with the row id, `rnd2` is random. The 
"single file" layout is
   the same table after compacting every range into one file. Queries read 
`id`, the filtered column,
   `s0` and `d0` with a single thread, and the filter is a range on one column. 
Apple M4, 16 GB,
   macOS, JDK 17, local SSD with a warm page cache. The time is the median of 7 
runs after a warm-up,
   reported as the median of 3 JVM runs, which differ by a few percent. A 
counting `FileIO` reports
   what is read from the Parquet data files. Matching rows and a checksum of 
their content are the
   same in all three layouts for every query.
   
   | filter | merged, off | merged, on | single file |
   |---|---|---|---|
   | `ts2` range, 0.1% of rows (updated, clustered) | 59 ms, 29.2 MB | 9 ms, 
1.2 MB | 9 ms, 1.2 MB |
   | `ts2` range, 1% | 63 ms, 29.2 MB | 10 ms, 1.8 MB | 9 ms, 1.8 MB |
   | `ts2` range, 10% | 72 ms, 29.2 MB | 50 ms, 12.3 MB | 38 ms, 12.3 MB |
   | `ts2` range, 50% | 199 ms, 73.0 MB | 175 ms, 59.0 MB | 152 ms, 59.0 MB |
   | `ts` range, 1% (not updated, clustered) | 29 ms, 14.6 MB | 8 ms, 1.8 MB | 
7 ms, 1.8 MB |
   | `rnd2` range, 1% (updated, random) | 224 ms, 121.8 MB | 228 ms, 121.8 MB | 
196 ms, 121.8 MB |
   | no filter | 327 ms, 116.8 MB | 327 ms, 116.8 MB | 284 ms, 116.8 MB |
   
   - The scan already skips the row id ranges whose statistics exclude the 
filter, so "merged, off"
     reads only the ranges overlapping the clustered filters (1 to 5 of the 8). 
This PR prunes inside
     them.
   - The gain follows the selectivity and the clustering: up to 6x faster and 
24x fewer bytes for a
     selective filter on a clustered column, nothing for a random column or 
without a filter.
   - The `ts` row shows that a filter on a column which was never updated also 
benefits.
   
   I/O pattern, number of reads on the input stream, seeks and opened streams, 
for the 1% `ts2`
   filter and the random filter:
   
   | | merged, off | merged, on | single file |
   |---|---|---|---|
   | `ts2` 1%: reads / seeks / opens | 34 / 16 / 4 | 3896 / 28 / 6 | 3268 / 22 
/ 2 |
   | `rnd2` 1%: reads / seeks / opens | 128 / 56 / 16 | 4264 / 72 / 24 | 13072 
/ 80 / 8 |
   
   Skipping pages turns a few large reads into many small ones, as it does for 
the single file layout,
   and computing the candidate rows opens each winning file once more (6 vs 4 
streams for 2 groups).
   When nothing can be pruned (`rnd2`) the cost is these extra opens and reads, 
within the run to run
   noise here (about 2% of the time). Both are cheap on a local disk. Their 
effect on an object store depends on the buffering of
   its `FileIO` and was not measured.
   


-- 
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