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]
