andygrove commented on PR #5538:
URL: 
https://github.com/apache/datafusion-comet/pull/5538#issuecomment-5876671147

   This is a light fully automated review since there are so many PRs open.
   
   1. The direct path skips the invalid UTF-8 decoding that the FFI path does. 
`ScanExec::import_column` (`native/core/src/execution/operators/scan.rs:173`) 
runs every JVM-imported column through `decode_string_arrays`, while 
`BlockScanExec::get_next` 
(`native/core/src/execution/operators/shuffle_scan.rs:250-283`) decodes with 
`skip_validation` and passes the columns straight on. Shuffle blocks are 
encoded natively from strings that were already decoded, but broadcast blocks 
are written by Arrow Java from whatever the exchange's child produced. With 
`spark.comet.exec.localTableScan.enabled=true`, a broadcast join against 
`VALUES (1, CAST(X'80' AS STRING)) AS v(k, s)` serializes the 
`CometLocalTableScanExec` output directly, so the raw `0x80` byte reaches 
native unchecked. Tracing it through, `length(v.s)` above the join goes to 
`SparkLengthFunc`'s `chars().count()` and returns 0, where Spark and the legacy 
path both return 1. The in-memory cache and `sparkToColumnar` leaves build the
 ir strings with the same `ArrowWriter`. Could the broadcast kind run 
`decode_string_arrays` on its columns, with a test along those lines?
   
   2. 
`spark/src/test/scala/org/apache/spark/sql/benchmark/CometBroadcastDirectReadBenchmark.scala:65-72`
 builds both tables from `spark.range`. `RangeExec` only becomes Comet through 
`spark.comet.sparkToColumnar.enabled`, which defaults to `false` and isn't set 
here, and `CometExecRule` only converts the projection and the 
`BroadcastExchangeExec` above it when their children are native. So I think 
both cases time the same all-Spark plan. Would it make sense to write the 
tables to Parquet first, the way `CometExecBenchmark` uses `prepareTable`, and 
check that the plan contains a `BroadcastScan` before posting numbers? Once 
this is merged with `main`, the `getSparkSession` override will also need 
`spark.comet.exec.onHeap.enabled=true` (or off-heap memory), since #6195 stops 
Comet from loading at all without one of them.
   
   3. `docs/source/contributor-guide/native_shuffle.md:236-259` still names 
`findShuffleScanIndices`, `shuffleScanIndices` and `ShuffleScanExec`, and says 
every other slot gets an `ArrowArrayStream`. 
`docs/source/contributor-guide/ffi.md:105-109` also describes a `ShuffleScan` 
slot as the only one that skips the FFI boundary. Could this PR update both for 
the broadcast path? It would also help if the new config's doc at 
`spark/src/main/scala/org/apache/comet/CometConf.scala:232` said that the block 
codec comes from `spark.comet.shuffle.compression.codec` and 
`spark.shuffle.compress`.
   


-- 
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