xudong963 opened a new issue, #24157: URL: https://github.com/apache/datafusion/issues/24157
### Is your feature request related to a problem or challenge? Interval-overlap joins and similar workloads commonly contain two cross-input inequality predicates, for example: ```sql left.start_at <= right.end_at AND right.start_at <= left.end_at ``` DataFusion can currently exploit a single range predicate with `PiecewiseMergeJoinExec`. With two range predicates, an unkeyed join may use `NestedLoopJoinExec`, while a join with equality keys may use `HashJoinExec` and evaluate both inequalities as residual filters. These plans can examine or materialize a large candidate set even when the final result is small. This is a concrete implementation follow-up to the broader range-join performance issue #8393. ### Describe the solution you'd like Add an experimental `IEJoinExec` physical operator for bounded inner joins driven by exactly two cross-input `<`, `<=`, `>`, or `>=` predicates. The initial implementation should: - Support unkeyed joins and optional equality keys. - Normalize range predicates to left-input expression, operator, and right-input expression. - Use two sorted orders, an inverse permutation, and a sparse bitmap to enumerate matches. - Recheck complete equality keys after hash grouping so hash collisions cannot affect correctness. - Preserve any additional join predicates as residual filters. - Exclude volatile expressions from IEJoin driver selection. - Emit bounded output batches, bound work performed by each stream poll, account for memory-pool usage, and return controlled resource-exhaustion errors. - Be selected by the physical planner only when an experimental configuration option is enabled; the option should default to `false`. The initial scope excludes outer, semi, anti, and null-aware joins, as well as spilling. ### Describe alternatives you've considered - Continue using `NestedLoopJoinExec` for unkeyed joins. This is general but may evaluate a quadratic candidate space. - Use `HashJoinExec` for equality keys and apply inequalities as residual filters. This helps keyed workloads but can still produce a large candidate set within each key. - Extend `PiecewiseMergeJoinExec`. It is designed around a single range predicate and does not use both inequalities to prune candidates. - Wait for a more general range-join framework. IEJoin is a focused first step for the common two-inequality pattern and can remain opt-in while its scope evolves. ### Additional context Draft implementation and initial benchmark results are available in #23741. The implementation spans the physical operator, keyed and residual-filter support, planner integration, configuration and documentation, end-to-end tests, and benchmarks. To keep each change independently reviewable, it will be delivered as a stack of smaller PRs; the proposed stack is documented in a follow-up comment. -- 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]
