SubhamSinghal opened a new pull request, #25839:
URL: https://github.com/apache/datafusion/pull/25839

   ## Which issue does this PR close?
   
   - Part of #18221 (and the PWMJ epic #17427).
   
   ## Rationale for this change
   
   There is no benchmark for `PiecewiseMergeJoin`'s classic joins (`INNER` / 
`LEFT` / `RIGHT` / `FULL`). `nlj`, `hj` and `smj` each have a suite under 
`benchmarks/sql_benchmarks/`, but PWMJ doesn't. It is also off by default 
(`enable_piecewise_merge_join = false`), so none of the existing suites reach 
it.
   
   The follow-up #PERF_PR changes two things in the classic stream: how it 
finds each streamed row's first match, and how it handles rows that can't 
match. That change needs a benchmark on `main` to compare against.
   
   ## What changes are included in this PR?
   
   A new `pwmj` suite in `benchmarks/sql_benchmarks/pwmj/`, plus one row in 
that directory's README. It doesn't change any library code.
   
   Every query is `SELECT lhs.payload, rhs.payload FROM lhs <JOIN> rhs ON 
lhs.key < rhs.key`:
   
   - **`lhs` (buffered):** the keys `1..=100000`, scrambled with a fixed 
multiplicative hash, so every run and every branch sees the same data. A 
streamed key `b` matches exactly the buffered keys `1..b`.
   - **`rhs` (streamed):** 2M rows. The five subgroups are match regimes. They 
separate the two costs per streamed row: finding its first match, and building 
its output.
   
   | subgroup | queries | streamed keys | what it exercises |
   |---|---|---|---|
   | `no_match` | Q01–Q04 | all below every buffered key | rows that can't 
match anything |
   | `null_heavy` | Q05–Q08 | half NULL, the rest as `no_match` | NULL keys |
   | `half_match` | Q09–Q12 | half `no_match`, half `selective` | a mix |
   | `selective` | Q13–Q16 | each matches only the 1–4 smallest buffered keys | 
a first match at the far end of the buffered side, with small output |
   | `all_match` | Q17–Q20 | all above every buffered key; only 2K rows, since 
the output is `lhs × rhs` = 2×10⁸ rows | output-bound; the first match is the 
first buffered row |
   
   Each subgroup runs `INNER`, `LEFT`, `RIGHT` and `FULL`, in that order. All 
20 files share `pwmj.benchmark.template`:
   
   - `load` builds both tables in memory from `range()`, so there's no data to 
generate and key generation is never timed.
   - `init` turns on `enable_piecewise_merge_join`, and an `assert` checks it 
took effect. Without the flag every query would plan as a nested-loop join.
   - An `assert` checks the join's row count against the value that follows 
from the keys. A wrong result fails here rather than showing up as a plausible 
timing.
   - `expect_plan` pins `PiecewiseMergeJoinExec`, `operator: Lt` and the 
`join_type`. If the planner swapped the inputs, the query would measure the 
mirrored join.
   
   The runner re-plans every iteration and counts rows as they stream, so the 
buffered side is rebuilt each time and `all_match`'s 2×10⁸ rows are never held 
in memory.
   
   There is no nested-loop arm. At these sizes it is 2×10¹¹ predicate 
evaluations.
   
   ## What is the testing strategy for this PR?
   
   ## Are there any user-facing changes?
   
   No.


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


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to