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]

Reply via email to