andygrove commented on code in PR #6110:
URL: https://github.com/apache/datafusion-comet/pull/6110#discussion_r4115837660
##########
docs/source/contributor-guide/jvm_shuffle.md:
##########
@@ -54,6 +54,20 @@ JVM shuffle (`CometColumnarExchange`) is used instead of
native shuffle (`CometE
[Supported partition key
types](native_shuffle.md#when-native-shuffle-is-used) for the exact
rules. Complex types are fully supported as data columns in both
implementations.
+4. **Partition keys native shuffle cannot serialize**: native shuffle
serializes the
+ partitioning expressions to protobuf, so a key expression Comet has no
serde for, or whose
+ serde reports it incompatible, keeps the exchange off the native path. JVM
shuffle has no
+ such requirement, because it evaluates the key on the JVM through
`UnsafeProjection` and
+ `LazilyGeneratedOrdering`, so these exchanges land here rather than on
Spark's shuffle. One
+ example is the `mapsort(...)` wrapper Spark 4.0 and later adds around a map
used as a
+ shuffle key: Comet cannot serialize it for array or struct map keys, so
such an exchange
+ becomes `CometColumnarExchange`.
Review Comment:
Thanks for the doc updates. One gap I noticed. This new case 4 says native
shuffle declines a key expression Comet can't serialize. The "When Native
Shuffle is Used" list in `native_shuffle.md` says native is chosen "when all of
the following conditions are met", but it has no such condition. Could you add
a fifth item there saying every hash partitioning expression and range sort
order has to convert through `exprToProto`, with the same `mapsort` example?
Then the two docs agree on when native is picked.
--
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]