Repository navigation
Commit ae8a42a
committed
[SPARK-59716][SQL] Skip the per-match
### What changes were proposed in this pull request?
`SortMergeAsOfJoinScanner` rescans the buffered right-side group once per left row, and did
`bestMatch = rightRow.copy()` for every right row that improved on the match so far, although only
the final one is ever used.
That copy is only needed when `ExternalAppendOnlyUnsafeRowArray` has switched to its spillable
backing store, whose iterator re-points a single `UnsafeRow` on every `next()`. While the buffer is
in memory the iterator returns the distinct rows it stores, so a retained match stays valid.
This PR exposes that as `ExternalAppendOnlyUnsafeRowArray.isSpillBacked` and copies a match only
when it is true. The flag is read per scan rather than cached, because `clear()` drops the spillable
backing store between equi-key groups.
### Why are the changes needed?
The copies are pure overhead on the common, non-spilled path, and there are `O(matches)` of them per
left row.
`AsOfJoinBenchmark`, `Best Time(ms)` of the "Sort-merge AS-OF join" case, from the regenerated
result files in this PR. `Improvement` is the reduction in best time.
`AS-OF Join (left=10000, right=10000, groups=100)`
| JDK | Before | After | Improvement |
|---|---|---|---|
| 17 | 53 | 44 | 17.0% |
| 21 | 56 | 49 | 12.5% |
| 25 | 53 | 53 | 0.0% |
`AS-OF Join (left=10000, right=10000, groups=10)`
| JDK | Before | After | Improvement |
|---|---|---|---|
| 17 | 170 | 110 | 35.3% |
| 21 | 181 | 120 | 33.7% |
| 25 | 180 | 122 | 32.2% |
`AS-OF Join no equi-key (left=10000, right=10000)`
| JDK | Before | After | Improvement |
|---|---|---|---|
| 17 | 1378 | 804 | 41.7% |
| 21 | 1466 | 870 | 40.7% |
| 25 | 1457 | 942 | 35.3% |
The `groups=100` case buffers only ~100 right rows per group, so the scan is a small part of an
end-to-end time dominated by the shuffle and sort, and the result sits inside its own stdev either
way. The larger-group cases, where the per-left-row scan actually dominates, improve by 32-42%.
Reusing one `UnsafeRow` and copying the winning row's bytes into it was measured too, and did not
help: it removes the allocation but keeps the memory copy, which is what costs here.
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Pass the CIs.
### Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Opus 5
Closes #58888
Closes #58971 from dongjoon-hyun/SPARK-59716.
Authored-by: Dongjoon Hyun <dongjoon@apache.org>
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>UnsafeRow copy when the AS-OF join right-side buffer is in memory1 parent 30dc3dc commit ae8a42a
6 files changed
Lines changed: 82 additions & 20 deletions
File tree
- sql/core
- benchmarks
- src
- main/scala/org/apache/spark/sql/execution
- joins
- test/scala/org/apache/spark/sql
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
6 | 6 | | |
7 | 7 | | |
8 | 8 | | |
9 | | - | |
10 | | - | |
| 9 | + | |
| 10 | + | |
11 | 11 | | |
12 | 12 | | |
13 | 13 | | |
14 | 14 | | |
15 | 15 | | |
16 | | - | |
17 | | - | |
| 16 | + | |
| 17 | + | |
18 | 18 | | |
19 | 19 | | |
20 | 20 | | |
21 | 21 | | |
22 | 22 | | |
23 | | - | |
24 | | - | |
| 23 | + | |
| 24 | + | |
25 | 25 | | |
26 | 26 | | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
6 | 6 | | |
7 | 7 | | |
8 | 8 | | |
9 | | - | |
10 | | - | |
| 9 | + | |
| 10 | + | |
11 | 11 | | |
12 | 12 | | |
13 | 13 | | |
14 | 14 | | |
15 | 15 | | |
16 | | - | |
17 | | - | |
| 16 | + | |
| 17 | + | |
18 | 18 | | |
19 | 19 | | |
20 | 20 | | |
21 | 21 | | |
22 | 22 | | |
23 | | - | |
24 | | - | |
| 23 | + | |
| 24 | + | |
25 | 25 | | |
26 | 26 | | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
6 | 6 | | |
7 | 7 | | |
8 | 8 | | |
9 | | - | |
10 | | - | |
| 9 | + | |
| 10 | + | |
11 | 11 | | |
12 | 12 | | |
13 | 13 | | |
14 | 14 | | |
15 | 15 | | |
16 | | - | |
17 | | - | |
| 16 | + | |
| 17 | + | |
18 | 18 | | |
19 | 19 | | |
20 | 20 | | |
21 | 21 | | |
22 | 22 | | |
23 | | - | |
24 | | - | |
| 23 | + | |
| 24 | + | |
25 | 25 | | |
26 | 26 | | |
Lines changed: 11 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
99 | 99 | | |
100 | 100 | | |
101 | 101 | | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
102 | 113 | | |
103 | 114 | | |
104 | 115 | | |
| |||
Lines changed: 16 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
376 | 376 | | |
377 | 377 | | |
378 | 378 | | |
| 379 | + | |
| 380 | + | |
| 381 | + | |
| 382 | + | |
| 383 | + | |
| 384 | + | |
| 385 | + | |
| 386 | + | |
| 387 | + | |
| 388 | + | |
| 389 | + | |
| 390 | + | |
379 | 391 | | |
380 | 392 | | |
381 | 393 | | |
| |||
385 | 397 | | |
386 | 398 | | |
387 | 399 | | |
| 400 | + | |
388 | 401 | | |
389 | 402 | | |
390 | 403 | | |
| |||
399 | 412 | | |
400 | 413 | | |
401 | 414 | | |
402 | | - | |
| 415 | + | |
403 | 416 | | |
404 | 417 | | |
405 | 418 | | |
| |||
418 | 431 | | |
419 | 432 | | |
420 | 433 | | |
| 434 | + | |
421 | 435 | | |
422 | 436 | | |
423 | 437 | | |
| |||
434 | 448 | | |
435 | 449 | | |
436 | 450 | | |
437 | | - | |
| 451 | + | |
438 | 452 | | |
439 | 453 | | |
440 | 454 | | |
| |||
Lines changed: 37 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
820 | 820 | | |
821 | 821 | | |
822 | 822 | | |
| 823 | + | |
| 824 | + | |
| 825 | + | |
| 826 | + | |
| 827 | + | |
| 828 | + | |
| 829 | + | |
| 830 | + | |
| 831 | + | |
| 832 | + | |
| 833 | + | |
| 834 | + | |
| 835 | + | |
| 836 | + | |
| 837 | + | |
| 838 | + | |
| 839 | + | |
| 840 | + | |
| 841 | + | |
| 842 | + | |
| 843 | + | |
| 844 | + | |
| 845 | + | |
| 846 | + | |
| 847 | + | |
| 848 | + | |
| 849 | + | |
| 850 | + | |
| 851 | + | |
| 852 | + | |
| 853 | + | |
| 854 | + | |
| 855 | + | |
| 856 | + | |
| 857 | + | |
| 858 | + | |
| 859 | + | |
823 | 860 | | |
0 commit comments