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]

Reply via email to