dwsmith1983 opened a new issue, #25923:
URL: https://github.com/apache/datafusion/issues/25923

   ### Is your feature request related to a problem or challenge?
   
   Comet runs Spark's `EXISTS` / `NOT EXISTS` with a residual predicate as a 
`SortMergeJoinExec` with join type LeftSemi or LeftAnti and a `JoinFilter`. On 
a nested TPC-H q21 at SF1000 that query is 37% slower with Comet than with 
Spark (apache/datafusion-comet#6467). Its two correlated subqueries join 
lineitem to itself on the order key with the filter `l3.l_suppkey <> 
l1.l_suppkey`, so the key groups hold 1 to 7 rows per side.
   
   A standalone benchmark of just these two joins on 55.1.0 points at two 
things in `sort_merge_join/bitwise_stream.rs`.
   
   **1. The per-call cost of the filter dominates on small groups.** 
`evaluate_filter_for_inner_row` runs once per inner row of a key group. Each 
call builds a `ScalarValue`, broadcasts it to the length of the outer group, 
builds a `RecordBatch` and evaluates the expression. With groups this small, 
that fixed cost is most of the join.
   
   Setup: 1M keys, outer side 2.4M rows, inner side 4.0M rows (semi) and 2.4M 
rows (anti), i64 key and i64 `suppkey`, filter `outer.suppkey <> 
inner.suppkey`, batch size 8192, one thread, median of 5 runs.
   
   | Join | SMJ with filter | SMJ without filter | HashJoinExec with filter |
   |---|---|---|---|
   | LeftSemi | 445 ms | 58 ms | 82-93 ms |
   | LeftAnti | 463 ms | 48 ms | 46-56 ms |
   
   That is about 275 ns per filter evaluation, at 1.57 evaluations per matched 
group for the semi join and 1.74 for the anti join. 54.1.0 gives the same 
numbers.
   
   As an experiment I batched the evaluation across key groups: queue the outer 
and inner row indices, evaluate the filter once per 8192 or so pairs, and OR 
the results back per outer row. That brought the semi join to 266 ms and the 
anti join to 238 ms with identical output. Evaluating once per key group over 
the cross product did not help (506 ms and 474 ms), because the cost per call 
stays. What remains after batching looks like overhead per group: the 
synchronous fast path is skipped for every group once a filter is present, and 
each group slices both inputs.
   
   #25581 (with #25584) summarizes the inner group for comparison filters, but 
only for groups of at least seven inner rows, so it does not reach groups this 
small.
   
   **2. Each buffered key group is charged for its whole parent batch.** The 
inner key buffer reserves `slice.get_array_memory_size()`, which for a slice 
reports the capacity of the parent buffers. A group of at most 7 rows, about 64 
bytes of data, is charged 131,264 bytes (a full 8192-row batch), and 262,528 
bytes when it spans two batches. Under a tight pool every matched group then 
spills to its own file. With a `GreedyMemoryPool` of 64 KB over 20K keys there 
were 18,128 spills, and the anti join went from 9.2 ms to 1,157 ms. With 200 KB 
only the groups that span a batch boundary spilled. In a real plan the join 
shares the pool with the sorts in the same task, so this can be reached well 
before the data needs to spill.
   
   ### Describe the solution you'd like
   
   - Evaluate the join filter for many outer and inner row pairs at once 
instead of once per inner row, and keep a fast path for small filtered groups.
   - Charge a buffered key group for the rows it holds rather than for its 
parent buffers, for example by sizing the sliced rows or compacting small 
groups before buffering them.
   
   ### Describe alternatives you've considered
   
   A planner can choose `HashJoinExec` for these joins, which handles the 
filter much faster here, but it cannot spill yet (#24768).
   
   ### Additional context
   
   The benchmark is a small Rust binary against datafusion-physical-plan 55.1.0 
and arrow 59.3.0. The data is shaped like lineitem rows per order: 1 to 7 rows 
per key, random `suppkey` values out of 10M, and 60% of the rows on the outer 
side, like q21's `l_receiptdate > l_commitdate`.
   


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