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 viaCometCreateArray(nativemake_array) and special-cases the empty-array case to dodge a known DataFusioncoerce_typesissue (#3338).Spark 4.0.1 (audited 2026-05-27): semantics unchanged.
Spark 4.1.1 (audited 2026-05-27): adds
contextIndependentFoldableoverride; runtime semantics unchanged.
array_append#
Spark 3.4.3 (audited 2026-05-27): standalone
BinaryExpression, evaluated directly. Comet routes viaCometArrayAppend.Spark 3.5.8 (audited 2026-05-27): identical to 3.4.3.
Spark 4.0.1 (audited 2026-05-27): now
RuntimeReplaceableand rewritten toArrayInsert(arr, Literal(-1), elem).CometArrayAppendis therefore unreachable; dispatch goes throughCometArrayInsert(which carries its ownIncompatiblenotes documented at thearray_insertentry).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 throughCometArrayFilter, which emits a call to DataFusion’s built-inarray_compact(fromdatafusion-functions-nested) viaCometScalarFunction("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.xSpark4xCometExprShimstrips the wrapper and emits the samearray_compactcall.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;inputTypesusesfindWiderTypeWithoutStringPromotionForTwo. Wired asCometScalarFunction("array_contains").Spark 4.0.1 (audited 2026-05-27):
NullIntoleranttrait replaced bynullIntolerant: Boolean;checkInputDataTypesadoptsDataTypeUtils.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)overArraySetLike; usesSQLOpenHashSet, which merges NaNs but keeps+0.0and-0.0distinct. Wired asCometScalarFunction("array_distinct").Spark 4.0.1 (audited 2026-05-27):
NullIntolerant->nullIntolerantfield 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 viaCometArrayExceptand unconditionally flagsIncompatible(“Null handling and ordering may differ from Spark”); also falls back forBinaryType/StructTypeelement types.Spark 4.0.1 (audited 2026-05-27):
nullIntolerant = truemoves intoArrayBinaryLike; the overflow path usesarrayFunctionWithElementsExceedLimitError.Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.
Float/double NaN and signed-zero canonicalization matches
array_distinct. (array_exceptstill 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:
CometArrayIntersectreportsIncompatiblebecause 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 reportedUnsupportedand 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 viaCometArrayJointo DataFusion’sarray_to_string.Spark 4.0.1 (audited 2026-05-27):
inputTypeswidened toAbstractArrayType(StringTypeWithCollation(supportsTrimCollation = true)); non-binary collations not propagated (#2190).Spark 4.1.1 (audited 2026-05-27): adds
contextIndependentFoldableoverride; runtime unchanged.Current status:
CometArrayJoinreportsCompatiblewhen 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 anIsNullguard, sincearray_to_stringreads a nullnull_stringas “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):
ArrayMaxskips NULL elements and returns NULL for an empty or all-NULL array.SQLOrderingUtiltreats 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):
NullIntolerantbecomes anullIntolerantfield. 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:
CometArrayMaxuses the nativeSparkArrayExtremaUDF. 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:
CometArrayMinshares the nativeSparkArrayExtremaimplementation and support boundary witharray_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 asarray_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-basedLongTypeposition, 0 if not found, NULL if either input is NULL.CometArrayPositionfalls 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->nullIntolerantfield 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_prependdoes not exist (added in Spark 3.5.0). The SQL file test carriesMinSparkVersion: 3.5so 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 seesArrayPrepend; dispatch goes throughCometArrayInsert(which carries its own notes documented at thearray_insertentry). 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):
ArrayPrependnow extendsArrayPendBasebut thereplacementis 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 toright. Falls back viaArraysBase.isTypeSupportedfor binary/struct/map/null child types.Spark 4.0.1 (audited 2026-05-27):
NullIntolerant->nullIntolerantfield refactor.Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.
Current status:
CometArrayRemovesends arrays whose elements hold aFLOATorDOUBLEat any depth to the nativespark_array_remove, which compares as Spark’sgenEqualdoes (-0.0equals0.0, all NaNs are equal, nested elements throughspark_equality) and keeps the bits of the elements it keeps. Other element types use DataFusion’sarray_remove_all. Null elements stay, and a null array or value gives null.Performance (tuned 2026-10-01, PR #6518): with a constant value,
FLOATandDOUBLEelements get one keep mask over all the values fromBooleanBuffer::collect_bool, with the inverted validity ORed in, one running popcount over the mask’s words for each row’s kept count, and onefilter. 41-76% less time than DataFusion’sarray_remove_allwith 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_LENGTHthrows at runtime. Wired asCometScalarFunction("array_repeat")againstdatafusion-spark’sSparkArrayRepeat, 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 asCometScalarFunction("array_union").Spark 4.0.1 (audited 2026-05-27):
nullIntolerant = truemoves intoArrayBinaryLike; overflow path usesarrayFunctionWithElementsExceedLimitError.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 viaCometArraysOverlapto the nativespark_arrays_overlapUDF, which implements the same three-valued logic.Spark 4.0.1 (audited 2026-05-27):
NullIntolerant->nullIntolerantfield 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 viaCometArraysZipand 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
IllegalArgumentExceptiontoSparkIllegalArgumentException("_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 labelmap_funcs. Comet supportsArrayTypeinput through nativeListExtractandMapTypeinput through nativemap_extract.Spark 4.0.1 (audited 2026-05-27):
NullIntolerant->nullIntolerantfield refactor; group label changes tocollection_funcs; ANSI default flips totrueso out-of-bound throws by default. Comet wiresfailOnErrorthrough to nativeListExtract.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 viaCometFlattenand falls back for child types containingBinaryType/StructType/MapType(limitation ofArraysBase.isTypeSupported).Spark 4.0.1 (audited 2026-05-27):
NullIntolerant->nullIntolerantfield 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 viaCometGetArrayItem, wiringfailOnErrorthrough 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):
inputTypestightened toSeq(ArrayType, IntegralType)(analysis-time only); runtime unchanged.
sequence#
Spark 3.4.3 (audited 2026-08-29):
Sequence(start, stop, stepOpt, timeZoneId);Sequence.implselects the implementation fromdataType.elementType, so the integral/temporal split is knowable at plan time. Codegen for the integral path checks boundaries with a plainIllegalArgumentException("Illegal sequence boundaries: ..."), then calls the staticSequence.sequenceLength, which raisesSparkRuntimeException(_LEGACY_ERROR_TEMP_2161)pastMAX_ROUNDED_ARRAY_LENGTHandinternalError("Unreachable code reached.")whenstop - startoverflows Long but the exact length is within the limit. Default step is per-rowstart <= 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 becomesCOLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER(now carrying the function name); addsthrowableoptimizer hint.sequenceLengthitself is unchanged.Spark 4.1.1 (audited 2026-08-29): byte-identical
Sequenceclass body to 4.0.1.Comet routes integral element types (
ByteType/ShortType/IntegerType/LongType) viaCometSequenceto the nativespark_sequencekernel (#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 throughSparkErrorand mapped per Spark version byShimSparkErrorConverter. Date/timestamp/timestamp_ntz sequences returnUnsupportedand 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 ini32::MAX. Spark itself has no equivalent limit because it stores each row as its ownlong[]. If the total is exceeded, or if the allocator refuses the reservation, the query fails with aSparkError::SequenceBatchTooLargemessage that namesspark.comet.batchSizeas the actionable knob (lower it to group fewer rows per batch). Thetry_reservepath guarantees the failure surfaces as a query error rather than an allocator abort.Argument-shape restriction:
CometSequencereportsUnsupportedfor anySequencewhosestart,stop, orstepis 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 aCASEbranch) would have discarded, and could raise where Spark would have returnedNULL.
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 Math3MersenneTwisterwithrandomSeed + partitionIndexand applies the “inside-out” Fisher-Yates fromRandomIndicesGenerator. Only the one-argumentshuffle(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), exposingshuffle(array, seed)in SQL (seed must be an integer/long literal).RandomIndicesGeneratorand the eval logic are unchanged.Spark 4.1.1 (audited 2026-07-02): identical to 4.0.1. Comet routes via
CometShuffleand a dedicated statefulShuffleExprthat reproduces the same MersenneTwister and inside-out Fisher-Yates, so results match Spark bit for bit.childTypesSupportLevelfalls 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 aLiteral(_: Boolean, BooleanType).Spark 4.0.1 (audited 2026-05-27): trait set changes substantively:
ArraySortLikeandNullIntolerantare removed,nullIntolerant = truebecomes an override, andascendingOrderis widened to accept any foldable boolean (not justLiteral).Spark 4.1.1 (audited 2026-05-27): identical to 4.0.1.
Current status: elements with a
FLOATorDOUBLEat any depth go to the nativespark_sort_array, which sorts as Spark’s generated code does: a stable sort in Spark’s SQL ordering, except that for an ascending sort ofFLOATorDOUBLEelements that cannot be null, the serde has it put-0.0before0.0, asjava.util.Arrays.sortdoes. This holds in strict floating-point mode too. Other supported element types use DataFusion’sarray_sort, and unsupported ones route through the codegen dispatcher.ascendingOrdermay be any foldable boolean, which the serde evaluates.Performance (tuned 2026-10-01, PR #6518):
FLOATandDOUBLErows 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 stablesort_by. Longer rows usesort_unstable_byand 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’sarray_sort. Benchmark:benches/float_arrays.rs.