Spark operator support
One row per Spark physical operator: whether VecRuntime converts it, what has to hold for the
conversion to happen, the spark.vecruntime.* key that turns it off, and the exact fallback reason the
planner records when it does not convert. Modelled on Comet's
Spark Operator Support page;
the companion for expressions is docs/expressions.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 a plan. VectorExecRule (VectorColumnarRule.scala) walks the physical plan bottom-up
and replaces an operator only when (a) its input already produces columnar batches of supported types
-- a vectorized Parquet scan, a Comet scan, an Iceberg scan through the adapter, or another
VecRuntime operator -- and (b) every expression it carries compiles. Anything else is left as
Spark's operator with the reason attached as a VectorFallback tag. There is no silent fallback: the
Vector Acceleration UI tab colours every node by the engine that runs it and shows the reason on
hover, spark.vecruntime.explainFallback.enabled=true logs each reason at INFO, the TPC-H runner records
them per query, and the test suites assert on the strings (checkFallback(..., reasonContains)). The
strings in the Fallback reasons column are therefore a contract: change one and the suites and
this page change with it.
Supported types are the VecType lanes the kernels compute on: boolean; int and date
(INT32); bigint, timestamp and decimal(p <= 18) as the unscaled value (INT64); double
(FLOAT64); string (UTF8). A sixth lane, DECIMAL128, holds decimal(19 <= p <= 38) as two
little-endian 64-bit limbs (Arrow's Decimal128 layout, scalar only -- no kernel uses a vector
species on it; #28, option 1). The lane is carried (#257: a wide decimal column is adapted from
Spark's byte strings, compacted by the filter and forwarded by the projection as an Arrow
Decimal128 vector), computed on by the expression kernels (#258: comparisons, IN, + - * /,
casts, abs, negation) and read by the hash aggregate as grouping keys, sum / avg / min /
max / count / first / last inputs and arithmetic over wide sums in the result projection
and carried by every mover, both hash joins and the window as keys, payloads and frames (#259, the
rows below name the one frame shape that still falls back). An operator whose input or output carries any other type -- float, short,
byte, binary, arrays, maps, structs -- records unsupported column type <type> for <column> (input)
or unsupported output type <type> for <name> (a projection or aggregate result), whatever else is
true of it. The two exceptions are the operators that forward columns without reading them: a filter
and a projection pass such a column through as Spark's own vector (borrowed, or viewed through a row-id
mapping when a selection is applied), so only an expression that reads it is refused (#19). This is the single most common reason in real plans.
spark.vecruntime.enabled=false turns the whole rule off. Every converted operator has its own key,
listed below; all default to true.
Supported
| Spark operator | VecRuntime operator | Requirements | Config key | Fallback reasons |
|---|---|---|---|---|
FilterExec |
VectorFilterExec |
Columnar child; the condition compiles (see the expression matrix). Columns of types without a lane (struct, array, map, ...) are forwarded as Spark's vectors -- remapped rather than compacted when rows are dropped -- so they never disqualify the filter; a wide decimal has a lane (DECIMAL128) and is compacted like any other column; only a condition reading one is refused. Between two VecRuntime operators the filter forwards a selection bitmap instead of compacting (spark.vecruntime.exec.selection.enabled). |
spark.vecruntime.exec.filter.enabled |
child <op> is not columnar; any expression-compiler reason, e.g. unsupported expression Like: ..., unsupported type struct<...> for st |
FileSourceScanExec / BatchScanExec (Spark's vectorized Parquet scan, Iceberg's JVM reader) |
VectorPrefetchScanExec (#403, lever 2) |
Off by default. Inserted between a columnar Spark file scan whose every column has a lane and the first operator of ours above it (a filter, a projection, an aggregate, a join side, a sort, an expand). Not an operator: a per-task helper thread pulls the reader's batches, converts every column into our Arrow vectors (ColumnVectorAdapters.adapt then ArrowOutput.copy; a row-id-mapped Iceberg batch is normalized and its live rows compacted) and hands them through a bounded queue, so the reader's waits overlap our kernels and the adapters above read zero-copy. The reader's vectors are recycled per batch, so a batch is converted whole before the next is requested; a task kill or the task's end stops the helper and releases every queued batch. Never wraps a Comet scan (read zero-copy already), our own operators or a scan with a column without a lane. Counted as plumbing (a conversion) by the acceleration view |
spark.vecruntime.scan.prefetch (queue depth; 0 off) |
-- (never a fallback: without it the scan is read as before) |
ProjectExec |
VectorProjectExec |
Columnar child; every computed expression compiles and has a supported result type; a bare literal slot of any lane type -- a boolean or a typed NULL included -- is a constant or null column (#273). A bare column reference (or a rename of one) of any type is passed through: struct, array and map columns ride beside computed ones as Spark's own vectors, wide decimals as DECIMAL128 lanes, borrowed when the batch is dense and viewed through a row-id mapping when a selection from the filter below is applied. A projection may be a bare literal (SELECT 1 FROM ...). Operators above still refuse such a column by type; ColumnarToRowExec reads it as any other. |
spark.vecruntime.exec.project.enabled |
child <op> is not columnar; <expr>: <compiler reason> (an accessor into a nested column, st.a, is one -- #50); <expr>: unsupported output type <type> for <name> (several failures are joined with ; ) |
MergeRowsExec |
VectorMergeRowsExec |
Columnar child (the merge's join, ours since #273: the struct _partition rides the streamed side); a lane-less child column is forwarded. The row-level operator of a MERGE INTO (#21): the two presence predicates become the matched / not matched / not matched by source group masks; within a group each clause's condition is evaluated over the rows the earlier clauses left, so clause order is Spark's; a Keep compacts its projection by the clause mask (one output batch per projection), a Split (an update as delete + insert) two, a Discard none, and rows no clause takes are dropped. Unconditional clauses (true) take every remaining row. Lane-less output columns (Iceberg's struct _partition) are passed through by row id. The cardinality check keeps the matched __row_ids in one roaring bitmap per partition and raises MERGE_CARDINALITY_VIOLATION before the batch is emitted. Output order is grouped by clause within a batch rather than interleaved in join order, which the delta writer does not depend on. spark.vecruntime.exec.mergeRows.enabled. Reasons: child <join> is not columnar (today always: the merge's join carries the struct _partition, which our joins do not pass through -- the prerequisite for the operator to fire on a real Iceberg merge), merge clause condition: <reason>, <expr>: <reason>, merge cardinality check without a __row_id column. |
||
HashAggregateExec (Partial, Complete) |
VectorHashAggregateExec |
The update modes: the functions read their input, Partial emits the aggregation buffers and Complete evaluates the result expressions in the same operator (Spark 4.1's batch planner never emits Complete; the streaming planners do). Columnar child. Functions: count, sum, min, max, avg over numeric lanes (a decimal sum or avg past 18 digits accumulates in 128 bits -- from the INT64 lane or, for a wide input, the two limbs of the DECIMAL128 lane with a sticky overflow flag (#259) -- and emits Spark's wide (sum, isEmpty) / (sum, count) buffer; a Final decimal avg emits Spark's own result expression over the merged buffer, wide or not), min / max / first / last over the wide lane too (two-limb compares, #259); grouping keys of type int/bigint/boolean/string/date/timestamp/decimal (a wide decimal key is hashed and compared on both limbs, #259) and never a literal; FILTER clauses apply here as a narrowed selection per function; DISTINCT is a marker only (Spark's rewrites group by the distinct column below). A keys-only aggregate (SELECT DISTINCT, UNION, the first aggregate of the distinct rewrite) is accepted and emits its keys. Result expressions must be plain attributes at Partial. Double sum/avg round like Spark (bit-identical) unless spark.vecruntime.exec.strictFloatingPoint=false selects the faster lane-parallel, interleaved sums. |
spark.vecruntime.exec.aggregate.enabled, spark.vecruntime.exec.strictFloatingPoint |
child <op> is not columnar; unsupported column type ...; distinct aggregates not supported; FILTER <sql>: <reason>; aggregation modes <modes> mix buffer and result output; aggregate without keys or functions; first over <type> [without ignoreNulls] not supported; unsupported aggregate function <Class>: <expr>; sum buffer <type> exceeds 18 digits; avg buffer <type> within 18 digits not supported; sum over <type> producing <type> not supported; avg producing <type> not supported; count with several arguments not supported; min/max over <type> not supported; aggregate over <type> not supported; aggregate over a literal; grouping key type <type> not supported; literal grouping key; result expression <expr> is not a plain attribute; unsupported result type ...; unsupported output type ... |
HashAggregateExec (PartialMerge, Final) |
VectorHashAggregateExec |
The merge modes: each function's state is merged from the child's buffer columns (inputAggBufferAttributes), PartialMerge emits the buffers again and Final evaluates the result expressions with evaluateExpression substituted. Both read an exchange, so only the types of the input matter: Spark inserts RowToColumnarExec under us over its row shuffle, and Comet's columnar shuffle is read directly. Spark's count(distinct) rewrite produces a PartialMerge stage (ours) and a stage mixing PartialMerge with a distinct Partial (falls back on the distinct function, not on the mix -- update and merge are decided per aggregate expression, buffers vs results per operator). FILTER never applies in the merge modes. A wide decimal sum buffer (Decimal(p > 18)) is read row by row from the shuffle's column; its result (the same wide type) is emitted ready-made by the merge and forwarded by the result projection, or stands in for the sum inside a larger result expression (sum(x) / 7.0, sum(a) - sum(b)) that the wide kernels then compute (#259). Wide keys and min / max / first / last buffers arrive as DECIMAL128 lanes. |
spark.vecruntime.exec.aggregate.enabled, spark.vecruntime.exec.aggregate.final.enabled, spark.vecruntime.exec.strictFloatingPoint |
As above, plus merging aggregation stages disabled by configuration; merging sum buffers of <type> not supported; sum with a <n>-column buffer not supported |
SortExec |
VectorSortExec |
Columnar child only -- over Spark's row shuffle the sort stays Spark's (converting rows to sort them gains nothing), so a global ORDER BY is ours only over Comet's shuffle or a local SORT BY above our operators. Every sort key compiles, is not a literal and has a lane -- a bare wide decimal column included (the DECIMAL128 lane is ordered limb by limb, and wide payloads are gathered like any other; #257). The permutation comes from SortKernels.sortIndices, LSD radix passes over 8-bit digits since #285 (uniform digits and already-ordered passes skipped); the partition is sorted in runs of spark.vecruntime.sort.runRows rows (default 1M) sealed as it arrives and k-way merged on output by RunMerge (a loser tree over one cursor per run, fixed-width keys compared through per-run normalised arrays, ties by run then position -- stable; a block of the winner's rows below the runner-up's key leaves in one step, which is what makes presorted and low-cardinality input cheaper in runs than in one sort), a dictionary column decoded once per run. Runs past spark.vecruntime.sort.spillBytes (default 1 GiB, #416) spill to local disk as Arrow IPC and merge back on output. |
spark.vecruntime.exec.sort.enabled |
child <op> is not columnar; sort without keys; <key>: literal sort key; <key>: sort key type <type> not supported; <key>: <compiler reason> |
TakeOrderedAndProjectExec (ORDER BY ... LIMIT n) |
VectorTakeOrderedAndProjectExec |
Columnar child only; the output may be any lane type, wide decimals included (#259). Per partition the sort iterator runs with a limit: each sorted run keeps only its first n rows for the merge (top-N by run, #285) and only the first n merged rows are gathered; those at most n rows per partition go as rows through Spark's own single-partition shuffle, Spark's ordering takes the final top n, Spark's projection applies the select list and the result is one columnar batch (on-heap vectors) -- the row stages see at most n x partitions rows. Same key rules as the sort; every output type must be supported; OFFSET is not. In memory per partition like the sort (#12). |
spark.vecruntime.exec.takeOrdered.enabled |
child <op> is not columnar; offset <k> not supported; sort without keys; <key>: <sort key reason>; unsupported output type <type> for <name> |
LocalLimitExec |
VectorLocalLimitExec |
Columnar child only, any lane type (#259). Whole batches pass through untouched; the batch that crosses the boundary is compacted to its first surviving live rows (selection-aware); no further batch is pulled once the limit is reached (the child's release runs from its task-completion listener). | spark.vecruntime.exec.limit.enabled |
child <op> is not columnar; negative limit <n> |
GlobalLimitExec (no offset) |
VectorGlobalLimitExec |
As the local limit over the single partition Spark's AllTuples requirement provides -- so only over a columnar exchange (Comet) or a single-partition columnar child; above Spark's row shuffle it stays Spark's. |
spark.vecruntime.exec.limit.enabled |
child <op> is not columnar; offset <k> not supported; negative limit <n> |
CollectLimitExec (LIMIT n at the top of a query, no offset) |
VectorCollectLimitExec |
Columnar child only, any lane type (#259). The first n rows of every partition stay columnar; those at most n rows per partition go as rows through Spark's single-partition shuffle, the first n are taken and re-materialised as one columnar batch (VectorRowStages, shared with the top-N operator). Every output type must be supported. LIMIT 0 never reaches us (folded to an empty relation); adaptive execution may drop a limit it proves a no-op. |
spark.vecruntime.exec.limit.enabled |
child <op> is not columnar; offset <k> not supported; negative limit <n>; unsupported output type <type> for <name> |
UnionExec |
VectorUnionExec |
At least one child columnar; the children's column types do not matter (the union forwards batches and reads nothing), and refusing on them was harmful: a union left to Spark over columnar children runs Spark 4.1.3's concatenating columnar union, and q66 returned every row twice once its channel aggregates emitted a wide decimal sum (#162). Spark's own union already runs columnar when every child is; ours is columnar whatever the children are, so a union with a row child (VALUES, a row shuffle) keeps the chain and Spark's transitions convert that child through RowToColumnarExec below us. No computation: the children's batches, with Spark's partitioning contract -- when every child is hash-partitioned alike (spark.sql.unionOutputPartitioning) the union reports that partitioning, EnsureRequirements plans no shuffle above it, and the i-th partitions of the children are read together (SQLPartitioningAwareUnionRDD); otherwise the RDDs one after the other. Output nullability merged as Spark's. A filter directly below compacts rather than forwarding a selection (pass-through operator). Note: Spark 4.1.3's own UnionExec reports the same partitioning but concatenates on its columnar path (fixed upstream later), so disabling this operator over co-partitioned columnar children -- three aggregated channels under UNION ALL, TPC-DS q33/q56/q60 -- returns one row per child for a key. |
spark.vecruntime.exec.union.enabled |
no columnar child; union with fewer than two children; unsupported column type <type> for <name> |
CoalesceExec (coalesce(n), /*+ COALESCE */) |
VectorCoalesceExec |
Columnar child only. coalesce(n, shuffle = false) over the child's batches; within an output partition the child partitions are drained one after the other and each child iterator releases at task end. Pass-through, as the union. |
spark.vecruntime.exec.coalesce.enabled |
child <op> is not columnar; coalesce to <n> partitions |
GenerateExec (explode, posexplode, explode_outer, posexplode_outer, LATERAL VIEW [OUTER]) over an array column |
VectorGenerateExec |
Columnar child; the generator is explode/posexplode over an array column or a struct field of one (Spark's NestedColumnAliasing projects st.inner first -- that projection passes the array through), with an element type that is a lane. The array reaches the operator as Spark's own vector (#19): a repeat index is built from the array lengths, the forwarded columns are gathered through it (ArrowOutput.gather; a column without a lane as a RemappedColumnVector view), the elements are copied once per array from the array vector, the position column is an iota per array; a null or empty array contributes one null row under OUTER and nothing otherwise. One output batch per input batch, whatever the arrays expand to. The size(arr) > 0 AND isnotnull(arr) filter Spark adds below a non-outer explode compiles (SizeExpr, NestedValidityExpr), so the chain stays columnar. |
spark.vecruntime.exec.generate.enabled |
child <op> is not columnar; generator <name> not supported (explode and posexplode over an array are) (inline, stack, json_tuple, UDTFs); explode over a map not supported; array element type <type> not supported in explode; explode over <expr>: nested access over ... not supported (a computed array) |
SampleExec (TABLESAMPLE (p PERCENT), BUCKET x OUT OF y, df.sample) without replacement |
VectorSampleExec |
Columnar child only, any lane type (#259). A selection producer: Spark's own BernoulliCellSampler(lowerBound, upperBound) seeded with seed + partitionIndex and drawn once per live row in input order -- exactly what SampleExec's row path and generated code do, so the same seed returns the same rows as Spark (GapSamplingIterator belongs to RDD.sample, not the SQL path). The bitmap is forwarded to a VecRuntime parent or compacted for Spark like the filter's. |
spark.vecruntime.exec.sample.enabled |
child <op> is not columnar; sampling with replacement (Poisson) not supported |
WindowExec -- row_number, rank, dense_rank, percent_rank, cume_dist, ntile (layer 1 of #58); whole-partition and running aggregates (layer 2); lag, lead, first_value, last_value, nth_value (layer 3) |
VectorWindowExec (spark.vecruntime.exec.window.enabled) |
Any child, accepted on types alone: Spark plans Window above Sort above an exchange, and without a columnar shuffle that sort is Spark's, so RowToColumnarExec is inserted below us and the chain above (the filter on the rank, the projection) is columnar again. The child arrives sorted by the partition keys then the order keys (the operator requires exactly what WindowExec does), so the computation is one walk: a row starts a partition when a partition key differs from the previous row's and a peer group when an order key does -- null-safe equality, as Spark's RankLike compares -- and the counters and the previous row's keys carry across batches, so a partition longer than a batch is one partition. Input columns are forwarded (borrowed, never copied), the ranks are new INT32 columns. Several specs in a query are several stacked operators, each ours. Whole-partition aggregates -- any function of the grouped aggregate family over the frame UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING, Spark's default for a spec without ORDER BY -- treat each partition as a group: the grouped aggregate states run over every batch with the partition ordinal as the group id, the value is the function's evaluateExpression over the partition's buffers exactly as the Final aggregate computes it, and rows are held (copied, in memory, no spill -- like the sort) until their partition ends, then emitted batch by batch with the values gathered per row. Decimal aggregates run over whole-partition and running frames since #259 (slice 4): the 128-bit accumulators of the hash aggregate, a running sum or average combining exact partial totals per row and finalising as the merge does (an overflowed running sum is null, or the ANSI error), so the TPC-DS window aggregates (q12, q20, q47, q53, q57, q63, q89, q98) and q51's decimal(27,2) running totals are ours; only a sliding ROWS frame over a decimal is still refused (over decimals in a sliding frame not supported). Partition and order keys may be wide decimals (two-limb equality). Running frames -- ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, and RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, Spark's default for an aggregate over a spec with ORDER BY, where peers share the value at the end of their group -- make each row or peer group a group and, at emission, combine each group's buffers with the running buffers before it in the partition, slot by slot (sums and counts add, a bigint sum with Spark's ANSI overflow error; min/max compare), so the same evaluateExpression yields the running value; sum, avg, count, min and max have that form, the other functions are refused (running frame for first not supported). Keys: every supported type but double (double keys not supported, as for joins and aggregates). Offset functions -- lag / lead with a literal offset and a literal (or absent) default, first_value, last_value, nth_value with a literal n, over the whole-partition and running frames -- are one row of the partition each, read from the held copies of its rows (a released batch stays readable until no later batch can address one of its partitions, since lag reads backwards across batch boundaries): the shifted row or the default outside the partition, the partition's first row, the frame's last row (the end of the peer group under the RANGE default), the frame's n-th row or null. percent_rank ((rank - 1) / (n - 1), 0 for a single row), cume_dist (rows up to the end of the peer group over n) and ntile (Spark's bucket walk: the first n % buckets buckets hold one row more) need the partition size, so they take the held-partition path too, and row_number / rank / dense_rank join them there when they share an operator with an offset function (peer starts and the partition's first peer are tracked per row). Sliding frames -- sum, avg, count, min, max over ROWS BETWEEN a AND b with literal bounds (n PRECEDING/FOLLOWING, CURRENT ROW, unbounded either side, so the suffix CURRENT ROW AND UNBOUNDED FOLLOWING too) -- are re-aggregated per row over the frame's rows in order on the held-partition path, exactly as Spark's sliding frames re-aggregate (double sums bit-identical; an empty frame gives null sums and 0 counts); a frame starting at the partition is advanced as its end moves, as Spark's unbounded-preceding frame is. RANGE frames with value offsets -- the same five aggregates over RANGE BETWEEN 5 PRECEDING AND CURRENT ROW, 1 FOLLOWING AND 3 FOLLOWING, n PRECEDING AND m FOLLOWING, an unbounded side with an offset on the other, over one integral (tinyint..bigint) or date order key, ASC or DESC, either null ordering -- are computed per partition by the WindowFrameKernels: the partition's order keys and the input are gathered from the held batches into primitive arrays reused across partitions, every row's frame is found by Spark's own two-pointer walk (rangeBounds: the bound is the key plus the offset in the key's arithmetic -- wrapped, or the ANSI overflow error, as Spark's Add; DateAdd for dates -- compared in the SortOrder's direction and null ordering, so a null key's frame is exactly the null peer group and a valued row never sees a null key), and the aggregate is re-run in row order over each frame as Spark's SlidingWindowFunctionFrame does (double sums bit-identical; a frame whose lower bound has not moved continues the previous accumulation; count and an unchecked bigint sum slide by subtracting the rows that leave; an ANSI sum re-adds with addExact so the same partial sum overflows; an empty frame is a null sum / 0 count). The output column is written straight into its Arrow buffers -- no per-row objects. The peer-bounded RANGE frames without an offset (CURRENT ROW AND UNBOUNDED FOLLOWING, CURRENT ROW AND CURRENT ROW) run on the row path over any key type. Aggregates of different frames, offsets and ranking functions may share one operator: they all take the held-partition path. Refused with a reason naming the function and frame: RANGE offsets over a decimal, timestamp (interval offsets) or string order key (RANGE offsets over a decimal(12,2) order key not supported), min/max over a string or boolean in such a frame, sliding or running frames of other functions (running frame for stddev not supported, over a sliding frame not supported), IGNORE NULLS, decimals in sliding frames, a ranking spec without ORDER BY; one unsupported function in a spec keeps that operator Spark's. |
||
WindowGroupLimitExec (Partial and Final) -- Spark's per-partition top-k under a ranking window, inserted for a filter rank <= k with k up to spark.sql.optimizer.windowGroupLimitThreshold (layer 4 of #58) |
VectorWindowGroupLimitExec (spark.vecruntime.exec.window.enabled) |
Any child on types alone: Partial sits below the exchange over whatever produced the rows (in TPC-DS our aggregate or project, which it no longer forces back to rows), Final above Spark's sort. Both modes require the window's ordering, so it is the ranking walk with a filter: a row survives while row_number / rank / dense_rank at that row is at most k -- ties at k all survive under rank, whole peer groups under dense_rank. Spark's own operator also passes the first row of the peer group past k (it tests the previous row's rank); both are pre-filters under the window that computes the real ranks, so the tighter set is exact for the query and smaller for the shuffle. Output rows are compacted (a top-k keeps a few rows per partition); a batch with every row kept is forwarded, one with none is dropped. Keys as for the window (no doubles). |
||
LocalTableScanExec (VALUES, small createDataFrame, local relations) |
VectorLocalTableScanExec |
Off by default -- there is nothing to accelerate, this only removes the RowToColumnarExec transition so small-table tests exercise our operators (Comet keeps its equivalent off for the same reason). One on-heap batch per partition through the VectorRowStages.toBatch row writer, min(max(rows, 1), leafNodeDefaultParallelism) partitions as Spark's. Output types must be supported. Note the optimizer folds filters and projections over a literal relation on the driver (ConvertToLocalRelation); exclude that rule, as Spark's own operator tests do, for the operators to appear at all. |
spark.vecruntime.exec.localTableScan.enabled |
unsupported type <type> for local table column <name>; streaming local table scan |
RangeExec (spark.range(...), the range() table-valued function) |
VectorRangeExec |
A leaf, so no input condition. Over Spark's row RangeExec nothing of ours converts until the first exchange (a row child is refused by every operator; the filter, projection and partial aggregate over range() stayed Spark's, and only a merging aggregate or a join above the shuffle was ours, behind a RowToColumnarExec); with this leaf the whole chain is columnar from the source. Spark's semantics exactly -- start, end, positive or negative step, numSlices (or leafNodeDefaultParallelism), an empty range as an RDD with no partitions -- and the same split: partition i holds elements [i * n / slices, (i + 1) * n / slices) of the n values start + k * step, the partition bounds clamped to Long as RangeExec.getSafeMargin clamps them and the row count derived from the clamped bounds as Spark's generated initRange derives it, so spark_partition_id(), monotonically_increasing_id() and the values near Long.MaxValue / Long.MinValue match row for row (VectorRangeSuite). outputOrdering (id ASC) and outputPartitioning (RangePartitioning(id, slices), SinglePartition for one slice) are Spark's, so the sort on id the planner elided and the exchange a GROUP BY id does not need stay valid. Each task fills one native INT64 BigIntVector of spark.sql.inMemoryColumnarStorage.batchSize rows (the size the replaced RowToColumnarExec produced) with SequenceKernels.range -- lane offsets plus a broadcast base per block -- and re-emits it for every batch under the columnar reuse contract (as Spark's RowToColumnarExec reuses its vectors); nothing is allocated per row or per batch, the vector and allocator are released at task end. Metrics numOutputRows, numOutputBatches, time; the task's recordsRead is incremented per batch as Spark's is per row. The single id column is bigint NOT NULL. |
spark.vecruntime.exec.range.enabled |
streaming range; range with <n> slices (a slice count below 1, which Spark itself rejects at execution) |
ExpandExec (ROLLUP, CUBE, GROUPING SETS, the count(distinct) rewrite) |
VectorExpandExec |
Columnar child only. Every projection slot must be a column reference or a literal (incl. NULL) of a lane type -- what grouping-set analysis emits once the keys are aliased below; a wide decimal key or sum passes as a DECIMAL128 lane and its nulled slot is a wide null column (#259). One output batch per projection per input batch with no data copy: retained columns borrowed, nulled keys as all-invalid constant columns, the grouping id as a constant column. The input batch is held until its last projection has been emitted and released; pass-through for the selection rule. |
spark.vecruntime.exec.expand.enabled |
child <op> is not columnar; expand without projections; unsupported output type <type> for <name>; projection <i>: literal of unsupported type <type>; projection <i>: <expr> is not a column or literal; projection <i>: <expr>: <compiler reason> |
BroadcastHashJoinExec |
VectorBroadcastHashJoinExec |
Streamed side columnar, or an exchange (its AQE stage, a shuffle read); its payload columns may be of any type -- one without a lane (struct, array, map) is passed through as a remapped view of the streamed batch, null-padded for unmatched build rows (#273); the build side's keys and the columns its condition reads need lane types, and a build payload column without a lane (array, map, struct, nested any depth over Spark's physical types) is copied out of the broadcast rows into a row store -- Spark's on-heap column vector plus one null row -- and read as a remapped view over the build row ids the probe emits, a padded outer-join row reading the null row (#547, spark.vecruntime.join.buildPayload); the lane columns -- a wide decimal key or payload is a DECIMAL128 lane on the streamed side and is laid out as two limbs from the broadcast rows (#259) -- Spark converts it below us with RowToColumnarExec as for the shuffled hash join; that is the shape adaptive execution leaves when it re-plans a shuffled join as a broadcast join at runtime, and refusing it left the chain above to Spark in eleven TPC-DS queries (#98). build side is Spark's own BroadcastExchangeExec/HashedRelation (kept so a Spark join over the same broadcast still works) and only its types matter; the relation is read into columns (strings as their UTF-8 bytes, never decoded) and its key table built once per executor, shared read-only by every task of every join over it (#11 is now only the row-by-row read itself). Equi-join keys of type int/bigint/boolean/string/date/timestamp/decimal(<=18)/double (doubles compare by bits after the NormalizeNaNAndZero pass Spark's optimizer puts on them), never a literal. Join types: inner, left outer / left semi / left anti / existence (EXISTS used as a value: every left row plus the planner's exists boolean, false -- never null -- where nothing matched; TPC-DS q10 / q35 / q45) with build right, right outer with build left; each with an optional non-equi condition. Null keys never match. Null-aware anti join (#265): Spark's plan for NOT IN (subquery) over nullable columns, a single-key left anti join with the right side broadcast and no condition, built with isNullAware. Its three regimes: an empty build side is Spark's EmptyHashedRelation singleton and the anti join over an empty table keeps every streamed row, null keys included; a build side holding a null key is the HashedRelationWithAllNullKeys singleton (it has no rows to read) and the task emits nothing; otherwise the streamed rows with a null key are dropped along with the matched ones (JoinSpec.dropNullStreamedKeys, applied only when the table has rows). Any other shape under the flag is refused. No full outer join over a broadcast -- Spark never plans one (JoinSelection excludes it) and the shared build side would emit its unmatched rows once per task: full outer join over a broadcast not supported (Spark plans it as a shuffled join). |
spark.vecruntime.exec.broadcastHashJoin.enabled, spark.vecruntime.join.maxBuildSize |
child <op> is not columnar; unsupported column type ...; build side estimated at <n> bytes exceeds spark.vecruntime.join.maxBuildSize=<m> (#86: the build side is held in memory per task; default 1 GiB or a per-core share of the off-heap budget; an unknown estimate converts); null-aware anti join not supported; join without equi-join keys; join key type <type> not supported; literal join key; join type <type> with build side <side> not supported; any compiler reason for the condition |
BroadcastNestedLoopJoinExec (non-equi joins: inequality joins, cross joins with a filter, non-equi EXISTS) |
VectorBroadcastNestedLoopJoinExec |
Streamed side columnar or an exchange (as for the hash joins); build side is our columnar broadcast when the exchange below is ours (#325), otherwise Spark's identity broadcast of the rows; either way it is read into columns once per executor (shared by its tasks, as the hash joins' table), and only its types matter. The same iterator as the hash joins with every build row a candidate of every streamed row: the streamed rows go in chunks sized to a fixed pair budget (32k pairs), each chunk's pairs are gathered, the compiled condition evaluated over them, and the verdicts turned into inner rows, semi/anti/existence decisions or outer null padding -- the product is never materialised. Join types: inner and cross with either build side; left semi / left anti / existence / left outer with build right; right outer with build left; each with an optional condition (a cross join has none). Refused: full outer, and an outer join whose preserved side is the broadcast one (both need a matched bitmap over the shared broadcast, which every task would emit) -- as Comet refuses them. | spark.vecruntime.exec.broadcastNestedLoopJoin.enabled, spark.vecruntime.join.maxBuildSize |
child <op> is not columnar; unsupported column type ...; build side estimated at <n> bytes exceeds spark.vecruntime.join.maxBuildSize=<m>; full outer nested loop join not supported (needs a matched bitmap over the broadcast side); `<LeftOuter |
BroadcastExchangeExec (join broadcasts) |
VectorBroadcastExchangeExec |
The child is one of our operators or a columnar source, every column with a lane; the mode a HashedRelationBroadcastMode or the nested-loop join's IdentityBroadcastMode (#325). Each partition's batches travel to the driver as an Arrow IPC stream and are broadcast as they are (Spark's row and byte limits apply); VectorBroadcastHashJoinExec and VectorBroadcastNestedLoopJoinExec build their shared table from them once per executor, with no ColumnarToRow below the exchange and no rows on the way. A Spark consumer of the same exchange (a Spark join over a reused exchange or with our join switched off, dynamic partition pruning's SubqueryBroadcastExec) gets Spark's own value -- the relation, or the array of rows -- built from the batches on the driver on first use and broadcast once; with none it is never built. Dynamic partition pruning keeps its filter. Stays Spark's: a child with a lane-less column, under the Comet mixed pass (Comet converts only its own exchange with the join above), and with spark.vecruntime.exec.broadcastExchange.enabled=false |
||
ShuffledHashJoinExec |
VectorShuffledHashJoinExec |
Both inputs are exchanges (or their AQE stages / AQEShuffleRead), the build side accepted on lane types like the Final aggregate and the streamed side on any type (a lane-less payload passes through, #273) (wide decimal keys hashed and compared on both limbs, wide payloads gathered at the 16-byte stride, #259); the build side of each partition is drained into columns first -- in memory up to spark.vecruntime.join.spillBytes, past it both sides split into spark.vecruntime.join.spillBuckets buckets on disk and are joined one bucket at a time (#416, GraceHashJoin), so no build size is refused. Same keys, join types and conditions as the broadcast join, plus full outer (TPC-DS q51 / q97) and, by the same flags, a left or right outer join built from its preserved side (#273; the sort-merge route may then build the smaller side of any outer join): per-build-row matched flags, then a trailing pass emitting the unmatched build rows with the streamed side null -- per partition, so partitions holding build rows but no probe rows still emit them; rows failing the condition are unmatched on both sides. Spark plans it only with spark.sql.join.preferSortMergeJoin=false or a SHUFFLE_HASH hint (#10). |
spark.vecruntime.exec.shuffledHashJoin.enabled, spark.vecruntime.join.maxBuildSize |
As the broadcast join (the build-size estimate is the shuffle stage's runtime statistics under AQE), plus skew join not supported |
SortMergeJoinExec |
VectorShuffledHashJoinExec (the default under spark.vecruntime.exec.sortMergeJoin.mode=auto wherever the join's order cannot show; hash forces it) |
Spark's default for large equi joins, re-expressed as our shuffled hash join (#10, the issue's second option): both operators require the same distribution of the same keys and give the same rows for every join type we plan (inner, left/right/full outer, semi, anti, existence, with a condition), so per partition the hash join is a drop-in and the two sorts Spark placed for the merge are dropped (exactly the required ones). What changes is the memory profile -- the smaller side becomes a per-task hash table -- hence the build side is the smaller by the runtime statistics of the two AQE stages and must fit spark.vecruntime.join.maxBuildSize; a side without statistics (adaptive execution off: Spark's shuffles carry none) is not assumed small (no size statistics). A merge join whose output ordering an ancestor relies on stays Spark's (output ordering required by the parent operator): a Window partitioned by the join key sits directly on it -- the rule decides this in a top-down pre-pass over the whole plan (markSortMergeJoins), since a hash join has no ordering to offer. A chain of merge joins on the same key with no shuffle between the links converts whole: the upper join reads the lower one's output, columnar once the lower converts, and the pre-pass judges each input by what it will be after the transform (an exchange, a converting merge join, a project or filter over one, another operator the rule replaces) rather than by Spark's row operator. That judgement is optimistic on purpose, and the transform verifies it: a join that has to stay after all gets its ordering back through a VectorSortExec over any converted child below it, so a wrong guess costs a sort, never a row (#102). Rows are identical; the order of tied rows under an ORDER BY and the rows an unordered LIMIT picks can differ from Spark's, whose merge join is order preserving -- three files of Spark's SQL suite (subquery/in-subquery/in-limit, in-order-by, in-set-operations) pin that order and differ under the rewrite, which is one reason it is opt-in. Skew joins (AQE's split partitions) are refused like the shuffled hash join's. |
||
SortMergeJoinExec |
VectorSortMergeJoinExec (spark.vecruntime.exec.sortMergeJoin.mode=merge, or auto -- the default since #311 -- where the pre-pass chooses it: a parent relies on the ordering, the row order can reach a limit or a sort without an exchange in between, or the hash rewrite is not allowed -- no statistics, the build side past spark.vecruntime.join.hashMaxBuildSize per task (one bucketing pass of the grace hash join), a build side larger than the streamed side, a skew join; no input size is too large since #416 (the sort below spills past spark.vecruntime.sort.spillBytes), so nothing is left to Spark's own operator; hash and the boolean flag (= auto, false = off) as above; the operator prints the decision, #287) |
A real, order-preserving merge join (#286). Spark's contract kept: both children clustered on the keys and sorted by them ascending, so the sorts Spark placed below stay (ours over a columnar child, Spark's over its row shuffle with a RowToColumnarExec inserted above), and the output ordering is the join's own -- the preserved side's for inner, left outer, semi, anti and existence joins, the right side's for a right outer join (run with the sides swapped), none for a full outer -- so a parent relying on the merge's ordering converts and no build side or statistics are needed; skew joins are fine (the partitions are already split). Execution streams both sides: the buffered side is read run by run (RunKernels.boundaries over the sorted keys, one lane-parallel pass per key) and the current run -- its rows, its keys, a matched bitmap for the outer types -- is the only buffered state: a run inside one batch borrows the batch, a run reaching a batch edge is copied into its own (confined) arena and continued while the key holds, so memory is bounded by the largest run; the streamed side is walked run by run within each batch and equal runs emit their cross product as index pairs in Spark's order (streamed row, then the buffered rows), 8192 pairs per chunk, gathered by the same kernels as the hash join with -1 padding a missing side; a condition is evaluated over the gathered pairs and the survivors set the matched flags semi, anti, existence and the outer types read. Null keys never match (a null run is padded for outer types, skipped otherwise). Every join type Spark's operator supports; keys of every lane type; every column must have a lane (no pass-through for a struct yet). |
spark.vecruntime.exec.sortMergeJoin.mode |
join without equi-join keys; join type <type> not supported; <key>: <compiler reason>; unsupported column type <type> for <column>; <condition reason> |
ShuffleExchangeExec above a VecRuntime operator |
VectorShuffleExchangeExec (#288) |
The vecruntime-shuffle jar and spark.shuffle.manager=org.apache.spark.sql.vecruntime.shuffle.VectorShuffleManager (row shuffles still go to Spark's sort shuffle manager underneath); hash, single, round-robin and range partitioning (range samples the child like Spark and binary-searches its bounds per row); every column of a type the Arrow layer carries. Partition ids come from a Murmur3 kernel equal to Spark's Pmod(Murmur3Hash), each map output is one partitioned Arrow IPC file with dictionary-encoded strings, and a reducer streams the other executors' blocks over Arrow Flight (spark.vecruntime.shuffle.backend=flight, default; block for Spark's block transfer) -- no ColumnarToRowExec / RowToColumnarExec around the exchange. AQE coalescing, skew splitting and local reads are honoured through ShuffleExchangeLike. Taken only when Comet's native shuffle is not taking the exchange. |
spark.vecruntime.shuffle.enabled, spark.vecruntime.shuffle.backend |
none recorded -- an exchange that is not ours is left as Spark's and shows as Spark in the UI |
ShuffleExchangeExec / Comet's JVM columnar shuffle above a VecRuntime operator |
Comet native shuffle over VectorToCometExec |
Comet on the classpath, spark.shuffle.manager=...CometShuffleManager, spark.comet.exec.shuffle.enabled=true; hash, single, round-robin and range partitioning (range also needs Comet's spark.comet.shuffle.native.partitioning.range.enabled and samples the child twice, like Spark); every column and range key of a type the bridge carries (decimals cross widened to 128-bit). Without Comet and without our shuffle manager the exchange stays Spark's row shuffle, with ColumnarToRowExec above our operator and RowToColumnarExec below the consumer. |
spark.vecruntime.comet.shuffle.enabled, spark.vecruntime.comet.shuffle.range.enabled |
none recorded -- an exchange that is not bridged is simply left as Spark's and shows as Spark in the UI |
| A Comet operator above a VecRuntime chain | Comet's operator over the sink leaf (CometUnion pass-through over VectorToComet), planned by Comet's own rule from ours (spark.vecruntime.comet.mixed.enabled, default off): projection, filter, limits, expand, union, and an aggregate half when Comet's buffer predicate allows; joins and final aggregates through Comet's native shuffle; Comet's decline reasons recorded as fallbacks (docs/comet.md, #280) |
Leaf scans are inputs, not conversions: a FileSourceScanExec whose vectorized Parquet reader
emits ColumnarBatches (#61), Comet's native scan, and a BatchScanExec over Iceberg through
IcebergVectorAdapter (or any other DSv2 source whose ColumnVectors the adapters copy, #62) all
count as columnar children. The UI shows them as "Spark columnar scan", which neither counts as
accelerated nor as a missed conversion. Row/columnar transitions (ColumnarToRowExec,
RowToColumnarExec) and AQE plumbing (AQEShuffleReadExec, ReusedExchangeExec) are shown the
same way.
Scan compatibility
What makes a FileSourceScanExec a columnar input is decided by Spark, not by this project: the rule
accepts a scan when supportsColumnar holds and every output column has a lane type. Pinned by
VectorScanSuite (#61).
| Case | What happens | Reason recorded |
|---|---|---|
| Parquet, vectorized reader (the default) | Columnar input. OnHeapColumnVector batches are copied once into native memory per operator chain (SparkColumnVectorBuffers.copy; reading the heap arrays in place was measured at half the kernel speed); dictionary-encoded strings stay dictionary encoded; dictionary-encoded numerics are decoded once per batch |
-- |
Parquet, spark.sql.columnVector.offheap.enabled=true |
Columnar input. Fixed-width columns (int, bigint, double, date, timestamp, decimal <= 18 stored as int/long) without a dictionary are handed over as views of the OffHeapColumnVector's native memory -- no copy (SparkColumnVectorBuffers.wrappedOffHeapColumns counts them); the validity bitmap is derived from Spark's byte-per-row nulls, strings and dictionaries take the copy |
-- |
Parquet, spark.sql.parquet.enableVectorizedReader=false |
A row scan; nothing above it converts until a shuffle (a merging aggregate reads the shuffle and stays ours over RowToColumnarExec) |
child FileSourceScan is not columnar |
| Parquet with a nested column (struct, array, map) in the scan output | Spark 4's nested vectorized reader keeps the scan columnar; a filter or projection above passes the column through as Spark's vector (#19) and reads struct fields from its child vectors (st.a, st.c.d; #50), while any other operator over the nested column itself (aggregate, join, sort, window, exchange bridge) falls back on the type until the column is pruned or only its fields are used |
unsupported column type <type> for <name> (operators other than filter/project) |
| Wide decimals (precision > 18) in the scan output | Columnar in Spark and here: the DECIMAL128 lane (#257); expressions (#258) and the hash aggregate (#259) compute on it, the sort orders by it; a wide key or payload in a join or window still falls back (#259) | unsupported column type decimal(p,s) for <name> from those operators |
| ORC, vectorized reader (the default) | Columnar input through the adapter seam's generic copy (OrcColumnVector is read through the ColumnVector getters) |
-- |
Cached table (CACHE TABLE, df.cache()) whose whole cached relation is boolean/byte/short/int/long/float/double |
Columnar input (#55, level 1): InMemoryTableScanExec decompresses Spark's CachedBatches into on-heap batches, copied once per batch like a Parquet scan; our operators sit directly on the scan. Pinned by VectorCacheSuite |
-- |
Cached table with a string, date, timestamp, decimal or nested column anywhere in the cached relation (or spark.sql.inMemoryColumnarStorage.enableVectorizedReader=false) |
A row scan: Spark's DefaultCachedBatchSerializer.supportsColumnarOutput is decided on the relation's full schema, not the projected columns, and only admits the primitive types above. Nothing above it converts until a shuffle. Caching in a format of our own (an Arrow CachedBatchSerializer, level 2 of #55) is the way to lift this |
child Scan In-memory table <name> is not columnar / child InMemoryTableScan is not columnar |
| CSV, JSON, text, Avro, JDBC | Row-based readers; out of scope | child <op> is not columnar |
Comet's native scan, Iceberg's BatchScanExec |
Columnar input, zero-copy through the registered adapters (see docs/comet.md, docs/iceberg.md) | -- |
| Any other DSv2 columnar source | Columnar input, copied once per batch (#88) | -- |
Planned
Tracked issues, in the order they unblock TPC-H:
| Spark operator | Issue | What is missing |
|---|---|---|
InMemoryTableScanExec -- cache in our own format |
#55 (level 2) | An Arrow CachedBatchSerializer (spark.sql.cache.serializer) so cached strings, dates, timestamps and decimals come out columnar with no decompression; level 1 (consuming Spark's columnar cache output) is done, see "Scan compatibility". Measure Spark's compression trade-off before deciding |
SortAggregateExec |
VectorHashAggregateExec (+ VectorSortExec) |
Spark plans a sort-based aggregate when a buffer holds a string (min/max/first/last/max_by over strings), which an UnsafeRow cannot mutate; our group table has no such limit, so the same hash operator is built from the identical fields. Two contracts kept: the sort Spark placed below is dropped when it is exactly the required one, and a result-emitting stage keeps Spark's output ordering (the grouping keys ascending) through a VectorSortExec above, since parents were planned on it. |
ObjectHashAggregateExec |
VectorHashAggregateExec |
The object aggregates whose buffer we can carry (#57): bloom_filter_agg (the runtime filter's build side), collect_list, collect_set. Spark carries their state between the Partial and Final stages as one BinaryType column (serialize / deserialize), so we drive Spark's own TypedImperativeAggregate object per group -- createAggregationBuffer / update (over the batch rows through ColumnarBatch.getRow, the selection skipping unselected ones) / merge (of the input binary buffer) / serialize (the emitted partial buffer) / eval (the Final's array<T> or serialized filter). The partial buffer is byte-identical to Spark's, so a Spark Final or the bloom probe reads it unchanged, and a merging stage reads Spark's partial buffers unchanged; the array / binary output is held in Spark's own on-heap column vector (the join payload store, #547), read back as the declared type. All four modes, grouped and ungrouped, mixed with ordinary aggregates in one node; a columnar child (the aggregated value may be a non-lane column -- it is read by row). Nulls ignored and collect_set dedup are Spark's, since Spark's object does them. Object aggregates are counted against the task's memory budget (a bloom filter's bit array, a collect buffer's grown elements) and spill past spark.vecruntime.agg.spillThreshold (#57): a buffer-emitting stage emits its partial binary buffers and starts over, a Final spills its groups' serialized buffers (UTF8-carried through the grace-hash path) and merges them back one bucket at a time, so a grouped collect_* over many keys is bounded rather than pinning memory. A Complete stage (streaming only; Spark 4.1's batch planner never emits it) has no mergeable buffer input to reload and stays in memory. Percentiles (percentile, percentile_approx), collect_top_k and the rest keep the fallback. |
WindowExec -- running or sliding frames (ROWS or RANGE) of functions other than sum/avg/count/min/max, decimal aggregates over sliding frames, RANGE offsets over decimal or timestamp (interval) order keys, IGNORE NULLS on the offset functions |
#58 (residual), #28 | Every layer of #58 is converted above, RANGE frames with value offsets included. The other functions' sliding forms need their update expressions re-run per frame; the sliding kernels are long and double (whole-partition and running decimal frames run since #259) -- the TPC-DS window aggregates (q12, q20, q47, q53, q57, q63, q89, q98) are all decimal; a decimal RANGE offset needs DecimalAddNoOverflowCheck arithmetic on the key and an interval offset the session time zone |
BatchScanExec -- DSv2 sources beyond Iceberg |
#62 | Verification per source (Delta, Hudi, built-in DSv2). Unknown columnar sources already work and are copied once per batch: VectorUnknownSourceSuite (#88) reads a test-only DSv2 source whose batches are a ColumnVector class no adapter knows, checks the rows against Spark under our filter, projection and aggregate, and asserts through the adapter seam's counters (ColumnVectorAdapters.copiedColumns / adaptedColumns) that every column took the copy path; a struct column falls back with unsupported column type struct<...> for <name> until #19 |
DataWritingCommandExec -- Parquet from Arrow batches |
#64 | Not converted |
| Arrow-based Python UDF operators | #65 | Not converted |
Spill for VectorSortExec and the hash joins |
#12, #85, #86, #416 | Done for the sort (runs past spark.vecruntime.sort.spillBytes merged from disk, #448) and the shuffled hash join (both sides bucketed on disk past spark.vecruntime.join.spillBytes, GraceHashJoin); the broadcast joins still refuse a relation estimated above spark.vecruntime.join.maxBuildSize at planning time |
| Columnar broadcast exchange | #325 | Done for the hash joins and the nested-loop join (VectorBroadcastExchangeExec, see above) |
Not planned
Recorded here so that an absence is a decision rather than an omission (#66). These are open to revisiting on demand, not permanent exclusions.
| Family | Reason |
|---|---|
Structured Streaming operators (StateStoreSaveExec, StateStoreRestoreExec, StreamingSymmetricHashJoinExec, FlatMapGroupsWithStateExec, ...) |
This project targets batch execution. The state store is row-oriented and its on-disk format is a compatibility surface across Spark versions and checkpoints. A streaming query still benefits where a micro-batch contains supported batch operators. |
CartesianProductExec |
Rare, and expensive for reasons a columnar engine does not change. A cross join with a broadcast side and a filter is a different plan -- see BroadcastNestedLoopJoinExec (#60). |
Pickled (non-Arrow) Python UDFs (BatchEvalPythonExec) |
The data has to become Python objects row by row; there is no columnar path. Arrow-based UDFs are the ones worth accelerating (#65). |
| Command and DDL operators | Nothing to accelerate. |
Keeping this page honest
Every operator issue names this file as part of its definition of done: the row, the requirements,
the config key and the fallback strings land in the same commit as the operator. Until the matrix is
generated from PlanAcceleration (the UI's node classification, which is the natural source), the
check is the suites: a fallback string that changes without its row changing fails a
checkFallback assertion before it reaches a user.