array_funcs Expression Audits#

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.

array#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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).

  • Spark 4.0.1 (audited 2026-05-27): semantics unchanged.

  • Spark 4.1.1 (audited 2026-05-27): adds contextIndependentFoldable override; runtime semantics unchanged.

array_append#

  • Spark 3.4.3 (audited 2026-05-27): standalone BinaryExpression, evaluated directly. Comet routes via CometArrayAppend.

  • Spark 3.5.8 (audited 2026-05-27): identical to 3.4.3.

  • Spark 4.0.1 (audited 2026-05-27): now 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).

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

array_compact#

  • Spark 3.4.3 (audited 2026-05-27): 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").

  • Spark 3.5.8 (audited 2026-05-27): identical to 3.4.3.

  • Spark 4.0.1 (audited 2026-05-27): the replacement is wrapped in KnownNotContainsNull(...). The 4.x Spark4xCometExprShim strips the wrapper and emits the same array_compact call.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

array_contains#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. ArrayContains(left, right) extends BinaryExpression with NullIntolerant with Predicate; inputTypes uses findWiderTypeWithoutStringPromotionForTwo. Wired as CometScalarFunction("array_contains").

  • Spark 4.0.1 (audited 2026-05-27): NullIntolerant trait replaced by nullIntolerant: Boolean; checkInputDataTypes adopts DataTypeUtils.sameType (collation-aware in 4.x).

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Float/double arrays containing NaN and signed zero match Spark; DataFusion canonicalizes them the same way as Spark’s SQLOrderingUtil.

array_distinct#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. ArrayDistinct(child) over ArraySetLike; uses SQLOpenHashSet, which merges NaNs but keeps +0.0 and -0.0 distinct. Wired as CometScalarFunction("array_distinct").

  • Spark 4.0.1 (audited 2026-05-27): NullIntolerant -> nullIntolerant field refactor.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Float/double element types run natively only on Spark 4.2.0, which normalizes the arguments in the plan (SPARK-54918); every other version falls back by default. See Array distinct and union.

array_except#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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.

  • Spark 4.0.1 (audited 2026-05-27): nullIntolerant = true moves into ArrayBinaryLike; the overflow path uses arrayFunctionWithElementsExceedLimitError.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Float/double NaN and signed-zero canonicalization matches array_distinct. (array_except still falls back by default for the null-handling/ordering reasons noted above.)

array_insert#

  • Spark 3.4.3 audited 2026-04-02

  • Spark 3.5.8 audited 2026-04-02

  • Spark 4.0.1 audited 2026-04-02 (pos=0 error message differs from Spark)

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

array_intersect#

  • Spark 3.4.3 audited 2026-04-24 (result element order may differ from Spark when the right array is longer than the left; DataFusion probes the longer side)

  • Spark 3.5.8 audited 2026-04-24 (same ordering incompatibility as 3.4.3)

  • Spark 4.0.1 audited 2026-04-24 (ordering incompatibility as above; collated strings now fall back to Spark)

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Current status: 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.

array_join#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. ArrayJoin(array, delimiter, nullReplacement). Comet routes via CometArrayJoin to DataFusion’s array_to_string.

  • Spark 4.0.1 (audited 2026-05-27): inputTypes widened to AbstractArrayType(StringTypeWithCollation(supportsTrimCollation = true)); non-binary collations not propagated (#2190).

  • Spark 4.1.1 (audited 2026-05-27): adds contextIndependentFoldable override; runtime unchanged.

  • Current status: 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).

array_max#

  • Spark 3.4.3 (audited 2026-08-22): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-08-22): ArrayMax skips NULL elements and returns NULL for an empty or all-NULL array. SQLOrderingUtil treats all NaNs as equal and greater than non-NaN values, and signed zeros as equal. The first equal maximum is retained. Nested arrays and structs compare lexicographically, with NULL fields or elements ordered first.

  • Spark 4.0.1 (audited 2026-08-22): NullIntolerant becomes a nullIntolerant field. Extrema semantics are unchanged; string ordering can use non-default collations.

  • Spark 4.1.1 (audited 2026-08-22): identical to 4.0.1.

  • Current status: CometArrayMax uses the native SparkArrayExtrema UDF. Typed float/double scans and recursive array/struct comparisons follow Spark’s ordering and preserve the original first equal element, including its zero sign and NaN representation. This path is used in both strict and non-strict floating-point modes without the JVM codegen dispatcher. Other scalar element types retain the existing DataFusion implementation. Non-UTF8_BINARY string collations, including nested fields, use Spark’s JVM codegen dispatcher inside the Comet pipeline by default. If the dispatcher is disabled, these cases fall back to Spark unless incompatible native execution is explicitly enabled (#4496).

array_min#

  • Spark 3.4.3 (audited 2026-08-22): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-08-22): mirrors ArrayMax, retaining the first equal minimum. The NULL, NaN, signed-zero, and nested comparison rules are the same.

  • Spark 4.0.1 (audited 2026-08-22): same trait refactor and collation support as array_max, with no change in floating-point extrema semantics.

  • Spark 4.1.1 (audited 2026-08-22): identical to 4.0.1.

  • Current status: CometArrayMin shares the native SparkArrayExtrema implementation and support boundary with array_max. Both floating-point modes use Spark-compatible native ordering, preserving the original first equal minimum. Non-default string collations use the same JVM codegen dispatch and dispatcher-disabled fallback as array_max (#4496).

array_position#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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).

  • Spark 4.0.1 (audited 2026-05-27): NullIntolerant -> nullIntolerant field refactor.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

array_prepend#

  • Spark 3.4.3 (audited 2026-06-24): array_prepend does not exist (added in Spark 3.5.0). The SQL file test carries MinSparkVersion: 3.5 so it is skipped here.

  • Spark 3.5.8 (audited 2026-06-24): 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]).

  • Spark 4.0.1 (audited 2026-06-24): ArrayPrepend now extends ArrayPendBase but the replacement is unchanged (ArrayInsert(left, Literal(1), right)).

  • Spark 4.1.1 (audited 2026-06-24): identical to 4.0.1.

array_remove#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. ArrayRemove(left, right); removes all occurrences equal to right. Falls back via ArraysBase.isTypeSupported for binary/struct/map/null child types.

  • Spark 4.0.1 (audited 2026-05-27): NullIntolerant -> nullIntolerant field refactor.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Current status: CometArrayRemove sends arrays whose elements hold a FLOAT or DOUBLE at any depth to the native spark_array_remove, which compares as Spark’s genEqual does (-0.0 equals 0.0, all NaNs are equal, nested elements through spark_equality) and keeps the bits of the elements it keeps. Other element types use DataFusion’s array_remove_all. Null elements stay, and a null array or value gives null.

  • Performance (tuned 2026-10-01, PR #6518): with a constant value, FLOAT and DOUBLE elements get one keep mask over all the values from BooleanBuffer::collect_bool, with the inverted validity ORed in, one running popcount over the mask’s words for each row’s kept count, and one filter. 41-76% less time than DataFusion’s array_remove_all with a constant value, and 47-95% less with a value per row. Benchmark: benches/float_arrays.rs.

array_repeat#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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).

  • Spark 4.0.1 (audited 2026-05-27): error message uses createArrayWithElementsExceedLimitError(prettyName, count); semantics unchanged.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

array_union#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. ArrayUnion(left, right) extends ArrayBinaryLike with ComplexTypeMergingExpression; result is left-side distinct elements followed by new right-side elements. Wired as CometScalarFunction("array_union").

  • Spark 4.0.1 (audited 2026-05-27): nullIntolerant = true moves into ArrayBinaryLike; overflow path uses arrayFunctionWithElementsExceedLimitError.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Float/double element types follow the same Spark 4.2.0-only gate as array_distinct. Result element ordering matches Spark (left-side distinct elements followed by new right-side elements).

arrays_overlap#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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.

  • Spark 4.0.1 (audited 2026-05-27): NullIntolerant -> nullIntolerant field refactor.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

arrays_zip#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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).

  • Spark 4.0.1 (audited 2026-05-27): the length-mismatch error switches from IllegalArgumentException to SparkIllegalArgumentException("_LEGACY_ERROR_TEMP_3235"); runtime unchanged.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

element_at#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. ElementAt(left, right, defaultValueOutOfBound, failOnError); group label map_funcs. Comet supports ArrayType input through native ListExtract and MapType input through native map_extract.

  • Spark 4.0.1 (audited 2026-05-27): 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.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

flatten#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. 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).

  • Spark 4.0.1 (audited 2026-05-27): NullIntolerant -> nullIntolerant field refactor.

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

get#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. GetArrayItem(child, ordinal, failOnError); inputTypes = Seq(AnyDataType, IntegralType). Comet routes via CometGetArrayItem, wiring failOnError through to the proto.

  • Spark 4.0.1 (audited 2026-05-27): semantics unchanged; ANSI default flips to true.

  • Spark 4.1.1 (audited 2026-05-27): inputTypes tightened to Seq(ArrayType, IntegralType) (analysis-time only); runtime unchanged.

sequence#

  • Spark 3.4.3 (audited 2026-08-29): 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.

  • Spark 3.5.8 (audited 2026-08-29): internal refactors only (DataTypeUtils.sameType, PhysicalIntegralType.integral); runtime semantics identical to 3.4.3.

  • Spark 4.0.1 (audited 2026-08-29): boundary error becomes 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.

  • Spark 4.1.1 (audited 2026-08-29): byte-identical Sequence class body to 4.0.1.

  • Comet routes integral element types (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.

  • Per-batch capacity ceiling: the native kernel writes every row’s generated elements into one Arrow child buffer whose offsets are 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.

  • Argument-shape restriction: 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#

  • Spark 3.4.3 (audited 2026-07-02): 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.

  • Spark 3.5.8 (audited 2026-07-02): identical to 3.4.3.

  • Spark 4.0.1 (audited 2026-07-02): adds the two-argument constructor Shuffle(child, seed: Expression), exposing shuffle(array, seed) in SQL (seed must be an integer/long literal). RandomIndicesGenerator and the eval logic are unchanged.

  • Spark 4.1.1 (audited 2026-07-02): identical to 4.0.1. Comet routes via 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.

sort_array#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8.

  • Spark 3.5.8 (audited 2026-05-27): baseline. SortArray(base, ascendingOrder) extends BinaryExpression with ArraySortLike; the second arg must be a Literal(_: Boolean, BooleanType).

  • Spark 4.0.1 (audited 2026-05-27): trait set changes substantively: ArraySortLike and NullIntolerant are removed, nullIntolerant = true becomes an override, and ascendingOrder is widened to accept any foldable boolean (not just Literal).

  • Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.

  • Current status: elements with a FLOAT or DOUBLE at any depth go to the native spark_sort_array, which sorts as Spark’s generated code does: a stable sort in Spark’s SQL ordering, except that for an ascending sort of FLOAT or DOUBLE elements that cannot be null, the serde has it put -0.0 before 0.0, as java.util.Arrays.sort does. This holds in strict floating-point mode too. Other supported element types use DataFusion’s array_sort, and unsupported ones route through the codegen dispatcher. ascendingOrder may be any foldable boolean, which the serde evaluates.

  • Performance (tuned 2026-10-01, PR #6518): FLOAT and DOUBLE rows sort within one copy of all the values, each row’s valid values copied without branching next to its nulls. Rows of up to 20 elements use the stable sort_by. Longer rows use sort_unstable_by and then put the zero and NaN runs back in their original order, because Rust’s stable sort is up to 1.8 times slower between 33 and 63 elements. From 2% more to 21% less time than DataFusion’s array_sort. Benchmark: benches/float_arrays.rs.