Skip to content

fix: use Spark-compatible shuffle for wide decimal hash keys - #6005

Open
sunchao wants to merge 5 commits into
apache:mainfrom
sunchao:dev/chao/codex/oss-wide-decimal-shuffle
Open

sunchao wants to merge 5 commits into
apache:mainfrom
sunchao:dev/chao/codex/oss-wide-decimal-shuffle

Conversation

@sunchao

@sunchao sunchao commented Sep 17, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Related to #3079. Spark-compatible native decimal hashing is tracked in #5994.

Rationale for this change

Native hashing of decimal keys with precision greater than 18 differs from Spark's partition assignments. A join can mix native and JVM shuffles because Comet selects the implementation independently for each exchange. Matching keys then reach different partitions and rows silently disappear: the Parquet/JSON regression returns zero rows without this guard and the expected four with it, with AQE enabled or disabled.

What changes are included in this PR?

  • Reject native hashing of wide-decimal keys, including nested leaves, when an exchange has multiple output partitions. auto uses Comet columnar shuffle when eligible; native and Celeborn use Spark shuffle.
  • Keep wide-decimal payloads, range keys, and single-partition exchanges native. One-partition hash exchanges already serialize as SinglePartition, so no decimal hashing occurs.
  • Exercise the aggregate-buffer repair merged in fix: revert unsafe partial aggregates after final fallback #5421 for collect_list and collect_set, preserving native input scans. Add mixed-input join coverage and extend the existing Celeborn fallback test.
  • Update affected fuzz expectations, plan snapshots, and user/contributor documentation. Cross-reference the SQL hash restriction without introducing a shared abstraction.
  • Add a focused COUNT benchmark to quantify this correctness guard independently of decimal AVG support. The native hashing fix in Support Spark-compatible native hash and xxhash64 for decimals with precision >18 #5994 is the path to removing the fallback and its conversion cost.

How are these changes tested?

Rebased onto main 98662215d, preserving main's native-library reload regression and new shuffle/AQE coverage while keeping the wide-decimal routing tests in CometNativeShuffleSuite.

At f682d0902, the root Maven reactor passed all 64 selected tests on JDK 21 / Spark 4.1.3: seven decimal-routing tests plus all 57 Celeborn planning tests. Coverage includes the precision boundary, single-partition exception, unaffected payload/range keys, nested keys, aggregate buffers, and the mixed Parquet/JSON join with AQE enabled and disabled. Spotless and whitespace checks passed.

Follow-up f79e2ced3 makes the benchmark row-count conversion explicit after the Spark 3.5 strict-warning job flagged numeric widening. The complete root reactor test-compile -Pspark-3.5 -Pstrict-warnings -DskipTests passes on JDK 21 after that change.

The native source and build-input trees match main exactly. Tests used the SHA-verified Linux native artifact 11309692999 from run 37219565655 at 98662215d, bundled before the JVM tests. No new native compilation was needed for this Scala-only rebase.

The focused COUNT benchmark is retained to isolate shuffle fallback from decimal AVG support. Current-head timing measurements and TPC-DS runtimes were not rerun. Previous timing tables have been removed because they used an older revision and native artifact. The plan snapshots and fuzz expectation updates remain in the PR but were not rerun locally in this pass. Fresh hosted CI is pending.

@github-actions github-actions Bot added bug Something isn't working area:shuffle Shuffle (JVM and native) area:aggregation Hash aggregates, aggregate expressions area:expressions Expression evaluation labels Sep 17, 2026
@sunchao sunchao added run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue correctness labels Sep 17, 2026
@andygrove

Copy link
Copy Markdown
Member

I started this review before realizing the PR was a draft. Feel free to ignore if this isn't helpful.

I checked this out and ran it locally, and I think the justification is stronger than what the description says. The description leads with decimal aggregate overflow, which is hard to pin down, because a deterministic hash still sends every row with a given key to one partition, so the per-group aggregation doesn't actually change. What does break is co-partitioning. applyCometShuffle decides native vs columnar per exchange with no harmonization across the plan, so the two sides of a join can end up with different shuffle implementations. The columnar path uses Spark's HashPartitioning.partitionIdExpression and the native path uses the Rust murmur3, and on a wide decimal key those disagree.

I built that plan on main at 9a4d5f283, joining a Parquet table to a JSON view on a DECIMAL(38,0) key:

CometSortMergeJoin [k#37], [k#40], Inner
:- CometSort [k#37, a#38], [k#37 ASC NULLS FIRST]
:  +- CometExchange hashpartitioning(k#37, 7), ENSURE_REQUIREMENTS, CometNativeShuffle
:     +- CometNativeScan parquet spark_catalog.default.wd_parquet[k#37,a#38]
+- CometSort [k#40, b#41], [k#40 ASC NULLS FIRST]
   +- CometColumnarExchange hashpartitioning(k#40, 7), ENSURE_REQUIREMENTS, CometColumnarShuffle
      +- FileScan json [k#40,b#41]

Comet returns zero rows. Spark returns four. Default spark.comet.shuffle.mode=auto, nothing exotic set. It reproduces with AQE on as well, as long as coalescing doesn't collapse the fixture into a single partition, which is what you'd get at any real data volume. That's probably why nobody has hit it. Your one line fixes it. Both sides move to the columnar path and the answer comes out right.

The encoding mismatch underneath is what you'd expect from reading the two implementations side by side. hash_array_decimal! hashes value.to_le_bytes(), a fixed sixteen-byte little-endian i128, while Spark hashes toJavaBigDecimal.unscaledValue.toByteArray(), which is minimal-length big-endian two's complement. Byte order and length both differ, and murmur3 mixes length into fmix, so every value diverges. For DECIMAL(38,0) at seed 42, unscaled 1 hashes to -386724586 in Spark and -680163996 in Comet, which is partition 4 against partition 6 out of 7.

I should own part of this. I closed #3079 as not a bug, on the grounds that we don't need to allocate the same partitions as Spark. That was wrong, and the join above is why. Our own columnar path uses Spark's partitioner, so the native path has to match Spark in order to stay consistent with its sibling. I'll reopen it.

Could the description lead with the join rather than the overflow argument? It's much easier for a reviewer to verify. And would you add it as a test? The six tests here are good, and I confirmed they're load-bearing, since reverting just the one line fails four of them. But they all assert plan shape or partition-id equality. None of them shows a wrong answer, which is the case that actually matters.

A few smaller things.

The new comment drops the link to #3079, so there's no longer a pointer from the code to the work that would lift the restriction. Could it reference #5994 instead, since that's the live issue for the Spark-compatible native encoding? Related, HashUtils.unsupportedReasonFor in serde/hash.scala already carries this exact rule for hash, xxhash64 and sha*, recursing through struct, array and map, for the same reason. Would a shared predicate be worth it, or at least a cross-reference in both comments? When #5994 lands, both places need to change together and it would be easy to miss one.

On that, the real fix looks small. Producing BigInteger.toByteArray() form from an i128 is stripping leading 0x00 or 0xff from the big-endian bytes while keeping one sign byte, maybe ten lines of Rust. It would remove this restriction, unblock hash and xxhash64 for wide decimals, and avoid the plan regression. I'm fine landing the guard first since correctness comes first, but it would be good to be explicit about the sequencing rather than leaving it implied.

The tuning.md paragraph reads well, but the contributor docs need the same treatment. native_shuffle.md enumerates what disqualifies a hash key and currently lists only collated strings, and the comparison table near the bottom says "Primitives only (Hash, Range)". jvm_shuffle.md has the mirror list. All three are now slightly wrong.

Two things tuning.md doesn't mention that I think users would want. In native mode the whole aggregate reverts to Spark, not just the shuffle, which the new test confirms by asserting two ObjectHashAggregateExec and zero CometHashAggregateExec. And under CometCelebornShuffleManager there's no columnar fallback at all, so the shuffle drops to Spark entirely. Is the Celeborn case worth covering in CometCelebornShufflePlanningSuite too?

Eleven TPC-DS plans move from CometExchange to CometColumnarExchange, which is a full columnar to row to columnar round trip. The new benchmark is the right instrument for measuring that, but the description says no timings were collected and the numbers quoted are from a different workload on a different revision. Could you run it on both revisions before taking this out of draft, along with a couple of the affected TPC-DS queries? Also, CometShuffleBenchmark already sweeps a type list containing DecimalType(10, 0) and has a nested-hash-key section. Would a DecimalType(38, 2) entry there cover this without a new class?

Last, the branch state. #5421 landed as 6065705c1, so eleven of the thirteen commits here are already on main. After a rebase there's a conflict in CometAggregateSuite.scala, and CometSparkSessionExtensionsSuite.scala won't compile because shouldOverrideMemoryConf was removed from main. The core is fine though. I ran the six new CometNativeShuffleSuite decimal tests against current main with the routing change applied and they all pass.

@andygrove

Copy link
Copy Markdown
Member

#5421 is merged now, so the prerequisite here is met. Could you merge main and drop the carried commits? That would leave just f509e5db8 plus the regression test, which makes the routing change reviewable on its own, and then this can come out of draft.

One thing worth adding to the rationale: this fixes more than the aggregate-overflow case. applyCometShuffle picks native or columnar per exchange with no harmonization across the plan, so the two sides of a join can end up with different shuffle implementations — native murmur3 on one side, Spark's partitionIdExpression on the other. On a DECIMAL(38,0) key those disagree and the join silently drops rows; I reproduced zero rows against Spark's four on main yesterday. Your guard covers that too, since once the native path declines the wide decimal both sides use Spark's partitioner and agree again. Details are in #3079.

@sunchao
sunchao force-pushed the dev/chao/codex/oss-wide-decimal-shuffle branch from b42be68 to d2231af Compare September 28, 2026 21:13
@sunchao
sunchao marked this pull request as ready for review September 28, 2026 21:15
@sunchao sunchao added run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue run-benchmark-check Run the benchmark compile and lint check on this pull request instead of waiting for the merge queue labels Sep 28, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Native and JVM shuffles could assign matching wide-decimal keys to different partitions, silently dropping join rows.
  • Design approach: Reject native multi-partition hashing when a key contains a decimal with precision greater than 18, using the existing fallback paths.
  • Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The precision boundary matches Spark’s hashing implementation. Recursive checks cover nested decimal leaves. The single-partition exception matches the writer’s existing SinglePartition serialization. Aggregate-buffer repair remains in place.
  • Key design decisions: The change preserves native handling of decimal payloads, range keys and single-partition exchanges. Extending the existing predicate keeps the implementation simple. Cross-references connect the shuffle and SQL hash restrictions without adding an unnecessary abstraction.
  • Implementation sketch: The routing guard is accompanied by partition-ID, mixed-input join, nested-key, aggregate-buffer and Celeborn regressions, updated fuzz expectations and plan snapshots, documentation, and a focused COUNT benchmark.
  • Behavioral changes worth calling out: auto uses JVM columnar shuffle when eligible. native and Celeborn fall back to Spark, potentially restoring incompatible aggregates too. The conversion overhead is documented. The author’s measurements report increases of 9–33% for native mode and 17–25% for auto on the measured workload. I did not independently rerun those timings.
  • Suggested improvements: No introduced P1/P2 issues found within this review. No substantiated unresolved P1/P2 concerns remain from the existing discussion.

Reviewed the full base-relative diff at 8b5915ae9878238654f76e38a2fe1d72e600af9a against b8a2358cc0a54818f5ac6db85ce7e5a3a77e9692, including all four PR commits. The apparent timezone-documentation and nested-null-test removals are base-only commits after the branch point, not deletions introduced by this PR. Read the supplied discussions before forming conclusions. The PR remains non-draft.

Routed skills: review-comet-pr, review-comet-shuffle-pr, and review-comet-expression-pr.

Exact-head CI: run 36487974865 is associated with the reviewed SHA and tests synthetic merge 71a8847. At the final check, 35 checks passed, 11 were skipped and 12 remained running, with no failures. Shuffle suites passed across all five Spark profiles. Inspected logs confirm the new regressions passed on Spark 3.4, 3.5 and 4.1. Native tests, benchmark compilation, TPC-H and TPC-DS checks also passed. Spark SQL and several broader suites remain pending.

Local validation: A standalone diagnostic using the extracted native hash function confirmed the encoding mismatch and precision boundary. For wide-decimal 1, Spark-compatible encoding selects partition 4 of 7 while native encoding selects 6. Whitespace checks passed. No full local Comet build or end-to-end suite was run because this checkout has no prepared JVM/native artifacts. The diagnostic does not replace integration testing.

@sunchao
sunchao requested a review from andygrove September 28, 2026 23:53

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Native and JVM shuffles assigned matching wide-decimal keys to different partitions, allowing mixed-input joins to silently lose rows.
  • Design approach: Reject native multi-partition hashing for decimal keys with precision greater than 18, including nested leaves, through the existing shuffle fallback paths.
  • Correctness / compatibility analysis: Verified the hashing boundary against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Recursive key checks cover nested decimals. The single-partition exception matches the writer’s SinglePartition serialization. Existing aggregate-buffer repair handles the newly triggered Spark fallback.
  • Key design decisions: Decimal payloads, range keys and single-partition exchanges retain native support. Extending the existing predicate keeps the change small. Cross-references connect the shuffle and SQL hash restrictions without adding an unnecessary abstraction.
  • Implementation sketch: Adds partition-ID, mixed-input join, nested-key, aggregate-buffer and Celeborn regressions, updates fuzz expectations and plan snapshots, documents fallback behavior, and adds a focused COUNT benchmark.
  • Behavioral changes worth calling out: auto selects JVM columnar shuffle when eligible. native and Celeborn use Spark fallback, potentially restoring incompatible aggregates too. Conversion overhead is documented. The author reports increases of 9–33% for native mode and 17–25% for auto in the measured workload. These timings were not independently rerun, and TPC-DS runtime impact remains unmeasured.
  • Suggested improvements: No introduced P1/P2 issues found within this review. No substantiated unresolved P1/P2 concerns remain from the existing discussion.

Reviewed full SHA 8b5915ae9878238654f76e38a2fe1d72e600af9a against base b8a2358cc0a54818f5ac6db85ce7e5a3a77e9692, covering all four PR commits and the full base-relative scope. The apparent timezone-documentation and nested-null-test removals correspond to base-only additions after the branch point. Read all supplied non-Copilot discussion and reviews. Live metadata confirms the PR remains non-draft at the requested head.

Routed skills: review-comet-pr, review-comet-shuffle-pr, and review-comet-expression-pr.

Exact-head CI: 48 checks passed, 13 skipped, none failed or pending. Run 36487974865 is associated with the reviewed SHA and tested synthetic merge 71a8847. Shuffle jobs passed across all five Spark profiles. Spark 4.1 SQL, native tests, benchmark compilation, and TPC-H/TPC-DS checks passed. Inspected logs confirm the new regressions executed successfully.

Local validation: Recompiled and ran a standalone diagnostic using the exact current native hash function. It confirmed precision-18 agreement and precision-19/38 encoding divergence, including positive, negative and boundary values. For wide-decimal 1, Spark-compatible encoding selects partition 4 of 7 while native encoding selects 6. Whitespace checks passed. No full local build or JVM integration suite was run because this checkout lacks prepared artifacts. The diagnostic supplements the CI evidence and does not replace integration testing.

msafonov added a commit to joomcode/datafusion-comet that referenced this pull request Oct 1, 2026
apache#6005)

Comet's native Murmur3 hashes a decimal with precision above 18 differently
from Spark (apache#3079, apache#5994). Comet picks the shuffle of
each exchange on its own, so the two inputs of a sort-merge join could be
partitioned by different hash functions and matching keys meet in different
partitions: rows were silently dropped (a decimal(38,10) join returned 108 of
2000 rows).

Native shuffle now refuses hash partitioning over more than one partition
whose keys contain such a decimal at any depth. In auto mode the shuffle
becomes Comet's columnar shuffle when that applies; in native mode, or with
Celeborn, a Spark shuffle. Wide decimals in the payload, range partitioning
and single-partition shuffles stay native. BoundaryFormats offers a native
shuffle only under the same condition, via
CometShuffleExchangeExec.hasWideDecimalHashKey.

Ports the PR's tests and approved TPC-DS plans, and adds a sort-merge join
of a native-scan and a Spark-scan input on decimal(38,10) with AQE and
boundary formats on and off.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 3, 2026
…VM shuffle

A typed Dataset conversion moves the shuffle above it from Comet's columnar
shuffle, which partitions with Spark's hash, to native shuffle. Native
shuffle hashes decimals wider than 18 digits differently from Spark
(apache#5994), so a join on such keys with an input that stays on columnar
shuffle, for example one with an array<int> column the conversion
declines, put matching keys in different partitions and returned 9 rows
instead of 100. Such a shuffle now stays on columnar shuffle unless it has
one partition, the same rule apache#6005 proposes for every native shuffle.

Also declare the benchmark's row count as a Long, which the strict
Spark 3.5 compile requires.
@sunchao
sunchao force-pushed the dev/chao/codex/oss-wide-decimal-shuffle branch from 8b5915a to f682d09 Compare October 4, 2026 17:40
@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Rebased onto main 9866221 and pushed f682d09. The carried prerequisite commits are already absent; this remains the four-commit routing fix. The description leads with the mixed Parquet/JSON join, and the regression checks all four matching rows with AQE on and off. Code/docs retain the #5994 cross-reference, native/auto/Celeborn fallback behavior, recursive wide-decimal boundary and the single-partition exception.

Current-head validation passed all 64 selected JDK 21 / Spark 4.1.3 tests: seven decimal tests and the full 57-test Celeborn planning suite. The exact-base Linux native artifact was verified and bundled first. Spotless and whitespace checks pass.

The separate COUNT benchmark is retained to isolate shuffle fallback from the existing AVG benchmark’s wide-decimal support constraints. Current-head benchmark timings and TPC-DS query runtimes were not rerun, so the description now states that limit and removes the old timing table. Fresh hosted CI is pending.

@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

The fresh Spark 3.5 strict-warning job caught an implicit Int-to-Long conversion in the new benchmark constructor. Fixed with rows.toLong in f79e2ce. The complete root Maven reactor now passes test-compile -Pspark-3.5 -Pstrict-warnings -DskipTests locally on JDK 21 / Scala 2.12. The earlier 64 decimal-shuffle/Celeborn runtime tests remain unchanged; the new hosted run is pending.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Native and JVM shuffles assigned matching wide-decimal keys to different partitions, allowing mixed-input joins to silently lose rows.
  • Design approach: Reject native multi-partition hashing for decimal keys with precision greater than 18, including nested leaves, through the existing fallback paths.
  • Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The precision boundary matches Spark’s decimal hashing. The single-partition exception matches the writer’s existing SinglePartition serialization. Aggregate-buffer repair preserves compatible execution when fallback occurs.
  • Key design decisions: Reusing the recursive type predicate and existing shuffle selection keeps the implementation small. Wide-decimal payloads, range keys and single-partition exchanges retain native eligibility.
  • Implementation sketch: Adds partition-ID, mixed-input join, nested-key, aggregate-buffer and Celeborn regressions. Updates fuzz expectations, eleven plan snapshots and documentation, and adds a focused COUNT benchmark.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, this intentionally changes affected shuffle routing. auto uses JVM columnar shuffle when eligible, adding conversion overhead. native and Celeborn use Spark fallback, which can also restore incompatible aggregates. The documentation explains this tradeoff.
  • Suggested improvements: No introduced P1/P2 issues found within this review. No substantiated unresolved P1/P2 concerns remain from the existing discussion.

Reviewed the full 21-file, five-commit diff at f79e2ced39fbe30930a70f8370b1e506ed13a873 against ac9ae94d057013095fdf1b70e2d4bc5b24a58d90. Read the supplied reviews and discussions. Live metadata confirms the PR remains non-draft at the requested head.

Routed skills: review-comet-pr, review-comet-expression-pr and review-comet-shuffle-pr.

Exact-head CI: 53 checks succeeded, 14 skipped, none failed or pending. Run 37222542918 is associated with the reviewed SHA and tested synthetic merge f20fa8b7e351bb71fde38f9d9231634f4e4e850b. Inspected logs confirm the new decimal and Celeborn regressions passed across all five Spark profiles. Spark 4.1 SQL, native tests, strict Scala warnings, benchmark compilation and TPC-H/TPC-DS checks passed.

Local validation: Compiled and ran an isolated diagnostic using the current native hash function across 36 boundary values. It confirmed precision-18 agreement and precision-19/38 encoding divergence. Wide-decimal 1 selects partition 4 of 7 with Spark’s encoding versus 6 natively. Whitespace checks passed. No local full build or JVM integration suite was run because this checkout lacks prepared artifacts. Current-head benchmark timings and affected TPC-DS runtime costs remain unmeasured.

comphead pushed a commit to comphead/arrow-datafusion-comet that referenced this pull request Oct 5, 2026
…che#6564)

* feat: run the operators above a typed Dataset operation natively

Typed Dataset operations (map, flatMap, mapPartitions, mapGroups, cogroup)
pass JVM objects between their operators, so they stay on Spark, and today
the operators above them stay on Spark too until the next shuffle. Every
typed operation ends in SerializeFromObjectExec, whose output is ordinary
rows. With the new spark.comet.convert.typedDataset.enabled, CometExecRule
puts a CometSparkToColumnarExec above it, so a partial aggregate, a
broadcast join or a native shuffle above the operation runs natively.

Spark inserts no columnar transitions below a RowToColumnarTransition, so
the rule adds them to the subtree under the conversion with Spark's own
ApplyColumnarRulesAndInsertTransitions. Without them the typed operation
would read its Comet child through Spark's interpreted columnar-to-row
path. CometExecRule no longer tags a ColumnarToRowTransition as an
unsupported operator, which it did to these transitions on its second
pass under AQE.

Off by default: it is 1.4-2.2x faster when an aggregate over many groups
sits above the typed operation, and slower when the work above is cheap.

Adds CometTypedDatasetSuite and CometTypedDatasetBenchmark.

* fix: keep wide-decimal shuffles above a typed Dataset conversion on JVM shuffle

A typed Dataset conversion moves the shuffle above it from Comet's columnar
shuffle, which partitions with Spark's hash, to native shuffle. Native
shuffle hashes decimals wider than 18 digits differently from Spark
(apache#5994), so a join on such keys with an input that stays on columnar
shuffle, for example one with an array<int> column the conversion
declines, put matching keys in different partitions and returned 9 rows
instead of 100. Such a shuffle now stays on columnar shuffle unless it has
one partition, the same rule apache#6005 proposes for every native shuffle.

Also declare the benchmark's row count as a Long, which the strict
Spark 3.5 compile requires.

* refactor: simplify the typed Dataset conversion's shuffle guard and tests

- Use Spark's existsRecursively(DecimalType.isByteArrayDecimalType) for
  the wide-decimal check, fold the duplicated conversion cases in
  readsTypedDatasetConversion, and check the config and earlier native
  shuffle reasons before walking the plan. Add a TODO to drop the guard
  once apache#5994 lands.
- Tests: assert the join's exchanges with checkCometExchange, drop
  conditions that cannot change the transition assertion, write the
  Parquet table once for both AQE settings, and take checkConverted's
  query by value.
- Benchmark: import spark.implicits instead of declaring encoders.

* fix: preserve typed Dataset limit short-circuiting

* fix: leave typed Dataset output unconverted below any reader that stops early

The limit guard covered only physical limits, so a mapPartitions function
such as `_.take(1)`, or code reading Dataset.rdd, still read typed rows
through an Arrow batch and ran the user function on rows that Spark never
reaches. Walk down from each reader that can stop early (a limit, a top-k
over sorted input, a MapPartitionsExec, and a DeserializeToObjectExec at the
plan root) and stop at an operator that reads all of its input first, so a
limit above an aggregate keeps the conversion.

* fix: recognize Dataset.rdd plans whose deserializer Spark removed

The guard for code reading Dataset.rdd looked for the DeserializeToObjectExec
that Dataset.rdd puts at the root of the plan. When the Dataset ends in a
typed operation such as map, Spark's EliminateSerialization drops that
deserializer together with the operation's serializer, so the root is the
operation itself, or a typed filter over it. The output of an earlier typed
operation was then still converted, and map(f).filter(...).map(g).rdd.take(1)
ran f on rows that Spark never reaches. The guard now takes any root that
produces objects, under a filter or a project, as code reading Dataset.rdd.
A Dataset's own plan ends in rows, so no other plan has such a root.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:aggregation Hash aggregates, aggregate expressions area:expressions Expression evaluation area:shuffle Shuffle (JVM and native) bug Something isn't working correctness run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue run-benchmark-check Run the benchmark compile and lint check on this pull request instead of waiting for the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants