Repository navigation
fix: keep a native hash join's streamed order when Spark drops the sort above it - #6810
dwsmith1983 wants to merge 18 commits into
Conversation
… is not kept Spark reports a hash join's streamed-side ordering and may drop a sort above the join on that basis. DataFusion keeps unmatched probe rows of a right join in order only when it sees the probe input sorted, which it does not when that order comes from outside the native plan. The new HashJoin output_ordering field carries the ordering, and the native planner adds a sort above the join when the join does not already produce it.
For LeftOuter with BuildRight and RightOuter with BuildLeft, Spark reports the streamed side's ordering, and DataFusion runs the join as a right join that may emit unmatched probe rows after the matched ones. The serde now sends that ordering so the native planner can restore it, and falls back to Spark when the ordering cannot be sorted natively.
The ordering the merge join produces reaches the hash join through the swapped right join and its projection, so the planner keeps it without sorting again.
…ordering tests A null-aware anti join runs unswapped, so DataFusion builds on Spark's streamed side and concatenates its batches in reverse; it now sends its ordering too and gets the same native sort. The tests cover multi-column and multi-batch orderings, the expression fallback, a natively sorted probe, join metrics with a sort on top, and the join class per hint. The compatibility guide describes how hash joins keep the order.
The null-aware anti join test now checks that its build spans several batches, a metric check that held for any batch size is gone, and a test comment claims only what the test checks. The guides say a NOT IN join is sorted too, and that only a NOT IN under OR, planned as an existence nested loop join, falls back to Spark.
…atches DataFusion builds a NOT IN join on Spark's streamed side and, on its hash path, concatenates those batches in reverse, so the planner sorts the output. The test feeds two sparse-key batches and checks the order.
A plain NOT IN runs as a native null-aware anti join; only one combined with another predicate through OR becomes an existence join that falls back to Spark.
|
@andygrove could you take a look when you get a chance? I'd also suggest |
sunchao
left a comment
There was a problem hiding this comment.
Reviewed the full seven-file diff from 00a4b422f0ed6560ef76dce004b94c3697613108 to 8482fc93b5b4a51fd506c44cc0adf731279504bf. The PR is not a draft. No introduced P1/P2 issues found within this review. Existing discussion contains no substantiated unresolved P1/P2 concern.
Summary
- Prior state and problem: Spark can remove a sort because a hash join reports its streamed input’s ordering. Native outer joins could move unmatched rows, while null-aware anti joins could reverse build-batch order, violating that promise.
- Design approach: Serialize the reported ordering for the affected join shapes and restore it with
SortExecwhen DataFusion’s equivalence properties cannot guarantee it. - Correctness: Checked Spark’s
HashJoin.outputOrdering,ShuffledHashJoinExecexceptions,RemoveRedundantSorts, and ordering tests against DataFusion’s execution paths. The sort follows the join’s output-restoring projection, so bound column indices remain valid. Existing sort serialization and normalization preserve direction, null placement, and floating-point handling. No P1/P2 issue identified. - Compatibility analysis: Checked the relevant sources for Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. The affected ordering contract is consistent across them.
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1has the original behavior, making this an intended correctness change and a backport candidate. The protobuf field is additive and defaults to empty. - Key design decisions: Restricting serialization to the two affected outer shapes and null-aware anti joins limits additional work. Unsupported ordering types or expressions fall back to Spark through existing compatibility checks.
- Implementation sketch:
CometHashJoin.doConvertbinds ordering expressions tojoin.output. The planner handles join swapping, checks ordering satisfaction, and conditionally wraps the result.additional_native_plansretains join metrics without double-counting output rows. - Performance: Unordered joins and joins whose native ordering already satisfies the requirement avoid the extra sort. Restoring missing order adds buffering, sorting, and two memory consumers, reducing other consumers’ shares under the fair pool. A bounded execution probe successfully exercised spilling. No throughput improvement or regression was inferred without measurements.
- Design: Reusing DataFusion’s spillable sorter keeps the correction local and follows the existing sort-aggregate approach. The implementation is straightforward to reason about.
- Abstraction & complexity:
sort_unless_orderedcleanly isolates the ordering check and wrapper construction. Reusing existing serialization, normalization, and metric aggregation avoids introducing another sorting abstraction. - Behavioral changes worth calling out: Affected joins now restore their advertised order across unmatched rows and batch boundaries. Unsupported orderings can newly trigger fallback. The
NOT INdocumentation correction changes neither configuration defaults nor execution routing. - Suggested improvements: No additional code change meeting the P1/P2 reporting bar was identified.
Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-memory-pr. Read AGENTS.md, applicable contributor guidance, the supplied discussion snapshot, and current GitHub reviews/comments. No Copilot feedback was used.
Exact-head CI: Label pull requests passed. Comet CI, Check PR Title, and CodeQL report action_required. There is no exact-head build/test verdict yet.
Validation: A disposable standalone probe using DataFusion 55.1.0 and Arrow 59.3.0 reproduced both outer-join reorderings and the sparse-key null-aware anti-join reordering, then verified that sorting restores order without changing rows. Both outer shapes avoided an extra sort with natively sorted probes. A 100,000-row case returned correctly ordered output under a 256 KiB pool with 14 spills. git diff --check passed.
Validation limits: The focused exact-head native test build failed during dependency compilation with No space left on device. Its artifacts were removed. Comet planner tests, CometJoinSuite, Spark SQL suites, and the actual Comet fair-memory-pool path were not independently executed. The standalone probe validates the DataFusion mechanism, not the complete JVM/protobuf/native integration. Author-reported suite passes remain unverified. Project code is unchanged.
comphead
left a comment
There was a problem hiding this comment.
I checked the DataFusion 55.1.0 ordering paths and Spark's outputOrdering for the affected shapes and found no correctness problem. The inline comments cover sharing the new sort helper with the sort aggregate path and folding one repeated test run into the existing matrix.
|
|
||
| /// Returns a sort of `plan` by the SortOrder exprs in `ordering`, or `None` when `ordering` | ||
| /// is empty or the plan's equivalence properties already satisfy it. | ||
| fn sort_unless_ordered( |
There was a problem hiding this comment.
This is the same check and wrap that the sort aggregate already does for ordered_by_grouping_keys (the ordering_satisfy block in the HashAgg arm, around line 1712). Could sort_unless_ordered take the PhysicalSortExprs, or a LexOrdering, instead of Spark Exprs, with the join mapping create_sort_expr at its call site? Then the aggregate arm could call it too, and the planner would keep one place that decides when an output sort is needed.
There was a problem hiding this comment.
Could
sort_unless_orderedtake thePhysicalSortExprs, or aLexOrdering, instead of SparkExprs, with the join mappingcreate_sort_exprat its call site?
Done. The helper takes the physical sort expressions, the hash join maps create_sort_expr where it calls it, and the sort aggregate arm now goes through the same helper.
| val (leftOuter, buildSide) = | ||
| checkJoinKeepsOrder(query("left_outer"), keyOrder(IntegerType, Ascending)) | ||
| assert(buildSide == BuildRight) | ||
| assert(leftOuter.nativeOp.getHashJoin.getOutputOrderingCount > 0) |
There was a problem hiding this comment.
The left_outer half of this test runs the same tables, cache, query and confs as the matrix case shuffle_hash left_outer ... condition=false, AQE=false, batchSize=8192 above, so it repeats that run. Could the getOutputOrderingCount > 0 assertion move into the matrix loop instead? That would check the field for every affected shape, including broadcast and right_outer, and this test could keep only the inner check.
There was a problem hiding this comment.
Could the
getOutputOrderingCount > 0assertion move into the matrix loop instead?
Moved. Every matrix case now asserts that the join sends its ordering, and the separate test keeps only the inner check.
…t aggregate sort_unless_ordered now takes physical sort expressions. The hash join builds them at its call site, and the sort aggregate arm calls the same helper instead of its own ordering_satisfy check and SortExec wrap.
The streamed order matrix and the null-aware anti join test assert that the native op carries an output ordering. The separate test keeps only the inner join check, since its outer half repeated a matrix case.
With runtime filters on, a probe-side filter or projection becomes CometFilterExec or CometProjectionExec before the join is built. Check that both keep the sort's ordering, so an outer hash join that reports its streamed ordering still adds no sort.
Which issue does this PR close?
Closes #6787.
Rationale for this change
Spark's
HashJoin.outputOrderingis the streamed side's ordering, and the Comet hash join execs report it too. When the streamed side already arrives sorted, Spark removes a sort above the join on that basis. The native join does not always keep that order:LeftOuterwithBuildRightandRightOuterwithBuildLeftas a DataFusionRightjoin, which keeps unmatched probe rows in order only when it sees the probe input sorted (datafusion-physical-plan55.1.0,hash_join/exec.rs:1574). When the order comes from outside the native plan, such as a cached sorted table, the unmatched rows come after the matched ones in each batch.NOT IN) runs unswapped, so DataFusion builds its hash table on Spark's streamed side. For sparse keys spanning 1024 or more values it takes its hash map path, which concatenates those batches in reverse (exec.rs:2344,2365), so across batches the output comes out in reverse order. Denser keys use an array map that keeps the order.sortWithinPartitionsandSORT BYabove such a join then return unsorted partitions.What changes are included in this PR?
HashJoingets anoutput_orderingfield. The serde fills it with Spark's reported ordering, bound to the join output, for the three shapes above, and only when that ordering is non-empty. It falls back to Spark when the ordering has a type the native sort does not support or an expression that cannot be serialized.SortExecabove the hash join when the join's equivalence properties do not already satisfy that ordering. The check lives in one planner helper, which the sort aggregate'sordered_by_grouping_keyshandling now calls too. When the streamed side is sorted inside the same native plan, for example by a native sort or a sort-merge join, DataFusion keeps the order and no sort is added. Float keys are normalized the way a native sort below would normalize them, so that comparison holds.The sort runs only when the order comes from the JVM, which takes a non-default conversion such as
spark.comet.convert.inMemoryCache.enabledorspark.comet.sparkToColumnar.enabled, or for aNOT INjoin whose streamed side reports an ordering. There it replaces unsorted output. The sort registers two memory consumers in the task, so under the fair pool the build side's share is smaller while it runs.branch-1.1has the same code path, so this is a backport candidate.Merge order: #6785 should merge first. Both change
CometJoinSuitenear the same tests, and once #6785 lands I'll merge main here and resolve that. Draft #6437 also takesHashJoinfield 11; whichever merges second renumbers.How are these changes tested?
CometJoinSuite, checking rows against Spark and every partition's order with Spark's own ordering:shuffle_hashandbroadcast,left_outer(BuildRight) andright_outer(BuildLeft), over a cached sorted streamed side, with Spark's sort above the join removed; plus a join condition, AQE on, and a small batch size so the sort merges several batches. The plan has no sort above the join, the join class matches the hint, and the join's build metrics still reach the Spark node.NOT INjoin over a cached sorted side whose build spans several batches.Planner tests: the two outer shapes get a sort and return rows in order; a probe sorted natively, by a sort or a sort-merge join, gets none, including on a float key; no ordering adds no sort; a
NOT INjoin over two sparse-key build batches comes out sorted. Withspark.comet.exec.join.dynamicFilter.enabledon, a filter or projection above that native sort becomes its Comet wrapper, and the join still adds no sort.The cached-side and
NOT INtests fail on main. Removing the shape restriction, the field, either fallback guard, theordering_satisfycheck or the float normalization each fails a test.CometJoinSuitepasses on Spark 3.4, 3.5, 4.0 and 4.1, Spark 4.2 compiles, and the strict warnings build passes. The Spark SQLsql_core-1shard (dev/local-ci.sh, Spark 4.1.3), which runsRemoveRedundantSortsSuiteand the inner, outer, existence and hint join suites, passes.