zhuxiangyi opened a new pull request, #10037:
URL: https://github.com/apache/paimon/pull/10037
### Purpose
Performance: a conditional V1 `UPDATE` on a data-evolution table currently
pays a self-join, a shuffle and a sort that the unconditional form does not,
because the `WHERE` condition keeps it off the self-merge shortcut.
**Why.** `UpdatePaimonDataEvolutionTableCommand` runs an `UPDATE` as a
`MERGE INTO t USING t ON t._ROW_ID = s._ROW_ID`.
`MergeIntoPaimonDataEvolutionTable` has a self-merge shortcut for exactly this
shape — `Scan → MergeRows → Write`, no join, no shuffle, no sort — but the
shortcut only recognises a source that is a plain projection of the table. The
`WHERE` condition was placed as a `Filter` on the source side (`USING (SELECT
_ROW_ID FROM t WHERE cond)`), so every conditional `UPDATE` fell through to the
general join path. What that costs, taken from the physical plans Spark
actually ran for `UPDATE t SET v = v + 1 WHERE id > 4` on a two-file table:
```
job 0 find touched files
HashAggregate(distinct UDF(_ROW_ID)) <- Exchange <- HashAggregate
<- Filter(id > 4) <- BatchScan t full scan
(filter pushed down)
job 1 write the column files
MapPartitions(DataEvolutionPaimonWriter)
<- Sort [_FIRST_ROW_ID, _ROW_ID] <- Exchange
hashpartitioning(_FIRST_ROW_ID)
<- MergeRowsExec
<- BroadcastHashJoin LeftOuter
:- BatchScan t PaimonSplitScan touched
files
+- BroadcastExchange <- Filter(id > 4) <- BatchScan t the
table scanned again
```
The shuffle in job 1 carries *every* row of every touched file, not only the
rows the condition selects, because a data-evolution column file must cover the
whole row-id range of its base file. With a large source side the broadcast
join becomes a sort-merge join and adds two more shuffles.
**What.** Carry the condition as the `WHEN MATCHED` condition of the
self-merge instead of as a source `Filter`:
```
MERGE INTO t USING t ON t._ROW_ID = s._ROW_ID
WHEN MATCHED AND cond THEN UPDATE SET ...
```
The source is then `Project(PaimonRelation)` and the shortcut applies. Rows
that fail the condition go through the shortcut's keep-copy instruction and are
written back unchanged, which is the `UPDATE ... WHERE` semantics; a NULL
condition value counts as not matched, as in a `Filter`. The same statement now
runs as one job:
```
MapPartitions(DataEvolutionPaimonWriter)
<- Project UDF(_FIRST_ROW_ID) <- MergeRowsExec <- BatchScan t
PaimonSplitScan
```
The `Filter` shape is kept, and the general path used, when the condition
cannot be evaluated inside `MergeRows`:
- a subquery (`WHERE id IN (SELECT ...)`), which can only be planned on a
regular scan;
- attributes that do not belong to the relation, e.g. the read-side `CHAR`
padding projection the analyzer inserts on top of it, which would be unresolved
in the merge plan;
- a condition without column references (a constant, `rand() < 0.1`): file
pruning has nothing to work with, and a constant-false condition would
otherwise rewrite every file as a no-op.
**Benefit.**
- One job instead of two, and the target is scanned once instead of three
times (touched-file discovery, the touched files themselves, and the filtered
scan on the source side).
- No shuffle or sort of the touched files' rows: the scan already yields
whole files in row-id order, which is what the column-file writer needs.
- File pruning comes for free: `selfMergeActionPredicate` pushes the action
condition to the snapshot reader, so files whose statistics rule out the
condition (including partition pruning) are never read. Before, the set of
touched files was only known after job 0 scanned the table and collected the
distinct first row ids on the driver.
- Snapshot pinning, row-id conflict detection and the conflict retry loop
are unchanged; they sit above the merge shape.
Two things noticed while testing that are pre-existing and out of scope
here: a V1 `UPDATE` on a data-evolution table with a `CHAR` column fails at
analysis on master (the padded attribute leaks into the aligned assignments in
`PaimonUpdateTable`), and a non-deterministic `WHERE` is rejected by
`CheckAnalysis` on the command node itself. Both behave identically before and
after this change.
### Tests
`RowTrackingTestBase`, run as `RowTrackingTest` on Spark 3.5 (66 tests)
together with `BlobUpdateTest`, `DataEvolutionDeletionTest` and
`DataEvolutionUpdateSnapshotTest`:
- `V1 update table with data-evolution` — the existing conditional case; its
assertion that a `Join` is present is flipped to `assertSelfMergeShortcut` (no
`Join`, `Sort` or `RepartitionByExpression`).
- `V1 update with condition prunes files through the self-merge shortcut` —
two files, `WHERE b >= 30 SET b = b + 1`: shortcut taken, `RESULTED_TABLE_FILES
== 1`, the condition sees the pre-update values, and a condition that matches
nothing in the surviving file copies it through unchanged.
- `V1 update with partition condition prunes partitions through the
shortcut` — partitioned table, `WHERE dt = 'p2' AND id = 4` scans one file.
- `V1 update condition on metadata column and NULL values` — `WHERE b > 10`
leaves the NULL row alone; `WHERE _ROW_ID = 1` takes the shortcut.
- `V1 update with constant false condition changes nothing` — `WHERE 1 = 0`
produces no snapshot.
- `V1 conditional update with user-specified snapshot uses the shortcut` —
`scan.snapshot-id` pin plus `WHERE`, rows inserted after the pin are untouched.
- `V1 update with subquery condition keeps the source filter` — `WHERE id IN
(SELECT ...)` still uses the general join path and returns the right rows.
- Existing `V1 update retries concurrent update conflicts` (`WHERE id = 1`
under 4 concurrent writers) and `DataEvolutionUpdateSnapshotTest` (conflict
detected after the snapshot is pinned) now exercise the shortcut and still pass.
`spotless:check` and `checkstyle:check` pass on both modules.
### API and Format
No changes.
### Documentation
No changes.
--
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]