andygrove commented on PR #6005:
URL:
https://github.com/apache/datafusion-comet/pull/6005#issuecomment-5780680447
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.
--
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]