andygrove opened a new pull request, #6726: URL: https://github.com/apache/datafusion-comet/pull/6726
## Which issue does this PR close? Closes #6722. ## Rationale for this change `CometSparkToColumnarExec.isTypeSupported` overrode the shared `DataTypeSupport` rule to decline every array except `ARRAY<STRING>` (#5954) and every map except `MAP<STRING,STRING>` (#6036). Every conversion checks it: the `spark.comet.convert.*` sources, typed Dataset output (#6564), and the shuffle-input conversion in #6607 once that merges. A column of any other array or map type left the operators above it on Spark. The override was not a bug fix. The operator never admitted other arrays or maps, and #1741 carried the rule over when it refactored `DataTypeSupport`. `CometLocalTableScanExec` has no such override and already writes arbitrary arrays and maps through the same `RowArrowReader` and `ArrowWriter`. ## What changes are included in this PR? - Delete the override, so the shared rule decides. It checks element, key and value types recursively, rejects collated strings at every level, and rejects structs with duplicate field names or no fields. - Replace the `CometExecSuite` test that pinned the gate with one for the types that are now admitted and the ones that stay declined. A declined collection must record a fallback reason. The old override declined without one, so the reason tells the shared rule's recursion from a blanket decline. - Add a `CometExecSuite` test for each of 17 array and map shapes. It reads the column from an RDD and from Parquet through both readers, and runs a native filter, project, `size`, element access and native shuffle over it. - `CometTypedDatasetSuite`: two tests used `array<int>` as the column that does not convert. They use an interval column (`java.time.Duration`) now, which the conversion and Comet's columnar shuffle both decline. So in the wide-decimal join test the right input shuffles through Spark's exchange rather than Comet's columnar shuffle, and the test expects one Comet exchange instead of two. - `CometNativeShuffleSuite`: the duplicate-field-name test checked structs nested in arrays and maps for `CometLocalTableScanExec` only, because the blanket decline would have failed its reason assertion. It covers `CometSparkToColumnarExec` too now. - Docs: `datasources.md`, `in-memory-cache.md` and the comment in `CometInMemoryCacheBenchmark` described the old gate. The issue names the last two. `datasources.md` listed `ARRAY<STRING>` and `MAP<STRING,STRING>` as the only supported collections. The benchmark itself is unchanged. Adding array and map columns to it would mean re-measuring the tables in `in-memory-cache.md`. Row input writes nested values one element at a time (#6721), so these columns convert more slowly than flat ones, but the operators above them can now run natively. ## How are these changes tested? The 128 probe queries from the #6036 review ran before #6566 changed how nested columns are written, so the shapes are re-run as the new tests above. All of them pass. - With the override restored, 19 tests fail: the gate test, the 17 shape tests and the extended duplicate-field-name test. The two typed Dataset tests pass either way, since an interval stays unsupported. - Spark 4.1: `CometExecSuite` in full, the in-memory cache, native shuffle, planner (`CometExecRuleSuite`, `RevertNativeForTransitionHeavyStagesSuite`, `CometScanRuleSuite`) and task metrics suites, `CometArrowStreamSuite`, `CometArrowWriterSuite`, and the Iceberg write, Parquet writer and native reader suites. - Spark 3.4, 3.5, 4.0 and 4.2: the `SparkToColumnar` tests in `CometExecSuite`, `CometTypedDatasetSuite`, and the duplicate-field-name test. - A throwaway probe, not included, sent the same 17 shapes through seven more consumers: element access, a broadcast build side, union, sort, top-k, group-by and coalesce. Every plan was fully native and matched Spark. - Scalafix (`CHECK`, Spark 3.5 and Scala 2.12), spotless and scalastyle pass. The Spark SQL tests were not run. The diffs under `dev/diffs` set no `spark.comet.convert.*` key, and the one conversion that is on by default, `OneRowRelation`, has no columns, so they never reach the changed check. -- 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]
