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]