Repository navigation
feat: add ASOF join physical operator - #23828
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #23828 +/- ##
==========================================
+ Coverage 81.19% 81.21% +0.02%
==========================================
Files 1110 1114 +4
Lines 388618 393461 +4843
Branches 388618 393461 +4843
==========================================
+ Hits 315531 319542 +4011
- Misses 54507 55113 +606
- Partials 18580 18806 +226 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
2010YOUY01
left a comment
There was a problem hiding this comment.
Thank you! This is a really good start. I have done a quick first pass and left some suggestions.
| vec![ChildStats::At(partition), ChildStats::Skip] | ||
| } | ||
|
|
||
| fn statistics_from_inputs( |
There was a problem hiding this comment.
Just an idea to make this PR smaller, could we use the default implementation here? We could implement it later in a follow-up PR.
There was a problem hiding this comment.
I kept this small override because ASOF has two exact facts the default would discard: the output row count equals the left row count, and unmodified left columns retain their statistics. Right-side column statistics remain unknown. I added a comment to make that scope explicit.
There was a problem hiding this comment.
I see, this makes sense.
There was a problem hiding this comment.
Thanks! I also updated it to preserve the same contract through the general projection.
|
There is no new commits after the previous review, you may have forgotten to push the local changes 🤔 @Xuanwo |
2010YOUY01
left a comment
There was a problem hiding this comment.
I went over the implementation in detail, and I think it's well organized.
Need to do before merging
Before merging, I think we could add a few basic tests that run the executor and assert the results, just to cover some different cases:
- No equality condition in
on. - Different comparison operators in the match condition, such as
<and>=. - Complex expressions in the
onor match conditions, such asMATCH_CONDITION (l_c1 + l_c2) > r_c1orON l_c1 = (r_c1 + r_c2). I think these should be supported.
Optional suggestions
The main implementation/design questions I have are:
- We need to buffer
batch_sizeoutput rows before emitting them, now it's implemented through thePendingRowsstruct, so we need to buffer all source batches to produce the final output. I think this uses more memory and makes the implementation more complex. An alternative is a) keep ain_progress_batch, and emit it after it reaches threshould b) only buffer one (left_batch, right_batch), and materialize its valid indices into thein_progress_batchwhen the cursor moves across it. This might be fast enough and simpler. - Each loop iteration only advances the right index by one. We can probably explore some fast-forwarding optimization here.
But I suggest we first implement this end to end, including planning, SQL support, more tests, and benchmarks, before exploring these optimizations further. It should be easier to validate the ideas afterward.
For now, I only suggest trying to simplify the existing implementation or adding more documentation to make future iterations easier. I left a few suggestions in the comments.
|
|
||
| fn input_distribution_requirements(&self) -> InputDistributionRequirements { | ||
| InputDistributionRequirements::new(vec![ | ||
| Distribution::UnspecifiedDistribution, |
There was a problem hiding this comment.
Is it the case the optimizer will insert round-robin repartition, if we declare this UnspecifiedDistribution? It should be fine if it's doing so, I was a little bit confused by this name 🤔
| vec![ChildStats::At(partition), ChildStats::Skip] | ||
| } | ||
|
|
||
| fn statistics_from_inputs( |
There was a problem hiding this comment.
I see, this makes sense.
Co-authored-by: Yongting You <2010youy01@gmail.com>
|
Thanks for the detailed review — this was very helpful. I added result-based coverage for no equality keys, both |
Checked a bit, seems to be a mid level fix. Let's fix that in a seperate PR #24375. We can get this in first! 💌 |
jayzhan211
left a comment
There was a problem hiding this comment.
I hope we could have ASOFJoinStream that handles the AsOfJoinStreamState and performs the actual join. We could take inspiration from the existing join implementations, such as HashJoinStream.
| #[derive(Default)] | ||
| struct PendingRows { | ||
| /// Distinct source batches referenced by `indices`. | ||
| sources: Vec<Arc<RecordBatch>>, |
There was a problem hiding this comment.
Do we need Vec<Arc<RecordBatch>> and not just Vec<RecordBatch>
There was a problem hiding this comment.
Yep, Arc is intentional here: it keeps per-row clones O(1) and gives PendingRows a stable identity for deduplication. Added a short comment.
| /// flush pending rows without resetting either cursor or the candidate | ||
| /// ``` | ||
| async fn next_batch(&mut self) -> Result<Option<RecordBatch>> { | ||
| fn poll_next_impl( |
There was a problem hiding this comment.
I haven't fully figure out what's the best practice to implement the state machines to make them easy to read and extend, probably it is:
- generator pattern [Status Update] Simplifing Streams to be more textbook-like and have less state while keeping the same perf #23974
- If we need to draw a state machine to understand its mechanism, then use explicit state management like
Not an issue for now, we could experiment in the future.
|
cc @jayzhan211, @xudong963, and @alamb, would you like to take another look and move forward? |
datafusion project typically wait for 24 hrs after approval to merge, so others can have a chance to look. This one is a larger PR, probably we could wait for 2 days, unless anyone need more time to review. https://datafusion.apache.org/contributor-guide/pr_review.html#pr-review-mechanics |
|
Got it, thank you @2010YOUY01 for the explanation 🙌 |
jayzhan211
left a comment
There was a problem hiding this comment.
I think we could try to rewrite it with more straightforward way in the follow up PR 🤔
| ] | ||
| } | ||
|
|
||
| fn maintains_input_order(&self) -> Vec<bool> { |
There was a problem hiding this comment.
| fn maintains_input_order(&self) -> Vec<bool> { | |
| // ASOF emits exactly one row for each left row and never reorders the | |
| // left input. The right input is scanned independently. | |
| vec![true, false] |
Given @2010YOUY01 and @jayzhan211 have reviewed and approved I am happy to have this merged without my review (I am currently focusing a bit lower with arrow/parquet and trying to keep them going along) I did however add this feature to the list of things to list in the next release when we get there Nice work |
|
Thank you @2010YOUY01 again for getting this PR merged! Do you want to review the whole stack? Can I ping you on the PR while I'm ready? |
|
Sure, I'm interested in helping with this project 🫡 |
## Which issue does this PR close? - Part of #318. - Umbrella PR: #23738. - Depends on #23828 (merged). ## Rationale for this change This is the logical-planning layer of the ASOF JOIN stack. It defines the logical contract and planner behavior separately from the SQL frontend and serialization formats. #23828 is merged, so this PR's diff against `main` is the isolated logical layer. It no longer depends on the optional floating-point follow-up #24375. ## What changes are included in this PR? - Add `LogicalPlan::AsOfJoin`, `AsOfJoin`, and `AsOfMatch`. - Validate deterministic expressions, input ownership, supported match operators, equality-key types, and USING constraints. - Add `LogicalPlanBuilder` entry points and schema construction that preserves both qualified `USING` keys while exposing one unqualified wildcard key. - Integrate ASOF joins with tree transforms, display, type coercion, projection pruning, row bounds, and physical planning. - Plan logical ASOF joins to the broadcast-based `AsOfJoinExec` from #23828. - Fail closed at proto, SQL unparser, and Substrait boundaries until their owning stack layers add explicit support. - Defer ASOF-specific functional-dependency refinement to #24799 and filter pushdown to #24801 so each optimization can be reviewed independently. ## Are these changes tested? Yes: - `cargo fmt --all` - `cargo clippy --all-targets --all-features -- -D warnings` - `cargo test -p datafusion-expr min_rows_of_joins --all-features` - `cargo test -p datafusion-substrait asof_join_fails_closed_until_substrait_has_an_extension --all-features` - The extended workspace test command from the contributor guide ## Are there any user-facing changes? This adds logical-plan and builder APIs for ASOF joins. SQL syntax, DataFrame APIs, and plan serialization are intentionally left to dependent stack PRs. Floating equality keys remain rejected by the merged physical operator unless the independent follow-up #24375 is also included. As with any new public `LogicalPlan` variant, downstream exhaustive matches must add an arm. The variant is appended so existing variants retain their `PartialOrd` ordering; maintainers should still treat the enum addition as a Rust source-compatibility break. This PR can be reviewed independently now that #23828 has merged. The optimization follow-ups #24799 and #24801 are not required by the core ASOF stack.
## Which issue does this PR close? - Part of #318. - Umbrella PR: #23738. - Builds on #23829 and #23828, both merged. ## Rationale for this change This is the SQL frontend layer of the ASOF JOIN stack. It adds syntax and unparsing on top of the merged physical and logical contracts. The PR now contains only the isolated SQL frontend diff. ## What changes are included in this PR? - Plan `ASOF JOIN ... MATCH_CONDITION (...)` with optional `ON` or `USING` equality keys. - Reject unsupported match shapes and non-equality `ON` predicates. - Unparse ASOF joins while preserving right-side candidate preselection and nested join scope. - Document the supported SQL syntax and semantics, including qualified `USING` keys, left-partitioned broadcast execution, the full-right memory requirement, repeated scans, and the absence of spill/repartitioned ASOF. - Add SQL integration and sqllogictest coverage for all four match directions, coercion, equality-free joins, USING, invalid contracts, EXPLAIN, boundedness, and optimized-plan round trips. - Verify the broadcast topology with a multi-partition left input: left partitioning is preserved while the right input is single-partitioned. ## Are these changes tested? Yes: - `cargo fmt --all` - `./ci/scripts/doc_prettier_check.sh --write --allow-dirty` - `cargo clippy --all-targets --all-features -- -D warnings` - `cargo test -p datafusion --test core_integration asof --all-features` - `cargo test -p datafusion-sqllogictest --test sqllogictests --all-features -- asof_join` - The extended workspace test command from the contributor guide ## Are there any user-facing changes? Users can express Snowflake-style ASOF joins in SQL with `MATCH_CONDITION`, optional equality keys, and `<`, `<=`, `>`, or `>=` match directions. With `USING`, wildcard output exposes one unqualified key while both qualified input keys remain addressable. The user guide also documents the initial broadcast strategy and its memory/no-spill limitations. #23829 and #23828 are merged. This is the next core layer in the ASOF stack and does not depend on the optional floating-point follow-up #24375.
* feat: add ASOF join physical operator (apache#23828) ## Which issue does this PR close? - Part of apache#318. - Umbrella PR: apache#23738. - Follow-up for floating-point equality keys: apache#24375. ## Rationale for this change This is the first layer of the ASOF JOIN stack. It establishes a broadcast-based physical execution contract independently so later floating-point equality, logical-plan, SQL, DataFrame, and serialization changes can be reviewed as smaller follow-up PRs. The initial implementation deliberately favors the simpler broadcast design: the right input must fit in memory and each left partition scans the shared right-side batches. A repartitioned implementation can be evaluated separately without changing the ASOF semantics introduced here. Floating-point equality keys are rejected in this base layer because Arrow's required sort order distinguishes `-0.0` from `+0.0` while join equality does not. apache#24375 adds the required ordering normalization as an independently reviewable layer. ## What changes are included in this PR? - Add `AsOfJoinExec` for left-preserving, Snowflake-style ASOF semantics. - Coalesce and collect the ordered right input once, then share it across all left partitions. - Keep the left input partitioned so each partition can scan independently and preserve the left-side output partitioning. - Preserve merge state across input and output batch boundaries. - Reserve each retained Arrow buffer exactly once, including when right-side batches are zero-copy slices, and expose build, match, and output metrics. - Define output properties and statistics for the broadcast execution model. - Reject floating-point equality keys until apache#24375 supplies a sort/equality contract that handles signed zero correctly. - Add physical operator tests covering match directions, equality groups, batch boundaries, unmatched rows, invalid contracts, shared-buffer memory accounting, multi-partition broadcast execution, and float-key rejection. ## Are these changes tested? Yes: - `cargo fmt --all` - `cargo clippy --all-targets --all-features -- -D warnings` - `cargo test -p datafusion-physical-plan joins::asof_join --all-features` - Extended workspace tests from the contributor guide - FFI integration tests ## Are there any user-facing changes? This adds a new physical operator API. The base operator deliberately rejects floating-point equality keys; apache#24375 adds full Float16, Float32, and Float64 support. SQL and DataFrame APIs are left to later dependent PRs. --------- Co-authored-by: Yongting You <2010youy01@gmail.com> * feat: add ASOF join logical semantics (apache#23829) - Part of apache#318. - Umbrella PR: apache#23738. - Depends on apache#23828 (merged). This is the logical-planning layer of the ASOF JOIN stack. It defines the logical contract and planner behavior separately from the SQL frontend and serialization formats. logical layer. It no longer depends on the optional floating-point follow-up - Add `LogicalPlan::AsOfJoin`, `AsOfJoin`, and `AsOfMatch`. - Validate deterministic expressions, input ownership, supported match operators, equality-key types, and USING constraints. - Add `LogicalPlanBuilder` entry points and schema construction that preserves both qualified `USING` keys while exposing one unqualified wildcard key. - Integrate ASOF joins with tree transforms, display, type coercion, projection pruning, row bounds, and physical planning. - Plan logical ASOF joins to the broadcast-based `AsOfJoinExec` from - Fail closed at proto, SQL unparser, and Substrait boundaries until their owning stack layers add explicit support. - Defer ASOF-specific functional-dependency refinement to apache#24799 and filter pushdown to apache#24801 so each optimization can be reviewed independently. Yes: - `cargo fmt --all` - `cargo clippy --all-targets --all-features -- -D warnings` - `cargo test -p datafusion-expr min_rows_of_joins --all-features` - `cargo test -p datafusion-substrait asof_join_fails_closed_until_substrait_has_an_extension --all-features` - The extended workspace test command from the contributor guide This adds logical-plan and builder APIs for ASOF joins. SQL syntax, DataFrame APIs, and plan serialization are intentionally left to dependent stack PRs. Floating equality keys remain rejected by the merged physical operator unless the independent follow-up apache#24375 is also included. As with any new public `LogicalPlan` variant, downstream exhaustive matches must add an arm. The variant is appended so existing variants retain their `PartialOrd` ordering; maintainers should still treat the enum addition as a Rust source-compatibility break. This PR can be reviewed independently now that apache#23828 has merged. The optimization follow-ups apache#24799 and apache#24801 are not required by the core ASOF stack. * feat: support ASOF JOIN SQL (apache#23830) - Part of apache#318. - Umbrella PR: apache#23738. - Builds on apache#23829 and apache#23828, both merged. This is the SQL frontend layer of the ASOF JOIN stack. It adds syntax and unparsing on top of the merged physical and logical contracts. The PR now contains only the isolated SQL frontend diff. - Plan `ASOF JOIN ... MATCH_CONDITION (...)` with optional `ON` or `USING` equality keys. - Reject unsupported match shapes and non-equality `ON` predicates. - Unparse ASOF joins while preserving right-side candidate preselection and nested join scope. - Document the supported SQL syntax and semantics, including qualified `USING` keys, left-partitioned broadcast execution, the full-right memory requirement, repeated scans, and the absence of spill/repartitioned ASOF. - Add SQL integration and sqllogictest coverage for all four match directions, coercion, equality-free joins, USING, invalid contracts, EXPLAIN, boundedness, and optimized-plan round trips. - Verify the broadcast topology with a multi-partition left input: left partitioning is preserved while the right input is single-partitioned. Yes: - `cargo fmt --all` - `./ci/scripts/doc_prettier_check.sh --write --allow-dirty` - `cargo clippy --all-targets --all-features -- -D warnings` - `cargo test -p datafusion --test core_integration asof --all-features` - `cargo test -p datafusion-sqllogictest --test sqllogictests --all-features -- asof_join` - The extended workspace test command from the contributor guide Users can express Snowflake-style ASOF joins in SQL with `MATCH_CONDITION`, optional equality keys, and `<`, `<=`, `>`, or `>=` match directions. With `USING`, wildcard output exposes one unqualified key while both qualified input keys remain addressable. The user guide also documents the initial broadcast strategy and its memory/no-spill limitations. stack and does not depend on the optional floating-point follow-up apache#24375. --------- Co-authored-by: Xuanwo <github@xuanwo.io> Co-authored-by: Yongting You <2010youy01@gmail.com>
## Which issue does this PR close? - Part of apache#318. - Umbrella PR: apache#23738. - All required ASOF feature layers (apache#23828 through apache#23832) are merged. ## Rationale for this change This is the final core benchmark layer of the ASOF JOIN stack. It uses the standard SQL benchmark framework so the workloads share the existing runner, result format, and Criterion mode. ## What changes are included in this PR? - Add an `asof_join` SQL benchmark suite, runnable with `./bench.sh run asof_join` or `benchmark_runner`. - Add seven workloads covering relative input sizes, available ordering, equality-group cardinality and skew, match direction, and payload width. - Describe each workload and its explicit ASOF match semantics in the query. - Generate pre-sorted Parquet inputs under `DATA_DIR/asof_join` with `./bench.sh data asof_join`. The benchmark load hook only registers them with `WITH ORDER`, keeping data generation and input sorting outside the measured query. - Assert that every workload plans to `AsOfJoinExec`. ## Are these changes tested? Yes: - `cargo fmt --all` - `./ci/scripts/doc_prettier_check.sh --write --allow-dirty` - `cargo clippy --all-targets --all-features -- -D warnings` - `cargo test -p datafusion-benchmarks --lib --bins` - `DATA_DIR=/tmp/asof_join CARGO_COMMAND='cargo run' ./bench.sh data asof_join` - `DATA_DIR=/tmp/asof_join CARGO_COMMAND='cargo run' ./bench.sh run asof_join 7` - `cargo run -p datafusion-benchmarks --bin benchmark_runner -- asof_join --iterations 1 --partitions 4 --path /tmp/asof_join` - `cargo run -p datafusion-benchmarks --bin benchmark_runner -- asof_join --query 7 --partitions 4 --path /tmp/asof_join --criterion` The pre-sorted workload's physical plan contains `AsOfJoinExec` without an input `SortExec`. These smoke runs validate the suite and are not presented as performance claims. ## Are there any user-facing changes? This adds ASOF workloads to the repository benchmark tooling. It does not change query semantics or runtime behavior. This is the final core item tracked by apache#23738. It does not depend on the optional floating-point follow-up apache#24375.
## Which issue does this PR close? - Part of #318. - Umbrella PR: apache#23738. - Follow-up to the merged physical operator in apache#23828. ## Rationale for this change `AsOfJoinExec` already supports an internal output projection, but physical projection pushdown cannot currently populate it. Queries that need only a subset of ASOF output columns therefore materialize unused join columns before the outer projection removes them. ## What changes are included in this PR? - Implement the standard physical projection-embedding hook for `AsOfJoinExec`. - Preserve empty-projection row counts and decline unsafe repeated embedding. - Preserve projected plan properties, including dropping output ordering when its sort expressions are no longer available. - Add focused execution coverage and update the affected ASOF physical-plan expectations. ## Are these changes tested? Yes. The focused ASOF tests, the complete SQL logic test suite, Clippy with all targets and features, and the extended workspace test command from the contributor guide pass. I also compared the merged benchmark Q03 on a dedicated EC2 `c7i.4xlarge` with four partitions. Each revision used its own release build, three warmups, and 40 measured iterations in ABBA order. | Revision | Median | | --- | ---: | | `main` (`133111f43`) | 107.6 ms | | This PR (`72ee47918`) | 106.1 ms | That is a modest 1.3% improvement on Q03. The deterministic benefit is visible in the physical plan: the required output columns are embedded in `AsOfJoinExec`, so unused join columns are not constructed. A residual `ProjectionExec` is retained only when aliases or expressions still require it. ## Are there any user-facing changes? No SQL or logical API changes. Physical plans can avoid materializing ASOF output columns that are not required downstream.

Which issue does this PR close?
Rationale for this change
This is the first layer of the ASOF JOIN stack. It establishes a
broadcast-based physical execution contract independently so later
floating-point equality, logical-plan, SQL, DataFrame, and serialization changes
can be reviewed as smaller follow-up PRs.
The initial implementation deliberately favors the simpler broadcast design:
the right input must fit in memory and each left partition scans the shared
right-side batches. A repartitioned implementation can be evaluated separately
without changing the ASOF semantics introduced here.
Floating-point equality keys are rejected in this base layer because Arrow's
required sort order distinguishes
-0.0from+0.0while join equality doesnot. #24375 adds the required ordering normalization as an independently
reviewable layer.
What changes are included in this PR?
AsOfJoinExecfor left-preserving, Snowflake-style ASOF semantics.left partitions.
preserve the left-side output partitioning.
batches are zero-copy slices, and expose build, match, and output metrics.
contract that handles signed zero correctly.
batch boundaries, unmatched rows, invalid contracts, shared-buffer memory
accounting, multi-partition broadcast execution, and float-key rejection.
Are these changes tested?
Yes:
cargo fmt --allcargo clippy --all-targets --all-features -- -D warningscargo test -p datafusion-physical-plan joins::asof_join --all-featuresAre there any user-facing changes?
This adds a new physical operator API. The base operator deliberately rejects
floating-point equality keys; #24375 adds full Float16, Float32, and Float64
support. SQL and DataFrame APIs are left to later dependent PRs.