Audit notes for expressions in this category that have been audited. Absence of an entry means the expression has not been audited yet, not that it is unsupported. See the user guide Spark Expression Support for current support status.
CreateArray(children, useStringTypeWhenEmpty); element type is the common type of children. Comet routes via CometCreateArray (native make_array) and special-cases the empty-array case to dodge a known DataFusion coerce_types issue (#3338).contextIndependentFoldable override; runtime semantics unchanged.BinaryExpression, evaluated directly. Comet routes via CometArrayAppend.RuntimeReplaceable and rewritten to ArrayInsert(arr, Literal(-1), elem). CometArrayAppend is therefore unreachable; dispatch goes through CometArrayInsert (which carries its own Incompatible notes documented at the array_insert entry).RuntimeReplaceable -> ArrayFilter(arr, IsNotNull(lambda)). Comet receives the rewritten form, dispatches through CometArrayFilter, which emits a call to DataFusion's built-in array_compact (from datafusion-functions-nested) via CometScalarFunction("array_compact").KnownNotContainsNull(...). The 4.x Spark4xCometExprShim strips the wrapper and emits the same array_compact call.ArrayContains(left, right) extends BinaryExpression with NullIntolerant with Predicate; inputTypes uses findWiderTypeWithoutStringPromotionForTwo. Wired as CometScalarFunction("array_contains").NullIntolerant trait replaced by nullIntolerant: Boolean; checkInputDataTypes adopts DataTypeUtils.sameType (collation-aware in 4.x).SQLOrderingUtil.ArrayDistinct(child) over ArraySetLike; uses SQLOpenHashSet so NaN and +0.0/-0.0 are canonicalized. Wired as CometScalarFunction("array_distinct").NullIntolerant -> nullIntolerant field refactor.SQLOpenHashSet.ArrayExcept(left, right) extends ArrayBinaryLike with ComplexTypeMergingExpression; result preserves left-side first occurrences not present in right. Comet routes via CometArrayExcept and unconditionally flags Incompatible (“Null handling and ordering may differ from Spark”); also falls back for BinaryType / StructType element types.nullIntolerant = true moves into ArrayBinaryLike; the overflow path uses arrayFunctionWithElementsExceedLimitError.array_distinct. (array_except still falls back by default for the null-handling/ordering reasons noted above.)CometArrayIntersect reports Incompatible because native element ordering can differ from Spark when the right array is longer than the left (DataFusion probes the longer side). By default it runs through the codegen dispatcher (Spark-correct) and uses the native path only when incompatible expressions are explicitly allowed. Non-default string collations are reported Unsupported and fall back to Spark.ArrayJoin(array, delimiter, nullReplacement). Comet routes via CometArrayJoin to DataFusion's array_to_string.inputTypes widened to AbstractArrayType(StringTypeWithCollation(supportsTrimCollation = true)); non-binary collations not propagated (#2190).contextIndependentFoldable override; runtime unchanged.CometArrayJoin reports Compatible when the delimiter and null replacement are literals or column reads; Spark short-circuits past those arguments and DataFusion does not, so anything else runs through the codegen dispatcher, as do non-default string collations (#2190). A nullable replacement is wrapped in an IsNull guard, since array_to_string reads a null null_string as “omit nulls” (#3178).ArrayMax(child) extends UnaryExpression with ImplicitCastInputTypes; skips NULL elements; for float/double Spark's SQLOrderingUtil treats NaN as greater than any non-NaN. Wired as CometScalarFunction("array_max").NullIntolerant -> nullIntolerant field refactor.ArrayMax with evalInternal returning the minimum. Same NULL-skip and NaN-ordering semantics. Wired as CometScalarFunction("array_min").array_max.array_max.ArrayPosition(left, right); returns 1-based LongType position, 0 if not found, NULL if either input is NULL. CometArrayPosition falls back for all-foldable args (constant folding handles those) and for unsupported element types (binary/struct/map/null).NullIntolerant -> nullIntolerant field refactor.array_prepend does not exist (added in Spark 3.5.0). The SQL file test carries MinSparkVersion: 3.5 so it is skipped here.ArrayPrepend(left, right) extends RuntimeReplaceable; replacement = new ArrayInsert(left, Literal(1), right). Comet never sees ArrayPrepend; dispatch goes through CometArrayInsert (which carries its own notes documented at the array_insert entry). NULL array yields NULL, NULL element is prepended. Type coercion casts the array to the tightest common type of element and array element (e.g. array_prepend(array(1, 2), 1.23D) -> [1.23, 1.0, 2.0]).ArrayPrepend now extends ArrayPendBase but the replacement is unchanged (ArrayInsert(left, Literal(1), right)).ArrayRemove(left, right); removes all occurrences equal to right. Wired as CometScalarFunction("array_remove"). Falls back via ArraysBase.isTypeSupported for binary/struct/map/null child types.NullIntolerant -> nullIntolerant field refactor.ArrayRepeat(left, right) extends BinaryExpression with ExpectsInputTypes; inputTypes = Seq(AnyDataType, IntegerType). NULL count yields NULL; count <= 0 yields empty array; count > MAX_ROUNDED_ARRAY_LENGTH throws at runtime. Wired as CometScalarFunction("array_repeat") against datafusion-spark's SparkArrayRepeat, which returns NULL for NULL count and repeats NULL elements (matching Spark).createArrayWithElementsExceedLimitError(prettyName, count); semantics unchanged.ArrayUnion(left, right) extends ArrayBinaryLike with ComplexTypeMergingExpression; result is left-side distinct elements followed by new right-side elements. Wired as CometScalarFunction("array_union").nullIntolerant = true moves into ArrayBinaryLike; overflow path uses arrayFunctionWithElementsExceedLimitError.array_distinct. Result element ordering also matches Spark (left-side distinct elements followed by new right-side elements).ArraysOverlap(left, right); three-valued logic (TRUE if any common non-null element, NULL if a null is present and no overlap is found in non-nulls, FALSE otherwise). Comet routes via CometArraysOverlap to the native spark_arrays_overlap UDF, which implements the same three-valued logic.NullIntolerant -> nullIntolerant field refactor.ArraysZip(children, names); returns an array of structs, padding shorter inputs with NULL. Comet routes via CometArraysZip and rejects unsupported child element types (anything outside primitives, decimals, dates/timestamps, strings, binary, and nested arrays/structs of those).IllegalArgumentException to SparkIllegalArgumentException("_LEGACY_ERROR_TEMP_3235"); runtime unchanged.ElementAt(left, right, defaultValueOutOfBound, failOnError); group label map_funcs. Comet supports ArrayType input through native ListExtract and MapType input through native map_extract.NullIntolerant -> nullIntolerant field refactor; group label changes to collection_funcs; ANSI default flips to true so out-of-bound throws by default. Comet wires failOnError through to native ListExtract.Flatten(child) extends UnaryExpression; returns NULL if any inner sub-array is NULL. Comet routes via CometFlatten and falls back for child types containing BinaryType / StructType / MapType (limitation of ArraysBase.isTypeSupported).NullIntolerant -> nullIntolerant field refactor.GetArrayItem(child, ordinal, failOnError); inputTypes = Seq(AnyDataType, IntegralType). Comet routes via CometGetArrayItem, wiring failOnError through to the proto.true.inputTypes tightened to Seq(ArrayType, IntegralType) (analysis-time only); runtime unchanged.Sequence(start, stop, stepOpt, timeZoneId); Sequence.impl selects the implementation from dataType.elementType, so the integral/temporal split is knowable at plan time. Codegen for the integral path checks boundaries with a plain IllegalArgumentException("Illegal sequence boundaries: ..."), then calls the static Sequence.sequenceLength, which raises SparkRuntimeException(_LEGACY_ERROR_TEMP_2161) past MAX_ROUNDED_ARRAY_LENGTH and internalError("Unreachable code reached.") when stop - start overflows Long but the exact length is within the limit. Default step is per-row start <= stop ? 1 : -1.DataTypeUtils.sameType, PhysicalIntegralType.integral); runtime semantics identical to 3.4.3.SparkIllegalArgumentException(_LEGACY_ERROR_TEMP_3243) and the length error becomes COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER (now carrying the function name); adds throwable optimizer hint. sequenceLength itself is unchanged.Sequence class body to 4.0.1.ByteType/ShortType/IntegerType/LongType) via CometSequence to the native spark_sequence kernel (#5349): one pass over the generated elements, child buffer reserved once per batch, no per-row allocation. The two-argument form is evaluated with Spark's per-row default step inside the kernel. Both error conditions and the internal-error edge are reproduced through SparkError and mapped per Spark version by ShimSparkErrorConverter. Date/timestamp/timestamp_ntz sequences return Unsupported and run on the JVM codegen dispatcher (CodegenDispatchFallback), pending the timezone/DST/legacy-calendar work.i32, so the sum of every row’s length in a single Arrow batch must fit in i32::MAX. Spark itself has no equivalent limit because it stores each row as its own long[]. If the total is exceeded, or if the allocator refuses the reservation, the query fails with a SparkError::SequenceBatchTooLarge message that names spark.comet.batchSize as the actionable knob (lower it to group fewer rows per batch). The try_reserve path guarantees the failure surfaces as a query error rather than an allocator abort.CometSequence reports Unsupported for any Sequence whose start, stop, or step is not a leaf expression, and routes those through the JVM codegen dispatcher (CodegenDispatchFallback). DataFusion evaluates each scalar-UDF argument over the whole batch before calling the outer kernel, so a non-leaf argument would run on rows that Spark's per-row null short-circuit (or a CASE branch) would have discarded, and could raise where Spark would have returned NULL.Shuffle(child, randomSeed: Option[Long]); inputTypes = Seq(ArrayType), dataType = child.dataType, non-deterministic and stateful. Seeds a Commons Math3 MersenneTwister with randomSeed + partitionIndex and applies the “inside-out” Fisher-Yates from RandomIndicesGenerator. Only the one-argument shuffle(array) form exists in SQL. NULL input returns NULL without advancing the RNG.Shuffle(child, seed: Expression), exposing shuffle(array, seed) in SQL (seed must be an integer/long literal). RandomIndicesGenerator and the eval logic are unchanged.CometShuffle and a dedicated stateful ShuffleExpr that reproduces the same MersenneTwister and inside-out Fisher-Yates, so results match Spark bit for bit. childTypesSupportLevel falls back for binary/struct/map element types, consistent with the other array expressions.SortArray(base, ascendingOrder) extends BinaryExpression with ArraySortLike; the second arg must be a Literal(_: Boolean, BooleanType). Comet CometSortArray flags Incompatible under strict floating-point and falls back for nested arrays whose innermost element is Struct or Null.ArraySortLike and NullIntolerant are removed, nullIntolerant = true becomes an override, and ascendingOrder is widened to accept any foldable boolean (not just Literal). Comet's CometSortArray still requires a Literal, so the new foldable form falls back at convert time.