Skip to content

fix: keep a streamed side sorted outside the native plan in order through outer hash joins - #6835

Open
andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:andygrove/comet-inmemorytablescanexec-audit
Open

andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:andygrove/comet-inmemorytablescanexec-audit

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6787. First item of the in-memory cache audit, #6834.

Rationale for this change

Spark's hash joins report the streamed side's order as their own, so Spark drops a sort above the join that this order satisfies. DataFusion's hash join keeps the unmatched probe rows of a right outer join in place only when the probe input declares an order, and otherwise appends them after the matched rows of each batch (append_right_indices in datafusion-physical-plan). Comet builds every hash join on the build side, so Spark's LeftOuter with BuildRight and RightOuter with BuildLeft both run as that right outer join. A streamed side that enters the native plan from the JVM arrives through a ScanExec, which declares no order, so the join returned those rows out of order with no sort left in the plan to fix them.

#6787 reproduced this through spark.comet.convert.inMemoryCache.enabled, which is opt-in. Since #5634 made the in-memory cache default-on, CometInMemoryTableScan reaches the same path with default settings: a relation cached after sortWithinPartitions reports its order, and a left outer CometHashJoin or CometBroadcastHashJoin over it returns unsorted partitions where Spark's own read of the same cache is sorted.

A Comet-side outputOrdering override on the join cannot fix this: RemoveRedundantSorts runs before Comet converts the join, so the sort is already gone by then. The native join has to honor Spark's claim instead.

What changes are included in this PR?

  • HashJoin gets a streamed_sort_orders field. CometHashJoin.doConvert fills it with the streamed side's outputOrdering for the two outer join shapes above, and sends the join back to Spark with a fallback reason if that order cannot be expressed natively. Every other join type keeps the probe order as it is, so nothing is sent for them.
  • The planner declares that order on the probe input through a new pass-through SortedInputExec when the input declares none. The declaration is placed on that one input rather than on every ScanExec because an order visible to the whole native plan would also change how DataFusion's aggregates group. The dynamic join filter's consumer forwards its input's properties, so the declaration survives it.

How are these changes tested?

  • CometInMemoryCacheSuite: "hash joins over a sorted cache keep the streamed side's order", for both CometHashJoin and CometBroadcastHashJoin over the native cache scan. The left outer join must keep the partitions sorted with no sort in the plan, and the inner join must still report the cache's order so that no sort is added.
  • CometJoinSuite: "outer hash joins keep a streamed side sorted outside the native plan in order", the Native hash joins lose the streamed side's order when it arrives sorted from outside the native plan #6787 shape through CometSparkToColumnar.
  • Both new tests fail on main. The full CometJoinSuite and CometInMemoryCacheSuite pass locally on Spark 4.1 (151 tests), as do the planner unit tests and clippy.

…ough outer hash joins

Spark's hash joins report the streamed side's order as their own, so Spark
drops a sort above the join that the order satisfies. DataFusion's hash join
keeps the unmatched probe rows of a right outer join in place only when the
probe input declares an order, and otherwise moves them after the matched
rows of each batch. Comet builds every hash join on the build side, so
Spark's left outer join built on the right and right outer join built on
the left both run as that right outer join, and a streamed side that enters
the native plan from the JVM, a relation sorted before it was cached for
instance, arrives through a ScanExec that declares no order. The join then
returned those rows out of order with no sort left in the plan to fix them.

Send the streamed side's sort order with the join for those two shapes, and
have the planner declare it on the probe input through a pass-through
SortedInputExec when the input declares none. The declaration is placed on
that one input rather than on every ScanExec because an order visible to the
whole native plan would also change how DataFusion's aggregates group. A
sort order that cannot be expressed natively sends the join back to Spark.

With the in-memory cache on by default since apache#5634, the native cache scan
reaches this with default settings.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the full seven-file diff from bfacc43a50aba13da4fa8a828db52309541191b4 to 207f06a45c24ed15be8a372a4d507f7f04a4787b. The PR is not a draft. No introduced P1/P2 issues found within this review. There are no existing review comments or unresolved review blockers.

Summary

  • Prior state and problem: A sorted streamed input arriving through ScanExec lost its ordering metadata. Outer hash joins could then append unmatched rows after matched rows, violating ordering that Spark relied on after removing a redundant sort.
  • Design approach: Serialize the streamed ordering for LeftOuter/BuildRight and RightOuter/BuildLeft, then declare it on the native probe input with SortedInputExec when native ordering is absent.
  • Correctness: Checked Spark’s HashJoin.outputOrdering, ShuffledHashJoinExec exceptions, and RemoveRedundantSorts against DataFusion 55.1.0’s join implementation. A standalone native reproduction demonstrated the old reordering and passed 144 variants using the exact-head wrapper, checking order and row multiplicity across both build sides, ascending/descending order, null and duplicate join keys, residual filters, empty builds, and batch boundaries.
  • Compatibility analysis: The relevant ordering contracts agree across Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and the 4.2 branch. Protobuf field 11 is additive and defaults to an empty list. Compared the affected paths with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, which retains the same ordering defect. This is an intended correctness fix.
  • Key design decisions: Preserve Spark’s existing ordering contract because redundant-sort removal has already occurred. Fall back when the required ordering cannot be serialized, and retain any ordering the native input already declares.
  • Implementation sketch: Shared hash-join conversion populates streamed_sort_orders. The planner wraps the appropriate input before build-side swapping. SortedInputExec forwards execution unchanged and supplies ordering properties. Added Scala tests cover broadcast and shuffled joins through both cache input paths.
  • Performance: The wrapper adds no sort, batch copy, or buffering. Affected joins use DataFusion’s order-preserving index construction. No benchmark speedup is established, and no reproducible P1/P2 performance regression was identified.
  • Design: The declaration is localized to the join’s probe input, avoiding a blanket change to every scan’s properties. Both hash-join implementations share the conversion logic, and unrelated join types retain their previous path.
  • Abstraction & complexity: A small pass-through execution node and one planner helper express the missing property without introducing another sorting implementation or configuration surface. No complexity issue meeting the reporting bar was identified.
  • Behavioral changes worth calling out: Sorted JVM and cache inputs now retain their streamed order through the affected outer joins. An unsupported streamed-order expression can now cause the join to fall back to Spark. Inner joins and unordered inputs are unchanged.
  • Suggested improvements: No additional code change meeting the P1/P2 reporting bar was identified. The repository-required Spark SQL validation verdict remains outstanding before queueing.

Routed skills: review-comet-pr and review-comet-expression-pr. Also read AGENTS.md and the relevant contributor guidance for operators, serialization, development, and CI.

Exact-head CI, checked at 2026-10-10 01:19 UTC: 18 checks succeeded, 7 were running, and 14 were skipped. Native compilation and all Spark-profile Java compilation checks succeeded. Rust tests, four Spark 4.1 test groups, and TPC-H/TPC-DS verification were still running. Spark SQL suites for all versions and macOS were skipped. No completed check reported failure.

Validation limits: The local reproduction compiled the exact-head SortedInputExec against cached DataFusion 55.1.0 libraries. It did not exercise JVM serialization, cache integration, or the complete Comet planner. I did not run a full Comet build, Scala suites, or Spark SQL suites locally. CI therefore remains pending, not green. Project code and GitHub state were unchanged.

@andygrove
andygrove requested review from comphead and viirya October 10, 2026 12:07
@andygrove

Copy link
Copy Markdown
Member Author

@dwsmith1983 you worked on the native hash join ordering in #6810 and #6785, so this touches code you know. If you have time, would you be able to review it? Thanks.

@dwsmith1983

Copy link
Copy Markdown
Contributor

@dwsmith1983 you worked on the native hash join ordering in #6810 and #6785, so this touches code you know. If you have time, would you be able to review it? Thanks.

Declaring the order on the probe input is the better fix for the two outer shapes. DataFusion's hash join only checks whether the probe declares an order (self.right.output_ordering().is_some()) before keeping unmatched probe rows in place, so the declaration costs nothing, where #6810 sorts the whole streamed side whenever it arrives from the JVM. With the in-memory cache on by default since #5634 that sort now runs with default settings, so I'd rather this approach covers those shapes.

The declared order goes past the join. The doc on SortedInputExec says the order is declared on the probe "and nowhere else", but a Right join keeps its probe's ordering in its own output, so everything above it in the native plan sees it too. An aggregate grouped by the key switches to sorted, streaming mode, and the sort aggregate's ordering_satisfy check treats the order as already there. That is a win when Spark's claim holds. It also means a claim that is wrong turns into wrong groups rather than only wrong order. Could the comment say so, and could one test put an aggregate on the key above the join?

NOT IN is not covered, and it is a claim that does not hold. Spark reports the streamed side's order for a null-aware anti join, but Comet runs it with Spark's streamed side as DataFusion's build side, and DataFusion concatenates the build batches in reverse, so its output is not in that order. Declaring an order cannot fix that; it needs a sort. Under this PR an outer join above it would wrap that output in SortedInputExec with the reported order.

On this branch it gives wrong groups. A cache repartitioned on k and sorted within partitions, filtered with k NOT IN (...), then a broadcast left outer join on k and a GROUP BY k, all in one native plan: the native plan puts SortedInputExec: ordering=[col_0@0 ASC] over the null-aware LeftAnti join, and both aggregates run with ordering_mode=Sorted. With 2000 rows over 40 keys it returns 75 rows instead of 37, with 29 keys split into two groups whose counts add up to the right one. It also fails at the default batch sizes, as long as a partition spans more than one batch. main and #6810 return the right rows for the same query. One trap if you add a test: dense keys take DataFusion's perfect hash build path, which keeps the batches in order and hides the bug, so the keys need to be spread out ((i % 40) * 1000 is enough).

#6810 already has the NOT IN sort and its tests. Since this branch turns that case from wrong order into wrong groups, the sort needs to be in before or with this change. I'll narrow #6810 to the NOT IN case now so it can land first, using field 12 for now or field 11 if you'd rather this PR rebase onto it. Once both are in, the outer-shape sort in #6810 goes away in favour of the declaration here. Let me know if you agree with narrowing #6810.

Tests.

  • Both new tests use left_outer with BuildRight. RightOuter with BuildLeft takes the path that is not swapped (planner.rs around the BuildLeft branch) and has no test.
  • collect().toSet hides duplicate rows, which is the failure mode above. Comparing the rows as a multiset would catch it.
  • A planner test for with_declared_order: it wraps the right side, and it skips the wrap when the input already declares an order.
  • The fallback when the order cannot be expressed natively has no test.
  • The repo asks for a Spark SQL suite verdict on planner and serde changes, and the PR tier only ran 4.1.

The branch also conflicts with main in CometJoinSuite.scala.

  test("NOT IN below an outer hash join over a sorted cache keeps each group once") {
    withSQLConf(
      SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
      SQLConf.SHUFFLE_PARTITIONS.key -> "2",
      SQLConf.COLUMN_BATCH_SIZE.key -> "8",
      CometConf.COMET_BATCH_SIZE.key -> "8",
      CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true") {
      // Sparse keys keep DataFusion off its perfect hash build path, which keeps batch order.
      withParquetTable((0 until 2000).map(i => ((i % 40) * 1000, i)), "big") {
        withParquetTable(Seq(1000, 10000, 20000, 100000).map(x => Tuple1(Option(x))), "excl") {
          withParquetTable((0 until 40 by 2).map(k => (k * 1000, k)), "small") {
            try {
              spark
                .table("big")
                .toDF("k", "v")
                .repartition(2, $"k")
                .sortWithinPartitions("k")
                .cache()
                .createOrReplaceTempView("cached")
              spark.table("cached").count()
              val query =
                """SELECT /*+ BROADCAST(b) */ a.k, count(*), sum(a.v), count(b._2)
                  |FROM (SELECT * FROM cached WHERE k NOT IN (SELECT _1 FROM excl)) a
                  |LEFT JOIN small b ON a.k = b._1
                  |GROUP BY a.k""".stripMargin
              def rows(): Seq[(Int, Long, Long, Long)] =
                sql(query)
                  .collect()
                  .map(r => (r.getInt(0), r.getLong(1), r.getLong(2), r.getLong(3)))
                  .toSeq
                  .sorted
              val cometRows = rows()
              val sparkRows = withSQLConf(CometConf.COMET_ENABLED.key -> "false")(rows())
              assert(
                cometRows.map(_._1).distinct.size == cometRows.size,
                "a key is grouped twice")
              assert(cometRows == sparkRows)
            } finally {
              spark.catalog.clearCache()
            }
          }
        }
      }
    }
  }

…orytablescanexec-audit

# Conflicts:
#	spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native hash joins lose the streamed side's order when it arrives sorted from outside the native plan

3 participants