viirya commented on code in PR #6726:
URL: https://github.com/apache/datafusion-comet/pull/6726#discussion_r4238895815
##########
spark/src/main/scala/org/apache/spark/sql/comet/CometSparkToColumnarExec.scala:
##########
@@ -137,14 +135,4 @@ object CometSparkToColumnarExec extends
CometSink[SparkPlan] with DataTypeSuppor
op: SparkPlan): CometNativeExec = {
CometScanWrapper(nativeOp, CometSparkToColumnarExec(op))
}
-
- override def isTypeSupported(
- dt: DataType,
- name: String,
- fallbackReasons: ListBuffer[String]): Boolean = dt match {
- case ArrayType(StringType, _) => true
- case MapType(StringType, StringType, _) => true
- case _: ArrayType | _: MapType => false
Review Comment:
With this case gone, the shared rule admits `CalendarIntervalType` inside
arrays and maps too. `IntervalMonthDayNanoWriter` multiplies microseconds by
1000 with `multiplyExact`. So an interval whose time part is longer than about
292 years fails the task with `ArithmeticException: long overflow`, where Spark
returns it. An RDD with an `ARRAY<INTERVAL>` column holding `new
CalendarInterval(1, 2, Long.MaxValue / 10)` reproduces it under
`spark.comet.convert.rdd.enabled=true`, and so does a map value. Before this PR
those columns fell back.
Top-level intervals already had this (#5279), and #6607 rules intervals out
of the shuffle-input conversion for the same reason. Could this sink decline
`CalendarIntervalType` with a reason pointing at #5279, at least inside arrays
and maps?
##########
docs/source/user-guide/latest/datasources.md:
##########
@@ -73,12 +73,13 @@ operator in
`spark.comet.sparkToColumnar.supportedOperatorList` by its Spark cla
### Spark-to-Comet conversion types
-Spark-to-Comet conversion supports `ARRAY<STRING>` and `MAP<STRING,STRING>`
with binary
-string semantics, both as top-level fields and inside supported structs.
Arrays and maps may
-be null; array elements and map values may also be null. Map keys must be
non-null.
-This applies to Spark row and columnar inputs when conversion is enabled for
the source.
-Other array element types, other map key/value types, nested collections, and
non-binary
-string collations remain unsupported at this conversion boundary. Source
defaults are unchanged.
+Spark-to-Comet conversion supports booleans, integers, floats, decimals,
strings, binary, dates,
+timestamps and calendar intervals, and arrays, maps and structs of them,
nested to any depth.
Review Comment:
This lists calendar intervals among the supported types, including inside
arrays and maps. Because of the overflow in #5279, could this either drop them
or say that an interval longer than about 292 years in its time part fails the
query?
##########
spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala:
##########
@@ -4150,43 +4150,190 @@ class CometExecSuite extends CometTestBase {
})
}
- test("SparkToColumnar admits only binary string arrays and maps through the
collection gate") {
- for (nullable <- Seq(false, true); containsNull <- Seq(false, true)) {
- for (dataType <- Seq(
- ArrayType(StringType, containsNull),
- MapType(StringType, StringType, containsNull))) {
- val field = StructField("tags", dataType, nullable)
- Seq(
- StructType(Seq(field)),
- StructType(Seq(StructField("nested", StructType(Seq(field))))))
- .foreach { schema =>
- assert(CometSparkToColumnarExec.isSchemaSupported(schema,
ListBuffer.empty))
- }
- }
- }
- val unsupported = Seq(
+ test("SparkToColumnar admits arrays and maps of every type the shared rule
supports") {
+ val struct = StructType(Seq(StructField("i", IntegerType),
StructField("s", StringType)))
+ val admitted = Seq(
+ ArrayType(StringType),
ArrayType(IntegerType),
+ ArrayType(DecimalType(38, 10)),
+ ArrayType(TimestampType),
ArrayType(BinaryType),
ArrayType(ArrayType(StringType)),
- ArrayType(StructType(Seq(StructField("s", StringType)))),
- MapType(StringType, IntegerType),
- MapType(IntegerType, StringType),
- MapType(StringType, ArrayType(StringType)),
- MapType(StringType, StructType(Seq(StructField("s", StringType))))) ++
- (if (isSpark40Plus)
+ ArrayType(struct),
+ ArrayType(MapType(StringType, IntegerType)),
+ MapType(StringType, StringType),
+ MapType(IntegerType, IntegerType),
+ MapType(DateType, DecimalType(20, 2)),
+ MapType(StringType, ArrayType(IntegerType)),
+ MapType(IntegerType, struct),
+ MapType(struct, StringType),
+ MapType(ArrayType(IntegerType), IntegerType),
+ MapType(StringType, MapType(IntegerType, StringType)))
+ for (nullable <- Seq(false, true); containsNull <- Seq(false, true);
dataType <- admitted) {
+ val collection = dataType match {
+ case a: ArrayType => a.copy(containsNull = containsNull)
+ case m: MapType => m.copy(valueContainsNull = containsNull)
+ }
+ val field = StructField("c", collection, nullable)
+ Seq(StructType(Seq(field)), StructType(Seq(StructField("nested",
StructType(Seq(field))))))
+ .foreach { schema =>
+ val reasons = ListBuffer.empty[String]
+ assert(CometSparkToColumnarExec.isSchemaSupported(schema, reasons),
s"$collection")
+ assert(reasons.isEmpty, reasons)
+ }
+ }
+
+ // The shared rule looks inside a collection and declines these for what
it finds there. A
+ // decline of the whole collection would record no reason, which is how
the two differ.
+ val duplicate = StructType(Seq(StructField("a", LongType),
StructField("a", LongType)))
+ val declined = Seq(
Review Comment:
Could one calendar-interval collection go in this list, such as
`ArrayType(CalendarIntervalType)`, once the sink declines it? Then the test
pins that decision the way it pins the other interval types.
##########
spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala:
##########
@@ -4150,43 +4150,190 @@ class CometExecSuite extends CometTestBase {
})
}
- test("SparkToColumnar admits only binary string arrays and maps through the
collection gate") {
- for (nullable <- Seq(false, true); containsNull <- Seq(false, true)) {
- for (dataType <- Seq(
- ArrayType(StringType, containsNull),
- MapType(StringType, StringType, containsNull))) {
- val field = StructField("tags", dataType, nullable)
- Seq(
- StructType(Seq(field)),
- StructType(Seq(StructField("nested", StructType(Seq(field))))))
- .foreach { schema =>
- assert(CometSparkToColumnarExec.isSchemaSupported(schema,
ListBuffer.empty))
- }
- }
- }
- val unsupported = Seq(
+ test("SparkToColumnar admits arrays and maps of every type the shared rule
supports") {
+ val struct = StructType(Seq(StructField("i", IntegerType),
StructField("s", StringType)))
+ val admitted = Seq(
+ ArrayType(StringType),
ArrayType(IntegerType),
+ ArrayType(DecimalType(38, 10)),
+ ArrayType(TimestampType),
ArrayType(BinaryType),
ArrayType(ArrayType(StringType)),
- ArrayType(StructType(Seq(StructField("s", StringType)))),
- MapType(StringType, IntegerType),
- MapType(IntegerType, StringType),
- MapType(StringType, ArrayType(StringType)),
- MapType(StringType, StructType(Seq(StructField("s", StringType))))) ++
- (if (isSpark40Plus)
+ ArrayType(struct),
+ ArrayType(MapType(StringType, IntegerType)),
+ MapType(StringType, StringType),
+ MapType(IntegerType, IntegerType),
+ MapType(DateType, DecimalType(20, 2)),
+ MapType(StringType, ArrayType(IntegerType)),
+ MapType(IntegerType, struct),
+ MapType(struct, StringType),
+ MapType(ArrayType(IntegerType), IntegerType),
+ MapType(StringType, MapType(IntegerType, StringType)))
+ for (nullable <- Seq(false, true); containsNull <- Seq(false, true);
dataType <- admitted) {
+ val collection = dataType match {
+ case a: ArrayType => a.copy(containsNull = containsNull)
+ case m: MapType => m.copy(valueContainsNull = containsNull)
+ }
+ val field = StructField("c", collection, nullable)
+ Seq(StructType(Seq(field)), StructType(Seq(StructField("nested",
StructType(Seq(field))))))
+ .foreach { schema =>
+ val reasons = ListBuffer.empty[String]
+ assert(CometSparkToColumnarExec.isSchemaSupported(schema, reasons),
s"$collection")
+ assert(reasons.isEmpty, reasons)
+ }
+ }
+
+ // The shared rule looks inside a collection and declines these for what
it finds there. A
+ // decline of the whole collection would record no reason, which is how
the two differ.
+ val duplicate = StructType(Seq(StructField("a", LongType),
StructField("a", LongType)))
+ val declined = Seq(
+ ArrayType(duplicate),
+ MapType(LongType, duplicate),
+ ArrayType(ArrayType(duplicate)),
+ ArrayType(NullType),
+ MapType(IntegerType, ArrayType(NullType)),
+ ArrayType(DayTimeIntervalType()),
+ MapType(StringType, YearMonthIntervalType())) ++
+ (if (isSpark40Plus) {
Seq(
- DataType.fromDDL("ARRAY<STRING COLLATE UTF8_LCASE>"),
- DataType.fromDDL("MAP<STRING COLLATE UTF8_LCASE,STRING>"),
- DataType.fromDDL("MAP<STRING,STRING COLLATE UTF8_LCASE>"))
- else Seq.empty)
- unsupported.foreach { dataType =>
+ "ARRAY<STRING COLLATE UTF8_LCASE>",
+ "ARRAY<ARRAY<STRING COLLATE UTF8_LCASE>>",
+ "ARRAY<STRUCT<s: STRING COLLATE UTF8_LCASE>>",
+ "MAP<STRING COLLATE UTF8_LCASE,STRING>",
+ "MAP<STRING,STRING COLLATE UTF8_LCASE>",
+ "MAP<INT,ARRAY<STRING COLLATE UTF8_LCASE>>").map(DataType.fromDDL)
+ } else Seq.empty)
+ declined.foreach { dataType =>
val schema = StructType(Seq(StructField("value", dataType)))
- assert(!CometSparkToColumnarExec.isSchemaSupported(schema,
ListBuffer.empty), dataType)
- assert(
- !CometSparkToColumnarExec.isSchemaSupported(
- StructType(Seq(StructField("nested", schema))),
- ListBuffer.empty),
- dataType)
+ Seq(schema, StructType(Seq(StructField("nested", schema)))).foreach {
declinedSchema =>
+ val reasons = ListBuffer.empty[String]
+ assert(!CometSparkToColumnarExec.isSchemaSupported(declinedSchema,
reasons), dataType)
+ assert(reasons.nonEmpty, s"$dataType was declined without a reason")
+ }
+ }
+ }
+
+ // A column of each array and map shape the conversion admits, as its type,
an expression over
+ // `id` that builds it, and one that reads an element of it. The column has,
in some rows, a
+ // null collection, an empty one, and elements that are null.
+ private val sparkToColumnarShapes: Seq[(String, String, String)] = {
Review Comment:
Every shape here builds its column with `cast(... AS ...)`, so array
elements and map values are always nullable. `CometLocalTableScanExec` widens
child nullability because non-null children crashed `slice` and `map_entries`
(#4789), and this sink doesn't widen.
I checked `slice`, `map_entries`, `array_insert` and `array_prepend` over
`containsNull = false` and `valueContainsNull = false` input from an RDD, a
typed `Seq[Int]`/`Map[Int, Int]` Dataset, and Parquet. They all match Spark
today. Could you add one shape whose schema has non-null children, for example
a typed `Seq[Int]` field, so this keeps working?
--
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]