andygrove opened a new issue, #6722:
URL: https://github.com/apache/datafusion-comet/issues/6722
### What is the problem the feature request solves?
`CometSparkToColumnarExec.isTypeSupported` overrides the shared
`DataTypeSupport` rule and declines every array except `ARRAY<STRING>` (#5954)
and every map except `MAP<STRING,STRING>` (#6036):
```scala
case ArrayType(StringType, _) => true
case MapType(StringType, StringType, _) => true
case _: ArrayType | _: MapType => false
case _ => super.isTypeSupported(dt, name, fallbackReasons)
```
Every `CometSparkToColumnarExec` conversion checks it: the
`spark.comet.convert.*` sources (Parquet, JSON, CSV, Range, Spark's cache, RDDs
and row data sources) and typed Dataset output (#6564). When one of them
produces any other array or map column, nothing converts it, and the operators
above stay on Spark. #6607 gates its shuffle-input conversion on the same
check, so once it merges, row-based shuffles carrying these columns would stay
on the JVM columnar shuffle too.
No bug fix introduced the rule: the operator has 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`.
### Describe the potential solution
Delete the override and let the shared rule decide. It already checks
element, key and value types recursively, rejects collated strings at every
level, and rejects structs with duplicate field names.
In the #6036 review, deleting the override passed 128 of 128 probe queries
on Spark 4.1. They covered 16 array and map shapes, read from an RDD and
through Parquet's row and vectorized readers, with native filters, projections,
element access and a native shuffle above. The shapes were maps with int, date,
`decimal(20,2)`, array and struct keys or values, and arrays of int,
`decimal(38,10)`, timestamp, binary, double, boolean, arrays and structs. The
#4789 non-null child shapes passed too (26 of 26). Those probes ran before
#6566 changed how nested columns are written from columnar input, so they need
to run again.
Also:
- Replace the `CometExecSuite` test that pins the gate ("SparkToColumnar
admits only binary string arrays and maps through the collection gate") with
one for the newly admitted types and the collated strings that stay declined.
- Fix the two places that describe the gate:
`docs/source/user-guide/latest/in-memory-cache.md` and the comment above the
struct columns in `CometInMemoryCacheBenchmark`. Both were already wrong for
`ARRAY<STRING>` and `MAP<STRING,STRING>`. The benchmark can then measure array
and map columns.
### Additional context
Part of #6565. 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 get to run natively.
--
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]