andygrove commented on code in PR #5889: URL: https://github.com/apache/datafusion-comet/pull/5889#discussion_r4185885090
########## spark/src/test/scala/org/apache/comet/serde/CometScalarSubquerySuite.scala: ########## @@ -0,0 +1,204 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.serde + +import org.apache.spark.sql.CometTestBase +import org.apache.spark.sql.catalyst.expressions.{Alias, Literal, NamedExpression} +import org.apache.spark.sql.execution.{ProjectExec, ScalarSubquery, SubqueryExec} +import org.apache.spark.sql.types._ + +import org.apache.comet.CometSparkSessionExtensions.{isSpark40Plus, isSpark41Plus} +import org.apache.comet.serde.QueryPlanSerde.supportedDataType +import org.apache.comet.shims.CometTypeShim + +/** Direct type-gate tests avoid optimizer folding and exercise types Parquet cannot store. */ +class CometScalarSubquerySuite extends CometTestBase with CometTypeShim { + + private def struct(dt: DataType): StructType = StructType(Seq(StructField("value", dt))) + + private lazy val emptyInput = spark.range(0).queryExecution.sparkPlan + + private def subquery(dt: DataType): ScalarSubquery = { + // Inspect the declared result type without executing or optimizing the typed NULL away. + val plan = ProjectExec(Seq(Alias(Literal.create(null, dt), "result")()), emptyInput) + ScalarSubquery(SubqueryExec("type-check", plan), NamedExpression.newExprId) + } + + private def supported(dt: DataType): Boolean = + CometScalarSubquery.getSupportLevel(subquery(dt)) == Compatible() + + private def versionSpecificTypes: Seq[DataType] = { + val strings = + if (isSpark40Plus) Seq(DataType.fromDDL("STRING COLLATE UTF8_LCASE")) else Seq.empty + val times = if (isSpark41Plus) Seq(DataType.fromDDL("TIME")) else Seq.empty + strings ++ times ++ variantType.toSeq + } + + private val scalarTypes: Seq[DataType] = Seq( + BooleanType, + ByteType, + ShortType, + IntegerType, + LongType, + FloatType, + DoubleType, + StringType, + BinaryType, + DateType, + TimestampType, + TimestampNTZType, + NullType, + DecimalType(1, 0), + DecimalType(10, 2), + DecimalType(38, 38)) + + // Freeze the pre-refactor shared predicate: existing callers must retain their accepted types. + private def legacySupported(dt: DataType, allowComplex: Boolean): Boolean = dt match { Review Comment: `legacySupported` and the first two tests check `QueryPlanSerde.supportedDataType` itself against a copy of the predicate from before #5025. Now that #5025 has merged with its own boundary test in `QueryPlanSerdeSuite`, these duplicate that coverage. The frozen copy also means a legitimate later change to the shared helper would fail a suite named for scalar subqueries. Could we drop them and keep the scalar-subquery checks, with the non-struct case comparing `supported(dt)` to `supportedDataType(dt)` directly? ########## docs/source/user-guide/latest/expressions.md: ########## @@ -705,6 +705,8 @@ Comet also accelerates a number of Catalyst expressions that have no Spark SQL f This list is illustrative, not exhaustive: the per-function tables are not the complete set of expressions Comet can accelerate. +Scalar subqueries can return structs, including those created when Spark merges multiple scalar subqueries. Struct results are transferred from Spark through Arrow IPC during native physical planning and retained as owned, immutable literals for that plan's execution. Supported fields include booleans, numeric types, default-collation strings, binary, dates, timestamps, nulls, and nested structs. Decimal fields require a non-negative scale no greater than their precision. Structs must be non-empty and have distinct field names at each level; arrays, maps, intervals, and other unsupported field types still cause fallback to Spark. Existing non-struct scalar-subquery paths are unchanged. Review Comment: Most of this paragraph describes the implementation rather than what a user sees. The Arrow IPC and owned-literal details don't give a user anything to act on, and "Existing non-struct scalar-subquery paths are unchanged" only makes sense next to this diff. One thing a user does need is missing. A string field that holds invalid UTF-8 now comes back with U+FFFD replacement where these queries used to fall back, so `hex(cast(s.v AS binary))` returns `EFBFBD28` where Spark returns `C328`. Could we trim this to what is supported and what falls back, and point to the "Strings with non-UTF-8 bytes" section of the compatibility guide? The list of boundaries in that section could name scalar subquery results too. -- 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]
