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]

Reply via email to