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]
