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]
