Spark expression support
One row per Spark expression: whether VecRuntime compiles it, which lane types it accepts, and the incompatibility that still makes it fall back where one exists. Modelled on Comet's Spark Expression Support page; the companion for operators is docs/operators.md.
For the tested-against-Spark subset — the exact cases that run on VecRuntime versus fall back, taken from the ported Comet suites — see the Compatibility matrix, and Testing & correctness for how it is all validated.
How to read it. ExpressionCompiler.compile (spark/src/main/scala/io/vecruntime/spark/expr/)
is one match over Catalyst expressions. Every case either produces a VectorExpr node with a kernel
behind it or returns a Left(reason); the reason is what the operator records (VectorFallback), what
the UI tooltip and spark.vecruntime.explainFallback.enabled show, and what the suites assert on with
checkFallback(..., reasonContains). An expression that is not matched at all records
unsupported expression <Class>: <sql> -- so a row missing from the Supported table below is a
not yet, and the Planned table says which issue tracks it.
Types are the VecType lanes, since that is what decides support here:
| Lane | Spark types |
|---|---|
| BOOL | boolean |
| INT32 | int, date |
| INT64 | bigint, timestamp, decimal(p <= 18) (unscaled value; the scale stays in the Spark type) |
| FLOAT64 | double |
| UTF8 | string (carried through filters, projections and as a grouping/join key; compared in UTF8_BINARY order by StringCompareKernels) |
| DECIMAL128 | decimal(p > 18) as two little-endian long limbs (#257): carried through filters, projections, sort keys, group and join keys, wide sum buffers; compared on the limbs -- = < <= > >= <> and IN against wide literals or a wide column -- and computed with + - * /, abs, negative and every Cast below (#258); %, pmod, round/bround/ceil/floor and a string source for a cast are still refused |
Anything else -- float, short, byte, binary, array, map, struct,
intervals -- has no lane: a column of such a type records unsupported type <type> for <name>, a
result of such a type unsupported output type <type> for <name>, whatever the expression. A wide
decimal expression (arithmetic, a cast) is still refused with its own reason until the rest of #258
lands (#26 / #27 / #28).
Literals are compiled as operands of a supported type: int, bigint, double, date,
timestamp, decimal(p <= 18), string. A typed NULL literal (CAST(NULL AS INT), the NULL AS col
a union coerces) of any lane type is an all-invalid column (NullLiteralExpr); the operators that need a
value literal -- comparison operands, IN lists, string arguments -- refuse a null there with their own
reason. A boolean literal records unsupported literal type boolean by design: booleans are bitmaps,
and every predicate over a boolean literal (b = true, <=>) takes a dedicated path before compiling
it. An expression made only of literals is
refused where it would be pointless as a kernel (comparison of two literals, arithmetic on two literals, cast of a literal, literal predicate, ...) -- Spark's optimizer normally folds those away
before we see them; a bare literal projection (SELECT 1 FROM t) is supported and materialised as a
constant column.
Spark 4 defaults to ANSI mode. Where ANSI changes the semantics (integer overflow, division by zero, decimal overflow) the row says whether the kernel raises the same error Spark does or falls back. Errors are raised only for rows that are active -- survivors of earlier conjuncts / the selection -- matching Spark's short-circuit behaviour.
Supported
| Expression | Types | Notes / fallback reasons |
|---|---|---|
Column reference (AttributeReference, BoundReference) |
all lanes | unbound attribute <name> if the attribute is not in the operator's input; unsupported type <type> for <name> otherwise |
Alias |
any | Transparent |
| Literal | INT32, INT64, FLOAT64, date, timestamp, decimal(<=18), string; a typed NULL of any lane type |
Operand, a whole projected column (a string literal becomes a constant UTF8 column; a typed null an all-invalid one), or a result column beside aggregates ('store' AS channel, sum(...)). unsupported literal type <type> (boolean, binary, ...), null literal of <type> for a type without a lane |
= < <= > >= and != / <> (Not(EqualTo)) |
INT32, INT64, FLOAT64 (incl. date, timestamp, decimal(<=18) as their lane), UTF8, DECIMAL128 (a wide column against a wide literal or column of the same type: signed high limb, then unsigned low limb, a scalar loop -- the lane has no SIMD path; #258) | Operands must have the same Spark type -- Spark's coercion inserts casts, which then have to compile (see Cast): comparison operands differ: <t1> vs <t2>. Strings compare in Spark's default UTF8_BINARY order (unsigned byte-wise, a prefix first) against a literal or another string column; a dictionary-encoded column is compared once per dictionary entry. Booleans: comparison not supported for boolean (#32). Doubles compare with Spark's ordering (NaN equal to NaN and greatest, -0.0 == 0.0); KnownFloatingPointNormalized is an identity, NormalizeNaNAndZero a real pass (see below). |
IN (v1, ..., vN) |
any comparable lane incl. UTF8 | Every element must be a non-null literal of the value's type (InExpr: the value is evaluated once, one equality pass per literal, through the dictionary for dictionary-encoded strings). NULL in IN list, IN list is not all literals, IN operands differ: ..., IN over a literal, empty IN list; above spark.sql.optimizer.inSetConversionThreshold literals Spark rewrites to InSet -- see its row |
InSet (an IN list above spark.sql.optimizer.inSetConversionThreshold, 10 by default) |
INT32, INT64, FLOAT64 (incl. date, decimal(<=18) as their lane), UTF8 | InSetExpr over PredicateKernels.inSet: the set's values sorted once at compile time, one binary search per row; doubles are matched on their bit images, which is Spark's boxed-set semantics (NaN matches NaN, -0.0 does not match 0.0). Strings take the IN path above (one equality pass per element, through the dictionary). NULL in IN set (Spark's result is then null for non-members), IN not supported for <type> |
<=> (EqualNullSafe) |
every comparable lane incl. UTF8 and BOOL | NullSafeEqExpr: the ordinary equality, then (eq AND valid) OR (both null) on the validity bitmaps -- never null. x <=> NULL compiles to x IS NULL. A boolean literal side falls back (<=> of a boolean literal); same-type rule as = |
isnan |
FLOAT64 | IsNaNExpr over a Vector API IS_NAN test packed into the bitmap; false (not null) where the argument is null, as Spark's. isnan over <type> for a non-double |
= <> < <= > >= on BOOL |
BOOL column vs boolean literal or BOOL column | BoolCompareExpr / BoolCompareScalarExpr on the packed words (false < true); null where an operand is. comparison against NULL |
BETWEEN |
as its comparisons | Spark rewrites it to two comparisons before planning; pinned by a test |
startswith, endswith, contains; LIKE 'p%', LIKE '%p', LIKE '%p%' (Spark's LikeSimplification rewrites these three shapes into the functions) |
UTF8 vs a string literal | StringMatchExpr over StringMatchKernels: byte-level under the default UTF8_BINARY collation (a UTF-8 pattern can only match at character boundaries), the empty pattern matches everything, once per dictionary entry on dictionary-encoded columns; prefix and suffix are one MemorySegment.mismatch over a fixed range, contains scans for the first byte and confirms with mismatch. null pattern, string pattern is not a literal, string pattern of type <t>, string match not supported for <t>, string match on a literal. A LIKE the optimizer does not simplify because it has several % wildcards -- '%a%b%', 'p%a%b', 'a%b%s', adjacent %% -- is a multi-token matcher (LikeTokensExpr over StringMatchKernels.matchTokens, #264): the prefix must occupy the start, each token is found left to right with the scan resuming after the previous match (the leftmost match leaves the most room, so greedy is exact), and the suffix must occupy the end without overlapping the last token; same dictionary handling. null LIKE pattern, LIKE pattern with a \_` wildcard not supported, LIKE pattern with an escape character not supported, LIKE pattern is not a literal, LIKE over |
instr, locate / position |
UTF8 haystack and needle (lanes or literals); INT32 start (literal or lane) | LocateExpr over StringSearchKernels: Spark's indexOf -- steps by code points and returns a 1-based code-point position, an empty needle is found at position 1; locate's rules -- a null start gives 0 (not null) before the other arguments are looked at, a start below 1 gives 0, otherwise the search begins at that code point |
replace(str, search[, replacement]) |
UTF8 lanes or literals | ReplaceExpr: Spark's byte-level find-and-copy of every non-overlapping occurrence; an empty haystack or search leaves the row unchanged; a missing replacement deletes |
translate(str, from, to) |
UTF8 subject; from and to string literals |
TranslateExpr: Spark's buildDict -- the first mapping per from code point wins, a from code point past the end of to is deleted -- applied per code point. Column from/to: translate with a non-literal from/to string not supported |
substring_index(str, delim, count) |
UTF8 subject and delimiter (lanes or literals); INT32 count | SubstringIndexExpr: Spark's subStringIndex byte algorithm -- count occurrences from the left (positive) or right (negative), the whole string when fewer are found, '' for count 0 or an empty delimiter |
split_part(str, delim, part) |
UTF8 subject and delimiter; INT32 part | SplitPartExpr: the byte split keeping empty parts (an empty delimiter gives one part), counted from either end, '' past the ends; part 0 raises Spark's INVALID_INDEX_OF_ZERO -- for active rows only, so a filtered row never raises. Matched on Spark's rewrite element_at(split(str, delim), part, '') |
find_in_set(word, set) |
UTF8 lanes or literals; INT32 out | FindInSetExpr: Spark's findInSet -- the 1-based position of word among the comma-separated set, 0 when absent or when the word contains a comma |
upper / ucase, lower / lcase, initcap |
UTF8 lane in the default UTF8_BINARY collation |
CaseMapExpr over StringCaseKernels: the kernel maps the ASCII rows in place (a byte-range test and a 0x20 flip) and flags the rows it cannot decide -- any byte at or above 0x80, and for initcap under Spark 4's default ICU case mappings anything but letters and spaces -- which are computed by Spark's own CollationSupport for this session's spark.sql.icu.caseMappings.enabled, so every row is Spark's result. Dictionary input is decided per entry. A collated column is not byte-mapped: upper over string collate utf8_lcase not supported (a collated string follows ICU rules) |
trim / btrim, ltrim, rtrim, trim(BOTH | LEADING | TRAILING [trimStr] FROM str) |
UTF8 lane; the trim string a non-null literal | TrimExpr: offset arithmetic -- scan in from each end while the code point is in the trim set (the single space by default; Spark's trim() trims only ASCII 32). A trim string from a column falls back: trim with a non-literal trim string not supported |
concat(s1, s2, ...) |
UTF8 lanes and string literals (at least one lane; string children only) | ConcatExpr over StringConcatKernels: each row's length is the sum across the inputs, one buffer from the prefix sum, a per-input copy loop; null when any input is null. Array and binary concat are not lanes: unsupported expression |
concat_ws(sep, s1, s2, ...) |
UTF8 separator (lane or literal) and UTF8 inputs; string children only | ConcatWsExpr: Spark's concatWs -- null only for a null separator, null inputs skipped, the separator only between live inputs, '' when none is live. Array arguments fall back |
elt(n, s1, s2, ...) |
INT32 index (lane or literal); UTF8 inputs | EltExpr: the input at the 1-based index; a null or out-of-range index and a null pick give null. Under ANSI an out-of-range index raises Spark's INVALID_ARRAY_INDEX -- only for rows still active, so a row a filter removed never raises |
length / len / char_length / character_length, octet_length, bit_length |
UTF8 subject; INT32 out | StringMeasureExpr over StringLengthKernels: length is Spark's numChars -- a walk over Spark's first-byte width table, so malformed input counts as Spark counts it; octet_length is the offset difference, bit_length that times eight. Dictionary-encoded input is measured once per dictionary entry and gathered. A BINARY subject is not a lane: length over binary not supported |
ascii |
UTF8 subject; INT32 out | Spark's rule: the code point of the first character (substring(0, 1).toString.codePointAt(0)), 0 for the empty string, U+FFFD for a malformed first sequence as Java's decoder gives |
chr / char |
INT32 or INT64 lane (the analyzer casts to bigint) | ChrExpr: Spark's rule -- a negative gives '', otherwise the character n & 0xFF: 0 is the NUL character, 1..127 one byte, 128..255 the two-byte UTF-8 encoding |
substring / substr (str FROM pos [FOR len]), left, right |
UTF8 subject; INT32 position and length as literals or lanes | SubstringExpr over StringSliceKernels: Spark's UTF8String.substringSQL per lane -- 1-based code points, a negative start counts from the end, out-of-range positions and negative lengths clamp to the empty string, substring(s, 0, n) is the first n code points -- with Spark's first-byte width table. left/right are Spark's own rewrites onto it. Two passes: per-row source ranges and lengths, one data buffer from their prefix sum, then the copies; dictionary input is read through the dictionary and the output is plain. A BINARY subject is not a lane: substring over binary not supported |
lpad, rpad |
UTF8 subject; INT32 length (literal or lane); UTF8 pad as a literal or a lane (default ' ') |
PadExpr: Spark's lpad/rpad -- a non-positive length gives '', a subject of at least len code points is cut to len, otherwise whole copies of the pad then a prefix of it fill the gap (an empty pad pads nothing). A literal length above 2^20 falls back (exceeds the batch output cap); the total output of one batch is capped at 1 GiB |
repeat, space |
UTF8 subject and INT32 count (literal or lane); space takes a lane count (a literal folds in Spark) |
RepeatExpr / SpaceExpr: a non-positive count gives ''. The same literal bound and per-batch cap as the pads -- a runaway count is declined rather than allocated |
overlay(input PLACING replace FROM pos [FOR len]) |
UTF8 input; UTF8 replacement as a literal or a lane; INT32 position and length | OverlayExpr: Spark's Overlay.calculate -- the first pos - 1 code points, the replacement, then the input from code point pos + len (the replacement's length when len is omitted or negative), pos + len in Spark's int arithmetic |
AND, OR, NOT |
BOOL | Operands must be non-literal booleans: boolean literal operand, expected boolean, got <type> |
IS NULL, IS NOT NULL |
all lanes | null test on literal |
+ - * on integers |
INT32, INT64 | Same-typed operands (arithmetic operands differ). Both modes: legacy wraps like Spark; in ANSI mode (Spark's default) the wrapped result is checked with an overflow lane mask (OverflowKernels: sign trick for + -, exact product for *) and Spark's ARITHMETIC_OVERFLOW is raised -- integer overflow / long overflow with the try_add / try_subtract / try_multiply hint -- only if an active row overflowed, so rows a filter removed or an earlier conjunct decided never raise. try_add / try_subtract / try_multiply (Spark's Add/Subtract/Multiply in TRY mode) clear the same overflow mask into the validity instead: the overflowing rows are null, the rest exact, combined with the operands' nulls. Decimal try_*: try_* arithmetic not supported (decimal try_* has no issue yet) |
+ - * / on doubles |
FLOAT64 | Bit-identical in ANSI and legacy mode; / by zero yields null (legacy) or raises DIVIDE_BY_ZERO (ANSI) with Spark's error context; try_divide (Spark's Divide in TRY mode, always floating point) nulls a zero divisor in either mode, unlike /. Integer division falls back: division not supported for <type> |
+ - * / on decimals |
INT64 (decimal(<=18) operands and result) | Both operands decimal (mixed decimal and non-decimal arithmetic otherwise); Spark's result type must fit 18 digits: decimal result <type> exceeds 18 digits -- decimal(15,2) * decimal(15,2) is decimal(31,4) and falls back, except directly under a decimal sum (the speculative narrow product below, #26). Only / can overflow and it checks; ANSI raises NUMERIC_VALUE_OUT_OF_RANGE, legacy yields null. Division is Spark-exact (round half up at the result scale). unexpected decimal result scale guards against an operand shape the kernel does not expect |
+ - * / on decimals with a wide operand or result |
DECIMAL128 (operands DECIMAL128 or INT64 lanes or literals; result DECIMAL128) | WideDecimalArithExpr over WideDecimalKernels (#258): the exact result then Spark's toPrecision(p, s, HALF_UP). On the limbs when the value fits: +/- rescale the narrower operand by a power of ten and add with carry and signed-overflow detection, * multiplies two values that fit a long into an exact 128-bit product; a value that would leave 128 bits, a result scale Spark capped below the exact one, and every / (Spark's own is BigDecimal.divide(38, HALF_UP) then toPrecision) take the exact BigInteger/BigDecimal path for that row -- never a wrong value, only a slower row. Overflow of the result precision is null (legacy) or NUMERIC_VALUE_OUT_OF_RANGE (ANSI) and a zero divisor null or DIVIDE_BY_ZERO, for the active rows only. Narrow operands with a wide declared result under sum/avg keep #26's speculative INT64 path (pinned by a test); as a projected value they compile onto the wide lane. try_*, %, pmod, unary minus and abs over the lane are not yet compiled: try_* arithmetic not supported, unsupported expression <name>. |
| Unary minus | INT32, INT64, FLOAT64, decimal(<=18) | ANSI mode raises ARITHMETIC_OVERFLOW (integer overflow / long overflow, no hint) for an active MIN_VALUE; doubles and decimals never overflow here |
Cast |
INT32 -> INT64, INT32 -> FLOAT64, INT64 -> FLOAT64; int/bigint/double -> decimal(<=18); decimal -> decimal / double / bigint / int; timestamp -> date under a UTC or fixed-offset session zone (TimestampToDateExpr: floorDiv(micros + offset, micros per day); Spark inserts this cast under year(ts) etc.); a cast to the operand's own type; narrowing bigint -> int, double -> int / bigint (NarrowCastExpr: Java's rule as Spark's legacy mode -- a long wraps, a double truncates toward zero and saturates, NaN is 0; under ANSI Spark's CAST_OVERFLOW by its floor(v) <= MAX && ceil(v) >= MIN test, so NaN and the infinities raise, for active rows only); boolean int / bigint / double -> boolean (v != 0, NaN true), boolean -> int / bigint / double (1 / 0), string -> boolean (StringToBooleanExpr: Spark's spellings t/true/y/yes/1 and f/false/n/no/0, trimmed, any case; anything else null, or under ANSI CAST_INVALID_INPUT for an active row); date -> timestamp under a UTC or fixed-offset zone (DateToTimestampExpr); int / bigint / double / boolean -> string (ToStringExpr: Java's toString, which is what Spark's cast calls, through the row writer); string -> int / bigint / double (StringToNumberExpr: Spark's own parsers per row -- legacy UTF8String.toInt/toLong (trimmed, signed, an all-digit fraction dropped, null otherwise or on overflow) and Double.parseDouble with Spark's special literals NaN/Infinity/inf; under ANSI Spark's exact parsers raise CAST_INVALID_INPUT for active rows only); date / timestamp -> string (DateTimeToStringExpr: the DateFormatter and fraction TimestampFormatter Spark's cast uses, per row, any session zone); string -> date / timestamp (StringToDateTimeExpr: Spark's own stringToDate / stringToTimestamp per row, so every form Spark accepts is accepted, any session zone; null or under ANSI CAST_INVALID_INPUT for an active row) |
unsupported cast <from> -> <to> / unsupported cast target <type> for pairs outside the lane types (binary, intervals, nested types, tinyint/smallint/float which have no lane). timestamp -> date and date -> timestamp under a zone with rules: ... needs a fixed-offset session zone, not <zone> (Not planned below). Decimal casts check the range; ANSI raises, legacy nulls. try_cast runs the same nodes with null instead of raise -- the narrowing casts null the flagged rows, the string, boolean and datetime casts take their legacy null path; try_cast into or out of a decimal: try_cast not supported (no issue yet) |
Cast with a wide decimal on either side |
decimal <-> decimal at any width (DECIMAL128 <-> DECIMAL128, DECIMAL128 <-> INT64), int / bigint / date / double -> DECIMAL128, DECIMAL128 -> double / bigint / int / string | WideDecimalCastExpr over WideDecimalCastKernels (#258): Spark's changePrecision -- the value rescaled half up to the target scale and checked against the target precision (on the limbs when the value fits a long and the scale grows within one; the exact BigInteger path otherwise); a double through BigDecimal.valueOf (Spark builds its Decimal from the double's shortest string), NaN and infinities invalid; to double via BigDecimal.doubleValue; to bigint / int truncating toward zero, the wrapped value in legacy mode and CAST_OVERFLOW in ANSI; to string as Spark's Cast prints it (plain notation under ANSI, BigDecimal.toString otherwise). A value that does not fit the target is null (legacy) or CAST_OVERFLOW / NUMERIC_VALUE_OUT_OF_RANGE (ANSI) for the active rows. Refused: a string source (unsupported cast string -> decimal(p,s), Spark's Decimal.fromString rules are not reproduced yet), try_cast. |
year, month, dayofmonth / day, dayofyear, quarter, dayofweek, weekday, extract(<field> FROM date) |
INT32 days -> INT32 | DateFieldExpr over DateKernels.field: branch-free civil-from-days per lane (no LocalDate), valid for negative days and every leap rule; Spark's numbering (dayofweek 1 = Sunday, weekday 0 = Monday). Over a timestamp Spark first casts to date (see Cast). date function over <type>, date function on a literal |
trunc(date, unit) |
INT32 -> INT32 | Units YEAR/YYYY/YY, QUARTER, MONTH/MON/MM, WEEK (Monday), case-insensitive, as a string literal. trunc unit '<u>' not supported, trunc unit is not a string literal |
date_add, date_sub, datediff |
INT32 lanes | ArithExpr add/subtract on days (Spark does not overflow-check these): date +/- int and date - date; either side may be a literal. date arithmetic over <type>, date arithmetic with <type> days, arithmetic on two literals |
hour, minute, second |
INT64 micros -> INT32 | Under a UTC or fixed-offset session zone only (TimeFieldExpr: local micros = micros + offset); a zone with rules falls back: time field needs a fixed-offset session zone, not <zone>. time field over <type> |
last_day, add_months(date, months), next_day(date, literal day), weekofyear |
INT32 date lane; INT32 months (literal or lane); the day name a literal | DateScalarExpr over DateKernels: last_day and add_months on the civil-date conversion (add_months clamps the day to the target month as java.time does), next_day with Spark's getNextDateForDayOfWeek and day codes (MO/MON/MONDAY, ...; an unknown or non-literal name falls back), weekofyear as the ISO week of the week-based year |
months_between(a, b[, roundOff]) |
timestamp lanes, or date lanes behind Spark's date -> timestamp cast | MonthsBetweenExpr: Spark's DateTimeUtils.monthsBetween -- whole months when the days of month agree or both are month ends, otherwise the day-and-second difference over a 31-day month, rounded to 8 places when roundOff. Under a UTC or fixed-offset session zone only |
make_date(y, m, d) |
INT32 lanes | MakeDateExpr: an invalid civil date gives null, or under ANSI raises Spark's DATETIME_FIELD_OUT_OF_BOUNDS -- for active rows only, so a filtered row never raises |
unix_date, date_from_unix_date, timestamp_micros, unix_micros |
INT32 / INT64 lanes | RelabelExpr: the same lane under Spark's other type |
date_format(ts, literal pattern), from_unixtime(seconds, literal pattern) |
INT64 micros / seconds, or a date lane behind Spark's date -> timestamp cast -> UTF8 | FormatInstantExpr: Spark's own TimestampFormatter for the session zone, built as DateFormatClass builds it and applied per row, the results written as one lane by the row writer. A timestamp input works under any session zone (the formatter owns the zone rules); a date input needs a UTC or fixed-offset zone, because its instant is our arithmetic. A non-literal pattern falls back |
unix_timestamp(ts | date), to_unix_timestamp(ts | date) |
INT64 micros / INT32 days -> INT64 | UnixTimestampExpr: Spark's truncating division into seconds; a date under a UTC or fixed-offset zone only. Parsing a string input is not ours and falls back |
date_trunc(literal unit, ts) |
INT64 micros -> INT64 | TruncTimestampExpr: Spark's truncTimestamp under a UTC or fixed-offset zone -- MICROSECOND/MILLISECOND/SECOND zone-free, MINUTE/HOUR/DAY a floor in local time, WEEK/MONTH/QUARTER/YEAR via the date kernel's truncation of the local day. An unknown unit (Spark gives null) falls back |
timestamp_seconds, timestamp_millis (integral input), unix_seconds, unix_millis |
INT32 / INT64 lanes -> INT64 | EpochScaleExpr: Spark's multiplyExact (long overflow on an active row) and floorDiv. timestamp_seconds over a double or decimal falls back |
monotonically_increasing_id() |
-> INT64 | MonotonicIdExpr: partitionIndex << 33 (from the task context) plus a running row number carried across the partition's batches -- Spark's contract, and identical values for the same plan since the numbering follows the rows that reach the expression: with a forwarded selection only the selected rows are numbered (SequenceKernels.iotaSelected). Rows outside a CASE branch's active mask are still numbered (Spark would skip them); the id is a tag, and uniqueness and monotonicity within the partition hold either way |
spark_partition_id() |
-> INT32 | SparkPartitionIdExpr: the task's partition index from the task context, read once per task and written with SequenceKernels.fillInt (one Vector API broadcast store per block) -- the value Spark's SparkPartitionID returns |
abs, positive, negative |
INT32, INT64, FLOAT64 | AbsExpr over MathKernels.abs; ANSI mode raises ARITHMETIC_OVERFLOW (integer overflow / long overflow, no hint -- negateExact's message) for an active MIN_VALUE, legacy wraps. positive is the identity, negative is UnaryMinus. abs of a literal, abs over <type> not supported (decimals: no issue yet) |
abs, negative over a wide decimal |
DECIMAL128 | WideDecimalUnaryExpr over WideDecimalCastKernels.abs / negate (#258): two's complement on the limbs; a decimal's range is symmetric, so neither overflows. (Spark itself raises negating a decimal(38,10) at the type's extreme under ANSI when its Decimal holds a rounded form of the value -- see the wide-decimal suite.) |
sign / signum |
FLOAT64 -> FLOAT64 | Spark casts the argument to double first (see Cast); NaN stays NaN, -0.0 stays -0.0 |
sqrt, cbrt, exp, expm1, sin, cos, tan, asin, acos, atan, sinh, cosh, tanh, asinh, acosh, atanh, cot, sec, csc, degrees, radians |
FLOAT64 -> FLOAT64 (Spark casts the argument) | UnaryMathExpr over TranscendentalKernels: one scalar call per lane, exactly the call Spark's generated code makes -- java.lang.Math for the trigonometric and hyperbolic functions, cbrt, the conversions, sqrt; java.lang.StrictMath for exp/expm1; Spark's own formulas for the inverse hyperbolics; sec/csc/cot as 1 / cos, 1 / sin, 1 / tan -- so every lane is bit-identical to Spark and no ulp tolerance enters (the tests compare with tolerance 0). Domain edges are Math's: NaN for asin(2), acosh(0.5), atanh(2); infinities for csc(0), cot(0). The Vector API's EXP/LOG/SIN... operators are 1-2 ulp off StrictMath and are deliberately not used; a measured switch is a separate decision |
ln / log, log2, log10, log1p |
FLOAT64 -> FLOAT64 | StrictMath.log/log10/log1p and log(x) / log(2) as Spark's UnaryLogExpression; null (not NaN) at or below the asymptote -- 0 for the logarithms, -1 for log1p -- by Spark's !(x <= asymptote) test, so NaN stays NaN |
pow / power, atan2, hypot, log(base, x) |
FLOAT64, FLOAT64 (a literal on either side) | BinaryMathExpr: StrictMath.pow (overflow is infinity, pow(0, 0) is 1), Math.atan2(y + 0.0, x + 0.0) with Spark's signed-zero folding, Math.hypot, StrictMath.log(x) / StrictMath.log(base) null when either operand is <= 0 (Spark's Logarithm). Two literals fold before reaching us; pi() and e() are folded constants |
% / mod, pmod |
INT32, INT64, FLOAT64 | DivideLikeExpr over MathKernels.remainder: Java's truncated remainder for %, Spark's r = a % n; r < 0 ? (r + n) % n : r for pmod (int arithmetic wraps like Spark's). A zero divisor never reaches the arithmetic: on an active row ANSI raises Spark's REMAINDER_BY_ZERO, legacy nulls the lane; rows a filter removed or an earlier conjunct decided never raise. Same Spark type on both sides. try_mod (Remainder in TRY mode) nulls a zero divisor in either mode, the legacy path. pmod / div in TRY mode have no SQL function and fall back (try_* <op> not supported). % over <type> not supported (decimals: no issue yet), % operands differ: ... |
div (IntegralDivide) |
INT32 / INT64 operands -> INT64 | Truncation toward zero, long result as Spark's; zero divisor as above but with DIVIDE_BY_ZERO; ANSI raises ARITHMETIC_OVERFLOW (Overflow in integral divide, hint try_divide) for an active Long.MIN_VALUE div -1. div over double not supported, decimals: no issue yet |
greatest, least |
INT32, INT64, FLOAT64 | PickExpr over MathKernels.pick: variadic, every child the same Spark type after Spark's casts, null operands ignored, all-null is null, doubles in Spark's ordering (NaN greatest, -0.0 == 0.0); literal children become constant columns. greatest operands differ: ..., greatest over <type> not supported, greatest of literals only |
nanvl |
FLOAT64 | NanvlExpr over MathKernels.nanvl with Spark's eval semantics: null where the first argument is null, the first argument where it is not NaN (whatever the second is), otherwise the second (null if it is null); the second argument is evaluated only on the rows where the first is NaN (an active mask), as Spark does, so nanvl(c, 1/c) never raises where c is kept. nanvl over <type>, nanvl of two literals |
ceil, floor (one argument) |
INT64 (identity), FLOAT64 -> INT64, decimal(p,s) -> decimal(p-s+1, 0) | CeilFloorExpr over RoundKernels: Spark casts an int argument to long first; a double becomes a long through Java's (long) Math.ceil(d) (NaN is 0, the infinities saturate); a decimal is an unscaled divide by 10^s in CEILING / FLOOR mode to Spark's bounded(p - s + 1, 0) type. ceil over decimal(18,4) -> decimal(19,0) not supported when the widened result passes 18 digits |
rint |
FLOAT64 | RintExpr, Math.rint |
round(x[, k]), bround(x[, k]) |
INT32, INT64, FLOAT64, decimal | RoundExpr -- HALF_UP for round, HALF_EVEN for bround, each as Spark's RoundBase: a double is BigDecimal(Double.toString(d)).setScale(k, mode).doubleValue() (the shortest decimal representation is rounded, so round(2.675, 2) is 2.68; NaN and the infinities pass through; scalar path); an integer with k >= 0 is itself, with k < 0 rounded to a power of ten -- a result outside the type raises ARITHMETIC_OVERFLOW (Overflow, no hint) in ANSI mode for active rows and wraps like BigDecimal.intValue() otherwise; a decimal is an unscaled-value operation to the type Spark computed ((p - s + 1 + min(s, k), min(s, k)), or (max(p - s + 1, -k + 1), 0) for a negative k). k must be an int literal (Spark requires it foldable). round over decimal(18,4) -> decimal(19,4) not supported when Spark's widened result passes 18 digits |
ceil(x, k), floor(x, k) |
decimal (int arguments arrive cast to decimal(10,0)) | RoundExpr in CEILING / FLOOR mode (RoundCeil / RoundFloor), same result-type rule as round. Spark casts a double argument to decimal(30,15) first, which stays a fallback (cast ... not supported; no issue yet) |
&, |, ^, ~ |
INT32, INT64 | BitBinaryExpr / BitNotExpr over BitKernels, Vector API lanewise ops; both operands the same integral type after Spark's coercion; a literal on either side. Bytes, shorts and booleans have no lane here (& over tinyint, tinyint not supported); boolean and/or are the logical path above |
shiftleft, shiftright, shiftrightunsigned (<<, >>, >>>) |
INT32, INT64 value, INT32 amount | BitBinaryExpr with Java's (= Spark's) semantics: the amount is masked by the lane width (x << 33 on an int is x << 1, negative amounts wrap the same way), >>> on INT32 is computed in 32 bits, the amount may be a literal or a column, and a literal value shifted by a column (1 << i) is handled too |
bit_count |
INT32, INT64 -> INT32 | BitCountExpr: java.lang.Long.bitCount of the value widened to a long, exactly as Spark computes it -- bit_count(CAST(-1 AS INT)) is 64, not 32. bit_count over boolean not supported (booleans have no INT lane to widen) |
bit_get / getbit |
fallback | Returns TINYINT, which has no lane in this project (no INT8): bit_get returns tinyint, which has no lane |
UnscaledValue, MakeDecimal |
INT64 (decimal(<=18)) | The optimizer's DecimalAggregates rewrite of sum(decimal(p <= 8)) and avg(decimal(p <= 11)), pinned end to end (every aggregate stage ours); MakeDecimal into more than 18 digits falls back (make_decimal into <type> exceeds 18 digits), overflow nulls or raises per nullOnOverflow |
CheckOverflow(child, decimal(p, s), nullOnOverflow) |
INT64 (decimal(<=18)) -> INT64 | Spark's toPrecision(p, s, HALF_UP): the identity when the child already has the declared type (a decimal lane holds values within its precision by construction), otherwise the decimal cast's rescale-and-check -- nullOnOverflow nulls the row, otherwise Spark's NUMERIC_VALUE_OUT_OF_RANGE for an active row only. Note: Spark 4.1's analyzer and optimizer leave no CheckOverflow in a batch plan -- the decimal operators check their own overflow (DecimalArithExpr mirrors that) and the aggregate form is CheckOverflowInSum (see sum); the node is constructed only by Dataset encoders, streaming joins and table-output coercion, so this case exists for completeness and is tested by building the node directly. More than 18 digits either side: check_overflow ... exceeds 18 digits |
KnownFloatingPointNormalized |
any | Identity marker |
NormalizeNaNAndZero |
FLOAT64 | A real pass (MathKernels.normalizeNaNAndZero): every NaN becomes the canonical NaN and -0.0 becomes 0.0, so the bit comparison of the group key table and the join tables agrees with Spark's equality on the double grouping and join keys the optimizer wraps in it (in plain comparisons the compare kernels already use that ordering). Over any other type: NormalizeNaNAndZero over <type> not supported |
col.field, col.a.b (GetStructField over a struct column or a chain of struct fields) |
the field's lane type: INT32, INT64, FLOAT64, BOOL, UTF8 (incl. date, timestamp, decimal(<=18)) | StructFieldExpr: the struct column reaches the operator as Spark's own vector (#19) and its fields are that vector's children (ColumnVector.getChild), so the field is read by walking the chain and adapting the leaf like an input column -- borrowed from off-heap memory, copied from on-heap, nothing assembled row by row; every ancestor struct's nulls are folded into the validity, as Spark's GetStructField is null when the struct is. A field whose type is itself a struct, array or map, projected as a value (st.c AS c, or Spark's _extract_inner alias below a generate), is passed through by the project as a view of the child vector with the struct's nulls folded in (NestedFieldColumnVector), like a whole nested column (#19); inside any other expression such a field is refused (not supported as a value); arr[i] / map[key] / arr.field are refused with array element access / map value access / array of struct fields ... not supported (#50). |
size(arr), size(map), cardinality (Size) over a column without a lane or a struct field of one |
INT32 out; the array/map read from Spark's vector | SizeExpr: the element count per row; a null array is null (or -1 under spark.sql.legacy.sizeOfNull). Spark adds size(arr) > 0 below a non-outer explode, so this keeps the generate's chain columnar. |
IS NULL, IS NOT NULL over a column without a lane or a struct field of one |
BOOL | NestedValidityExpr: the non-null rows of Spark's vector, ancestors' nulls included, as the tested validity. |
ScalarSubquery -- and any reference-free expression over one: the GetStructField Spark's MergeScalarSubqueries leaves over a struct-valued subquery, a CASE over such fields, arithmetic with one |
the result's type: INT32, INT64, FLOAT64, date, decimal(<=18), BOOL, UTF8 | SubqueryLiteralExpr: a literal by execution time. Spark runs the subquery before the operator (waitForSubqueries) and stores the result in the node, which travels to the executors with the expression as in Spark's interpreted path; read on first use, materialised as a constant lane (null result: a null lane). In filters, projections and aggregate inputs; subquery result of type <t> not supported otherwise. A merging aggregate whose function instance was rewritten between stages (a ReusedSubquery in its input) binds its buffers by position, as Spark does, where exprIds no longer agree |
xxhash64(c1, ..., cN[, seed]) |
INT32, INT64, FLOAT64 (incl. date, timestamp, decimal(<=18)), BOOL, UTF8 | XxHash64Expr: Spark's XXH64 steps per lane in Spark's order and type rules (hashInt for ints/dates, hashLong for longs/timestamps/unscaled decimals, hashLong(doubleToLongBits) with -0.0 normalised, `hashInt(1 |
hash(c1, ..., cN) (Murmur3, seed 42) |
the xxhash64 lane types |
Murmur3HashExpr: the same per-child chain calling Spark's own Murmur3_x86_32 steps -- a null child leaves the running value alone, -0.0 is normalised, ints/dates hash as ints, longs/timestamps/short decimals as longs, booleans as 1/0, strings as their bytes -- so the value is Spark's by construction (it decides shuffle partitions) |
md5, sha / sha1, sha2(str, bits), crc32 |
UTF8 lane through Spark's string -> binary cast (a reinterpretation of the bytes); sha2's bit length a literal |
DigestExpr: a per-row MessageDigest / CRC32 loop inside the operator with lowercase hex out through the row writer -- not vectorisable, present so a query containing one keeps the operator. sha2 accepts 224/256/384/512 and 0 (= 256); another bit length gives null, a non-literal one falls back. A BINARY column has no lane |
BloomFilterMightContain(filter, xxhash64(key)) -- Spark's runtime join filter (spark.sql.optimizer.runtime.bloomFilter.*) |
the probe value INT64 | BloomProbeExpr: the filter is a binary scalar subquery read as bytes and deserialised once through Spark's own BloomFilter; mightContainLong per lane; a null filter gives null lanes as Spark's does. The probe sits in the filter directly above the large side's scan, which is now ours; the build side's bloom_filter_agg is an ObjectHashAggregateExec which is now ours too (#57), byte-identical to Spark's |
CASE WHEN ... THEN ... [ELSE ...] END |
result of any lane; conditions BOOL | CaseWhenExpr over SelectKernels: each condition is evaluated only on the rows no earlier branch took, its winning rows are condition is true (a null condition counts as false), each branch value only on its winning rows; the result is null where the winner is null or no branch matched and there is no ELSE. Branches must share the result's Spark type (branch type <t> differs from <t>); NULL and literal branches -- string and boolean literals included -- are materialised as constant columns. unsupported result type <type> for <sql> for a wide decimal or nested result |
IF(c, a, b) |
as CASE WHEN |
The one-branch case with an ELSE |
COALESCE(a, b, ...), NVL, NVL2, NULLIF, IFNULL |
as CASE WHEN |
COALESCE is IS NOT NULL conditions over the operands with the last as ELSE (an operand is evaluated once for its test and once for its value); NVL/NVL2/NULLIF/IFNULL arrive as COALESCE/IF through Spark's own rewrites |
Aggregate functions
Compiled by VectorAggregates for HashAggregateExec in every mode (see
docs/operators.md for the operator's own conditions).
| Function | Types | Notes / fallback reasons |
|---|---|---|
count(*), count(x) |
any lane | count with several arguments not supported |
sum(a * b), sum(a + b), sum(a - b) and their nestings, under sum or avg, where the decimal arithmetic's declared result exceeds 18 digits (sum(l_extendedprice * (1 - l_discount)): decimal(12,2) * decimal(14,2) is decimal(27,4); sum(x + 1) over decimal(18,4) is decimal(19,4)) |
INT64 operands (decimal(<=18) each) or speculative expressions themselves; for * the declared scale exactly s1 + s2 and precision at most 38, for +/- the declared scale exactly max(s1, s2) (Spark's precision cap lowers the scale -- and rounds -- only past 32 integer digits; that shape is refused); not try_* |
SpeculativeDecimalMulExpr / SpeculativeDecimalAddExpr (#26): Spark's result precision is a static rule, not a statement about the data, so the product is computed in the INT64 unscaled lane and checked per row with Math.multiplyHigh; a row whose product leaves 64 bits is escalated -- its exact BigInteger product is added to the 128-bit sum accumulator beside the lane (WideDecimalSumAgg.Escalation), so the total is exact whatever the data and no wrong value ever leaves. Spark never rounds or raises on such a product (the declared precision holds every product of the operands' widths), so the two engines agree row for row, ANSI or not; a total past the sum's own declared precision is the sum's existing null / NUMERIC_VALUE_OUT_OF_RANGE. An operand may itself be such a product ((price * (1 - disc)) * (1 + tax), declared decimal(38,6)): its escalated rows are multiplied exactly, so the whole tree escalates as one; where Spark capped the declared precision at 38 a product may not fit it, and like Spark's Multiply that row is null in legacy mode and NUMERIC_VALUE_OUT_OF_RANGE in ANSI mode (only an escalated value can be that large). +/- rescale both operands to the common scale in the lane (the shift by 10^(s - si) checked with Math.multiplyHigh) and add with a sign-trick overflow check; a row that leaves 64 bits at either step is escalated exactly, and an operand's own escalated rows are combined exactly -- (a * b) + c, a * b - c * d, (a * b + 1) * c all escalate as one tree. An uncapped +/- result can never exceed its declared precision (one more digit than the wider operand); the capped one is checked like the product. Only the wide sum and average consume the speculative form: a wide product or sum anywhere else (a projection, max) keeps the exceeds 18 digits refusal, and a wide operand that is neither a lane nor a speculative expression (a wide CAST) is refused with neither a lane nor a speculative product or sum. SpeculativeDecimals.escalatedRows() counts escalations (an operator metric is a follow-up). |
sum |
INT32, INT64, FLOAT64, decimal | ANSI sum(bigint) is overflow-checked (a sign-trick overflow lane, Math.addExact on the grouped path). try_sum over an integral input is Spark's Sum in TRY mode (TrySumLongAgg / TrySumLongMergeAgg): the (sum: bigint, isEmpty) buffer, rows added exactly, and the first overflow poisons the whole group -- its sum is null from then on through every later row and merge, as Spark's Add(sum, child, TRY) leaves a null no later add revives; Spark's result If(isEmpty, null, sum) compiles over the emitted buffer. try_sum over doubles is the plain double sum (Spark tracks no isEmpty and doubles never overflow); over a decimal it falls back (try_sum over a decimal not supported; no issue yet). sum(decimal(p <= 8)) arrives as MakeDecimal(sum(UnscaledValue)) and is supported end to end. A wider decimal (decimal(12,2), decimal(18,4)) sums in the update modes into a 128-bit accumulator per group and emits Spark's (sum: decimal(p+10, s), isEmpty) buffer, the sum as a wide Arrow decimal column; the merge modes combine those buffers with Spark's rules (isEmpty && other.isEmpty, a null sum on a non-empty row means an earlier overflow) reading the wide column row by row into an exact per-group total -- a merge sees one buffer row per partition per group -- and Final applies If(isEmpty, null, CheckOverflowInSum(sum)) at emission: null when empty, Spark's overflow error in ANSI mode (ARITHMETIC_OVERFLOW for an overflowed buffer, NUMERIC_VALUE_OUT_OF_RANGE for a total past the declared precision) or null otherwise. The wide buffer and that result are the only places a decimal above 18 digits enters or leaves this plugin; any other expression over it falls back. sum over <type> producing <type> not supported for a double into a decimal. Double sums add in lane-parallel, interleaved order (vecruntime.agg.interleave, default 4): results can differ from Spark's in the last bits; interleave=1 reproduces Spark's rounding |
min, max |
INT32, INT64, FLOAT64 (incl. date, timestamp, decimal(<=18)), BOOL, UTF8 | Numeric lanes through the kernels' accumulators; booleans and strings through a comparison per row (OrderedMinMaxAgg, Spark's binary order for strings). A string buffer makes Spark plan a SortAggregateExec, which is converted too (see the operator matrix). bool_and / bool_or / every / any / some arrive as min / max over a boolean after Spark's rewrite |
count_if |
as count |
Arrives as count over Spark's rewrite of the condition (nullif(cond, false), a CASE WHEN over booleans) |
bit_and, bit_or, bit_xor |
INT32, INT64 | One operation per row (BitAgg), null over no rows; the merge is the same operation |
last(x[, ignoreNulls]) / last_value |
INT32, INT64, FLOAT64, date, boolean, decimal(<=18), UTF8 | LastAgg with Spark's (last, valueSet) buffer: the last row in partition order, the last non-null one with ignoreNulls; order-dependent across the shuffle exactly as Spark's own |
max_by, min_by |
value: INT32, INT64, FLOAT64, date, boolean, decimal(<=18), UTF8; ordering: any of those | MaxMinByAgg, Spark's (value, ordering) buffer; rows with a null ordering are skipped and a tie takes the later row, as Spark's predicate keeps the held value only when strictly better |
avg, try_avg |
INT32, INT64, FLOAT64 producing double; decimal: p <= 11 through the optimizer's rewrite to a double average, wider through Spark's own (sum: decimal(p+10, s), count) buffer |
try_avg is avg: Spark sums an average over integrals in doubles, which never overflow, so TRY changes nothing outside decimals. A decimal beyond 11 digits (avg(l_quantity) over decimal(15,2): buffer decimal(25,2), result decimal(19,6)) sums in the update modes into the 128-bit accumulator per group beside a count and emits Spark's buffer, the sum as a wide Arrow decimal column (a total past the buffer precision is null, as Spark's buffer writer leaves it); the merge modes add the wide column row by row and the counts, a null sum poisoning the group as Spark's DecimalAddNoOverflowCheck does; Final emits the result in the sum slot as Spark's own expression evaluated over the merged (sum, count) -- If(count = 0, null, DecimalDivideWithOverflowCheck(sum, count, decimal(p+4, s+4))), so the 39-digit half-up quotient, the rounding to the result scale, the null / NUMERIC_VALUE_OUT_OF_RANGE past the result precision and the ARITHMETIC_OVERFLOW on a buffer that overflowed are Spark's to the last digit -- and the projection forwards it (a result wider than 18 digits is a pass-through column above the aggregate, avg(x) * 2 over it falls back). An ungrouped Final divides the exact total even past the buffer precision, as Spark's generated code does (its sum lives in a local no writer re-checks); a grouped one sees the buffer nulled. try_avg over a decimal is the same function with nullOnOverflow. avg(a * b) with a declared-wide product is speculative like sum(a * b) below (#26) |
first(x[, ignoreNulls]) / first_value |
INT32, INT64, FLOAT64, date, boolean, decimal(<=18), UTF8 | FirstAgg / FirstMergeAgg with Spark's (first, valueSet) buffer: the first row in row order per group, the first non-null one with ignoreNulls (a null first value is a real value without it); across the shuffle Spark's own result is order-dependent too. Spark's distinct rewrite wraps every plain aggregate in the ignoreNulls form |
any of the above with FILTER (WHERE ...) |
the predicate must compile as a boolean | FilteredAgg: applied in the update modes only (Partial, Complete) as a narrowed selection for that function -- Spark drops the clause in the merge modes. FILTER <sql>: <reason> when the predicate does not compile |
any of the above with DISTINCT |
as the function | The isDistinct flag is only a marker in a physical plan: Spark's planAggregateWithOneDistinct (one distinct group) groups by the distinct column in the two inner stages and the flagged function runs over deduplicated input as a plain one; several distinct groups go through Expand with a keys-only first aggregate, FILTER (WHERE gid = k) on every function and first(..., true) around the plain ones -- all of which compile. A distinct with FILTER folds its condition with max over a boolean in the first aggregate, which falls back (min/max over boolean not supported, #45) |
stddev / stddev_samp / std, stddev_pop, variance / var_samp, var_pop, skewness, kurtosis |
FLOAT64 (Spark casts the argument) | MomentsAgg: Spark's own streaming Welford step (CentralMomentAgg: n, avg, m2, and m3 / m4 for the higher moments) applied row by row in partition order, so a Partial over the same rows produces Spark's doubles; the merging modes apply Spark's mergeExpressions (the parallel form) to the incoming buffers. Buffers are emitted as Spark's doubles and the result is Spark's own expression compiled over them: null over no rows, the sample forms null over one row, skewness / kurtosis null over a constant, NaN propagates. Agreement is to a relative tolerance, not bit equality -- the merge order across the shuffle is the shuffle's, as for Spark's own two runs |
covar_samp, covar_pop, corr, regr_sxy, regr_r2, regr_slope, regr_intercept |
FLOAT64, FLOAT64 | The same over Spark's Covariance (n, xAvg, yAvg, ck), PearsonCorrelation (+ xMk, yMk) and the (covariance, variance of x) pair of regr_slope / regr_intercept; a pair with a null is skipped. corr over a constant column divides by sqrt(0): Spark's own evaluation raises DIVIDE_BY_ZERO under ANSI and so does ours |
regr_count, regr_avgx, regr_avgy, regr_sxx, regr_syy |
as count / avg / the moments |
Spark's rewrites: count(y, x) (rows with both non-null, CountAllAgg), avg(if(...)), and RegrReplacement (a variance state whose result is m2) |
collect_list(x), collect_set(x) |
element: any lane type (int/bigint/double/date/timestamp/decimal(<=18)/boolean/string); output array<element> |
SparkObjectAgg (#57), planned by Spark as an ObjectHashAggregateExec: driven through Spark's own CollectList / CollectSet object per group, so nulls are ignored, collect_set dedups, and the partial buffer (serialize, an UnsafeProjection over the collected array) interoperates with a Spark merge unchanged. The array<T> output is Spark's GenericArrayData held in an on-heap column vector. The collected order is non-deterministic after a shuffle, exactly as Spark's; a wide-decimal or non-lane element type falls back through the operator's columnar-child gate. Held in memory (no spill) |
bloom_filter_agg(x[, numItems[, numBits]]) |
INT64 child (the runtime filter feeds xxhash64(key)); output binary |
SparkObjectAgg (#57): Spark's own BloomFilterAggregate object per group -- BloomFilter.create(numItems, numBits) from the folded literal expressions, putLong per row (the child's own eval), mergeInPlace for the merge -- so the partial buffer is byte-identical to BloomFilterImpl.writeTo and a Spark Final or the runtime probe (might_contain, BloomProbeExpr) reads it unchanged. The Final emits the serialized filter, or null when no bit is set, as Spark's eval does. This is the build side of the runtime join filter whose probe was already ours |
| any other function | -- | unsupported aggregate function <Class>: <sql> (percentile_*, collect_top_k, listagg and the other TypedImperativeAggregates Spark plans as ObjectHashAggregateExec are not converted) |
aggregate over a literal, aggregate over <type> not supported (non-numeric input) and
literal grouping key / grouping key type <type> not supported (unsupported types; doubles group by bits after their normalisation pass) are
the remaining reasons on the aggregate's inputs.
Planned
| Family | Issue |
|---|---|
LIKE with _ wildcards or escape characters, rlike |
no issue yet (the LikeSimplification shapes and every %-only pattern are Supported above, #264) |
Decimal results wider than 18 digits on narrow operands beyond the speculative * + - trees under sum / avg above: a wide product or sum as a projected value, escalation with a wide output column (128-bit DecimalVector), +/- where Spark's precision cap lowers the scale (it rounds there) |
#26 (slices 1-4 are the aggregate shapes; the rest is #28's wide lane) |
Genuinely wide declared decimals (p > 18) as operands of the conditional (CASE WHEN / IF with a wide result), the round family, % and as a scalar subquery's value -- the lane exists (#257--#259); these kernels do not yet |
follow-up to #258 |
bit_get / getbit (needs an INT8 lane) |
#36 |
Array element and map value access (arr[i], map[key], arr.field over an array of structs); struct fields are Supported above |
#50 (only if profiling says nested reads matter; a constant index over a list is the cheap case) |
| A row-based escape hatch for expressions with no kernel, including Scala UDFs | #51 |
Not planned
Recorded here so that an absence is a decision rather than an omission (#66). Open to revisiting on demand, not permanent exclusions. Comet's list is the reference for what a columnar engine reasonably leaves out.
| Family | Reason |
|---|---|
Probabilistic sketches (approx_count_distinct / HyperLogLog++, approx_percentile, count_min_sketch, hll_* functions) |
Sketch state is a Spark-defined binary buffer with its own merge semantics; the buffers are the compatibility surface, not the arithmetic, and Spark's implementation is already tight |
Nested results -- building a struct, array or map (named_struct, struct, array, map, map_from_arrays, ...), and the array_funcs, map_funcs and lambda (transform, filter, aggregate, exists, ...) families |
No nested lane: the engine reads into nested columns (struct fields, #50) and passes whole nested columns through (#19), but never assembles one (#50) |
| Geospatial functions | Outside the project's data model; no lane type and no plan to add one |
Avro / Protobuf codecs (from_avro, to_avro, from_protobuf, to_protobuf) |
Schema-driven record (de)serialisation, row by row by construction |
JVM reflection (java_method, reflect) |
Calls arbitrary JVM code per row; nothing to vectorise |
Niche string validators and encoders (is_valid_utf8, make_valid_utf8, validate_utf8, try_validate_utf8, luhn_check, soundex, levenshtein) |
Byte-level per-value algorithms with little SIMD upside; Spark's implementations are adequate. The common string functions are planned (#37-#41) |
Zone-dependent datetime arithmetic under a session zone with rules (DST): cast(timestamp AS date), hour/minute/second, months_between, date_trunc, a date's instant for date_format/unix_timestamp |
The kernels convert with one fixed offset; a zone with transitions needs a per-row rules lookup, which is what Spark already does. Only date_format over a timestamp runs Spark's own formatter and so takes any zone. Decided under #44 |
try_to_number, try_to_binary |
The to_number / to_binary formats they wrap are not kernels (format-string parsing per row); they follow whenever those land. Decided under #47 |
| Pickled (non-Arrow) Python UDFs, Scala UDFs as kernels | The data has to become objects row by row. The row-based escape hatch (#51) is the path for these, not a kernel |
Keeping this page honest
Every expression issue names this file as part of its definition of done: the row lands in the same
commit as the kernel, the scalar reference and the Spark comparison test. The reason strings above are
quoted from ExpressionCompiler, VectorAggregates and VectorHashAggregateExec; a string that
changes without its row changing fails a checkFallback assertion in the suites before it reaches a
user. The intended end state is to generate the Supported table from the compiler itself -- walking
Spark's FunctionRegistry against a dummy schema and recording each Left(reason) -- so that a
regression shows up as a diff of this file rather than as a stale row.