Configuration reference
Every spark.vecruntime.* key is a session-level Spark SQL configuration: set it with --conf on
spark-submit, in spark-defaults.conf, or at run time with spark.conf.set(...). The planner
reads them from the session's SQLConf when it plans a query, so a change applies to the next query
(the two UI keys and the shuffle keys are read from the SparkConf at start-up, see the notes in
their tables). The defaults were set from measurements on TPC-DS at 1 TB and TPC-H at SF10 -- see
docs/results.md for the runs behind each threshold. Size-valued keys accept Spark's size strings
(512m, 1g) as well as byte counts. Where a spark.vecruntime.* key is read is named in
spark/src/main/scala/io/vecruntime/spark/VectorConf.scala (the planner and operator keys),
spark/src/main/scala/org/apache/spark/sql/vecruntime/AggregateSpill.scala (the aggregate spill keys) and
the shuffle module (the columnar shuffle keys).
Main switch and diagnostics
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.enabled |
true |
boolean | Main switch: false leaves every operator to Spark and the plugin only attaches its UI tab. |
spark.vecruntime.explainFallback.enabled |
false |
boolean | Logs why each operator was left to Spark (the same reasons the Vector Acceleration tab shows); turn it on when a plan is not as accelerated as expected. |
spark.vecruntime.ui.enabled |
true |
boolean | Attaches the Vector Acceleration tab to the Spark UI; read from the SparkConf by the driver plugin before any session exists, so it must be set at start-up. |
spark.vecruntime.ui.retainedExecutions |
100 |
positive int | How many SQL executions the Vector Acceleration tab keeps; also read from the SparkConf at start-up. |
spark.vecruntime.iceberg.dvWriter.enabled |
false |
boolean | The columnar Iceberg v3 deletion-vector writer (#20): on a format-version-3 merge-on-read table (delete file format PUFFIN), writes the DELETE half of DELETE/UPDATE/MERGE from the _file/_pos lanes as deletion vectors instead of Spark's row-by-row WriteDeltaExec; inserts and the commit stay Iceberg's. v2 tables decline to Spark's writer. Covers DELETE, UPDATE and MERGE (the insert half through Iceberg's own writer), partitioned tables and repeated deletes. Default off: correct in every Iceberg merge-on-read suite, but on a CDC MERGE the gain is small (see docs/iceberg-dv-writer.md). Needs the vecruntime-iceberg-bridge jar on the classpath. |
Operator switches
Each converts one Spark operator into its VecRuntime counterpart when the conditions in
docs/operators.md hold; set one to false to keep Spark's operator for that kind while the rest
of the plan stays columnar.
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.exec.filter.enabled |
true |
boolean | Convert FilterExec. |
spark.vecruntime.exec.project.enabled |
true |
boolean | Convert ProjectExec. |
spark.vecruntime.exec.aggregate.enabled |
true |
boolean | Convert HashAggregateExec (and the SortAggregateExec Spark plans for string buffers). |
spark.vecruntime.exec.aggregate.final.enabled |
true |
boolean | Also convert the merging modes (Final, PartialMerge), whose input is a shuffle; false keeps Spark's Final aggregate over our partial one. |
spark.vecruntime.exec.objectAggregate.enabled |
true |
boolean | Convert ObjectHashAggregateExec for the object aggregates we carry (#57): bloom_filter_agg, collect_list, collect_set, driven through Spark's own function object per group so the partial buffer and result are byte-identical. Needs exec.aggregate.enabled (and exec.aggregate.final.enabled for the merge modes). Other object aggregates (percentile*, collect_top_k, ...) keep the fallback. |
spark.vecruntime.exec.aggregate.rollupRewrite.enabled |
true |
boolean | A partial aggregate over our Expand (ROLLUP, CUBE, GROUPING SETS) aggregates the finest grouping once and rolls the partials up instead of hashing every input row once per grouping set (#383); false disables the rewrite. |
spark.vecruntime.exec.mergeRows.enabled |
true |
boolean | Convert Iceberg's MergeRowsExec (the row-level operator of a MERGE INTO, #21) when the join below it is ours. |
spark.vecruntime.exec.sort.enabled |
true |
boolean | Convert SortExec over a columnar child; a global ORDER BY over Spark's row shuffle stays Spark's regardless. |
spark.vecruntime.exec.takeOrdered.enabled |
true |
boolean | Convert TakeOrderedAndProjectExec (ORDER BY ... LIMIT n) over a columnar child; the per-partition top-n is columnar, the final merge goes through Spark's single-partition shuffle. |
spark.vecruntime.exec.limit.enabled |
true |
boolean | Convert LocalLimitExec / GlobalLimitExec / CollectLimitExec over a columnar child; OFFSET falls back. |
spark.vecruntime.exec.union.enabled |
true |
boolean | Convert UnionExec when at least one child is columnar (row children go through RowToColumnarExec); keep it on with Spark 4.1.3, whose own columnar union concatenates co-partitioned children (#128). |
spark.vecruntime.exec.coalesce.enabled |
true |
boolean | Convert CoalesceExec (a shuffle-free coalesce(n)) over a columnar child. |
spark.vecruntime.exec.expand.enabled |
true |
boolean | Convert ExpandExec (grouping sets, the count(distinct) rewrite) over a columnar child: one borrowed-column batch per projection, no data copy. |
spark.vecruntime.exec.sample.enabled |
true |
boolean | Convert SampleExec without replacement over a columnar child, running Spark's own Bernoulli sampler per partition so a seed returns Spark's rows. |
spark.vecruntime.exec.generate.enabled |
true |
boolean | Convert GenerateExec for explode / posexplode (and the _outer forms) over an array column or a struct field of one. |
spark.vecruntime.exec.window.enabled |
true |
boolean | Convert WindowExec (ranking functions, whole-partition, running and sliding aggregates, the offset functions) and the WindowGroupLimitExec Spark plans below a rank <= k filter; the child may be Spark's row sort. |
spark.vecruntime.exec.localTableScan.enabled |
false |
boolean | Convert LocalTableScanExec (VALUES, small local relations) into one batch per partition; off because there is nothing to accelerate, it only lets small-table tests run our operators. |
spark.vecruntime.exec.range.enabled |
true |
boolean | Convert RangeExec (spark.range(...), the range() table function) into native INT64 batches of spark.sql.inMemoryColumnarStorage.batchSize rows -- Spark's rows in Spark's partitions, one reused vector per task -- so the filter, projection and partial aggregate above range() are ours from the leaf (over Spark's row leaf they stayed Spark's until the first exchange); false keeps Spark's row leaf. |
spark.vecruntime.exec.broadcastHashJoin.enabled |
true |
boolean | Convert BroadcastHashJoinExec when the streamed side is columnar or an exchange; the build side is our columnar broadcast when the exchange below is ours, otherwise Spark's broadcast HashedRelation. |
spark.vecruntime.exec.broadcastExchange.enabled |
true |
boolean | Replace a join's BroadcastExchangeExec (hash or nested-loop) over one of our operators or a columnar source by VectorBroadcastExchangeExec (#325): the build side travels as Arrow batches and our broadcast joins build their table from them. A Spark consumer of the same exchange gets Spark's value (the relation or the rows), built from the batches on first use. false keeps Spark's exchange. |
spark.vecruntime.exec.broadcastNestedLoopJoin.enabled |
true |
boolean | Convert BroadcastNestedLoopJoinExec (joins without equi-keys) when the streamed side is columnar or an exchange. |
spark.vecruntime.exec.shuffledHashJoin.enabled |
true |
boolean | Convert ShuffledHashJoinExec; over Spark's row shuffle both inputs go through RowToColumnarExec. |
spark.vecruntime.exec.sortMergeJoin.mode |
auto |
off, hash, merge, auto |
What becomes of SortMergeJoinExec (#286, #287): off leaves it to Spark, hash re-expresses it as our shuffled hash join, merge plans our order-preserving merge join over Spark's sorted inputs, auto decides per join -- the merge join where a parent relies on the ordering, where the row order can reach a limit or a sort, or where the hash rewrite is not allowed (no statistics, both sides past spark.vecruntime.join.hashMaxBuildSize, a skew join), the hash rewrite where a side's statistics fit. Since #416 nothing is left to Spark's own operator under auto. |
spark.vecruntime.exec.sortMergeJoin.enabled |
unset | boolean | Compatibility alias for the mode: false reads as mode=off, anything else as auto; ignored when spark.vecruntime.exec.sortMergeJoin.mode is set. Prefer the mode key. |
Execution semantics
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.exec.strictFloatingPoint |
true |
boolean | Double sum / avg round exactly like Spark's (one accumulator per group, rows added in order); false uses lane-parallel and interleaved partial sums that differ in the last bits, buying back about 7% of aggregate kernel time (2.5% of TPC-H Q1 at SF10) but able to make an equality between two double sums fail (TPC-H Q15). The benchmark configurations run with false, Comet's default. |
spark.vecruntime.exec.selection.enabled |
true |
boolean | Pass a selection bitmap between our operators instead of compacting each filtered batch (compaction still happens below vecruntime.selection.minFraction survivors and at the boundary to Spark); a debugging switch. |
spark.vecruntime.agg.dictionaryKeys |
true |
boolean | The grouped aggregate emits its UTF8 keys dictionary-encoded, with ids over the group table's own dictionary (#377), so the shuffle writer stages ids and the final aggregate maps dictionary entries rather than rows; false emits plain strings. |
Memory budgets and spilling
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.agg.spillThreshold |
1g |
size, 0 = never |
Hard cap on one grouped aggregate table; below it the operator acquires its real footprint from Spark's task memory manager and acts on a refusal (#363, #367). Past the budget a buffer-emitting mode emits its table and starts over, a merging mode spills into hash buckets; 0 keeps everything in memory. |
spark.vecruntime.agg.spillBuckets |
16 |
int (at least 2) | Buckets a merging (Final) aggregate spills its table into and merges one at a time; each bucket is merged in memory, so the buckets, not the input, must fit. |
spark.vecruntime.agg.passThroughRatio |
1.5 |
double, 0 = never |
A partial aggregate whose full table reduced its input by less than this factor stops aggregating and passes each batch on to the exchange (#376); 0 keeps aggregating whatever the ratio. |
spark.vecruntime.sort.runRows |
1048576 |
positive int | Rows per sorted run (#285): the sort orders each run as the partition arrives and k-way merges the runs on output, bounding its JVM scratch to the run rather than the partition. |
spark.vecruntime.sort.spillBytes |
1g (1 GiB) |
size, 0 = off |
The sort's memory budget per task (#416): sealed runs past it are written to local disk in sorted order and merged from there, which is what lets the merge join take inputs of any size. Measured alone at 1 TB in the full run's shape, the 32 MB of #451 cost q14a 30% and q4 50% against 1 GiB because spilled runs share the node disk with the shuffle files, hence the 1 GiB default. |
spark.vecruntime.join.maxBuildSize |
1 GiB, or spark.memory.offHeap.size / spark.executor.cores when off-heap is enabled |
size | The largest build side, by statistics, the two broadcast joins convert for (#86); they hold the relation in memory per task for the whole query. The shuffled hash join is no longer gated by it (#416). |
spark.vecruntime.join.spillBytes |
256m |
size, 0 = never split |
Build bytes a shuffled hash join holds in memory before it splits both sides into buckets on disk (#416, the grace hash join). At 1 TB on the eight join-heaviest TPC-DS queries 256 MB and 1 GiB were the same within the band and neither spilled, so the smaller one leaves the headroom to the sort and the operators beside the join. |
spark.vecruntime.join.spillBuckets |
32 |
int | Buckets a split shuffled hash join writes each side into and joins one at a time (#416); 1 or less turns the split off. |
spark.vecruntime.join.hashMaxBuildSize |
spillBuckets x spillBytes (8 GiB) |
size, 0 = no cap |
The most a sort-merge join's build side may weigh per task, by statistics, for mode=auto to make it the hash join: within spark.vecruntime.join.spillBytes it builds in memory, within this cap it splits into buckets once, past it -- or without an estimate -- the merge join over the spilling sort takes it (#416). |
spark.vecruntime.join.denseKeys |
true |
boolean | A hash join on one INT32/INT64 key whose build values span at most 10x the distinct keys (and at most 4 M values) is probed through an array indexed by the key, not the hash table (#546): a range check and one load per row, 18-23x the hash probe in JoinProbeBenchmark. Spark's own LongHashedRelation switches to a dense array under the same rule. |
spark.vecruntime.join.buildPayload |
true |
boolean | A broadcast hash join whose build side carries payload columns without a lane (arrays, maps, structs) keeps them in a row store beside the lane columns and converts (#547); keys and the join condition must still read lane columns only. Off: every build column needs a lane, as before. The shuffled hash join's build side still needs lanes. |
Scan
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.scan.prefetch |
0 (off) |
int 0-8 | Depth of the prefetching converter's queue (#403, lever 2): 1 or 2 inserts VectorPrefetchScanExec between a Spark vectorized file scan (Parquet, Iceberg's BatchScanExec; not a Comet scan) and the first operator of ours above it, whose helper thread pulls and converts the reader's next batch while the task thread works on the previous one. Memory grows by that many converted batches per task; larger values are accepted and capped at 8. |
Columnar shuffle (the vecruntime-shuffle jar)
spark.vecruntime.shuffle.enabled is a session key; the rest are read from the SparkConf by the
shuffle manager, the writer and the Flight server at start-up, so set them on spark-submit. None
of them has an effect unless spark.shuffle.manager names VectorShuffleManager (see the last
section) and the shuffle jar is on the classpath.
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.shuffle.enabled |
true |
boolean | Replace a ShuffleExchangeExec above one of our operators with VectorShuffleExchangeExec (Arrow IPC record batches per reduce partition, #288) for hash, round-robin, single and range partitioning; Comet's native shuffle takes precedence where it is configured. |
spark.vecruntime.shuffle.backend |
flight |
flight, block, or a class name |
How a reducer fetches a remote map output: flight opens one Arrow Flight DoGet per remote executor carrying all of its blocks for the reducer, block uses Spark's block transfer over the same files (use it where TLS is required, since the Flight server has none), or the class name of a VectorShuffleBackend from another jar (where a push-based service would plug in). |
spark.vecruntime.shuffle.compression |
zstd |
zstd, lz4, none |
Body compression of the shuffle's record batches; lz4 is Arrow's pure-Java codec and an order of magnitude slower, none writes about 1.8x Spark's lz4 bytes. |
spark.vecruntime.shuffle.batchRows |
8192 |
int | Rows a map task holds per reduce partition before writing them as one record batch. |
spark.vecruntime.shuffle.batchBytes |
1m |
size | Bytes a map task holds per reduce partition before writing them as one record batch (whichever of rows or bytes fills first). |
spark.vecruntime.shuffle.bufferBytes |
64m |
size | Cap on what one map task holds across all reduce partitions before flushing, measured on the writer's allocator rather than estimated (#340). |
spark.vecruntime.shuffle.flushBytes |
1m |
size | Serialised bytes a map task keeps in memory per reduce partition before spilling that partition's stream to a temporary file merged into the data file at commit. |
spark.vecruntime.shuffle.writer.memoryLimit |
1g |
size | The hard limit of one map task's Arrow allocator; the writer flushes long before it, this is the backstop behind bufferBytes, and hitting it fails the task with a serializable memory error. |
spark.vecruntime.shuffle.writer.dictionaryMaxRatio |
0.5 |
double 0-1 | A string column is dictionary-encoded on the wire only when distinct values / rows in the record batch is at most this (#356); 1 always encodes, 0 never. |
spark.vecruntime.shuffle.rebalance.rowSizing |
true |
boolean | Session key. A rebalance exchange (the one a data source such as Iceberg asks for ahead of its write) hands AQE partition sizes proportional to rows instead of our compressed bytes, whose bytes per row vary several-fold between kinds of rows (#485). Reducers still fetch by the real sizes. |
spark.vecruntime.shuffle.rebalance.advisoryScaling |
true |
boolean | Session key. A rebalance's advisory partition size is sized for Spark's shuffle (Iceberg: target file size x Spark's expected shuffle compression); once the map stage has written, it is scaled by our uncompressed bytes per row over an estimate of Spark's UnsafeRow bytes for the same rows, so AQE's pieces hold about as many rows as Spark's (#20: the CDC MERGE's write stage took 40-46 s at Iceberg's 384 MiB and Spark's ~24 s at a hand-set 128 MiB). A size set with spark.sql.iceberg.advisory-partition-size is kept as given, in our bytes; one set as an Iceberg write option or table property is scaled like the default -- set this key to false to keep it exact. |
spark.vecruntime.shuffle.aqe.mapSizeScaling |
true |
boolean | Session key; false restores AQE on our real bytes. Every exchange of ours other than a rebalance reports its map output sizes to AQE times Spark's estimated on-disk bytes per row over ours (UnsafeRow estimate with measured string bytes, over sparkCompressionRatio, bounded to [1/16, 16]), so coalescing, skew detection and join groups see about Spark's bytes for the same rows. Only the statistics change; reducers fetch the real bytes (#511: our shuffle is 1.6-4.3x smaller than Spark's at TPC-DS 1 TB, so AQE packed up to 2x Spark's rows into a task -- q67's final aggregate ran 150 tasks where Spark's ran 300, and spilled). |
spark.vecruntime.shuffle.aqe.sparkCompressionRatio |
0 |
double | The compression mapSizeScaling expects of Spark's shuffle (uncompressed UnsafeRow bytes over bytes on disk). 0 or negative, the default, assumes Spark's shuffle compresses as well as ours and takes the uncompressed ratio. On TPC-DS 1 TB (Graviton4, advisory 128m) that gave 1,835.6 s against 1,986.4 s unscaled; 2.5 (Spark's measured 2.57) gave 1,973.8 s, because the larger factor stopped AQE merging the short queries' small partitions. |
spark.vecruntime.shuffle.flight.bindHost |
the executor's block manager host | host name | The address the executor's Flight server binds to; change it when the executor's advertised host is not the one it can bind. |
spark.vecruntime.shuffle.flight.threads |
max(4, available processors) |
int | Serving threads of the executor's Flight server. |
Comet integration
These matter only with Comet's jar on the classpath; see docs/comet.md.
| Key | Default | Type / values | What it does |
|---|---|---|---|
spark.vecruntime.comet.shuffle.enabled |
true |
boolean | Rewrite a Spark exchange (or Comet's row-based columnar exchange) above one of our operators into Comet's native shuffle over the zero-copy VectorToCometExec bridge; requires Comet's shuffle manager and spark.comet.exec.shuffle.enabled=true. |
spark.vecruntime.comet.shuffle.range.enabled |
true |
boolean | Also hand range-partitioned exchanges (global ORDER BY) to Comet's native shuffle, which then samples the child with Spark's RangePartitioner as Spark's own exchange does; gated as well by Comet's spark.comet.shuffle.native.partitioning.range.enabled. |
spark.vecruntime.comet.mixed.enabled |
false |
boolean | Mixed chains (#280): a Spark operator left to Spark whose children are ours is offered to Comet's native operator through the sink leaf, for the operator kinds spark.vecruntime.comet.preferComet names. |
spark.vecruntime.comet.preferComet |
empty | comma-separated operator kinds, optionally qualified | The allowlist of the mixed pass (#281): kinds (filter, project, sort, sortMergeJoin, hashJoin, broadcastHashJoin, window, expand, union, limit, all) optionally qualified (project:wideDecimal, filter:strings, sort:estimatedRows>1000000); a listed operator is offered to Comet first and ours steps aside, one Comet declines is ours after all. Empty means mixed plans allowed, none requested; the #281 study found no entry that meets docs/comet.md's three-part rule. |
Non-spark.vector keys the plugin needs
- JVM options on driver and executors. The kernels use the incubating Vector API and the
FFM API, so both JVMs need
--add-modules=jdk.incubator.vector --enable-native-access=ALL-UNNAMEDinspark.driver.extraJavaOptionsandspark.executor.extraJavaOptions(--sun-misc-unsafe-memory-access=allowsilences theUnsafedeprecation warnings from Spark and Arrow on JDK 25). Every batch our operators produce, the shuffle's record batches and the Flight server's buffers are Arrow direct memory, bounded by-XX:MaxDirectMemorySize, which defaults to the heap size: set it to the executor's memory overhead less what the JVM itself needs (the 1 TB campaign ran a 30 GB heap and 20 GB of overhead per 13-core executor, the cluster manifests setting the bound to the overhead less 2 GB; see "Memory tuning" inREADME.md). - Registering the plugin.
spark.plugins=io.vecruntime.spark.VectorPluginregisters the session extension and attaches the UI tab; alternativelyspark.sql.extensions=io.vecruntime.spark.VectorSparkSessionExtensionsinjects the planner rule alone. - The columnar shuffle manager.
spark.shuffle.manager=org.apache.spark.sql.vecruntime.shuffle.VectorShuffleManager(from thevecruntime-shufflejar) is what makesspark.vecruntime.shuffle.enabledtake effect; the manager serves our dependencies with the Arrow IPC writer and reader and delegates every other shuffle to Spark's sort shuffle. Withspark.authenticateon, the Flight server requires Spark's shuffle secret as a bearer token; underspark.ssl.rpc.enabledit refuses to start, so usespark.vecruntime.shuffle.backend=blockthere. - Comet's shuffle instead. To feed Comet's native shuffle from our operators set
spark.shuffle.managerto Comet'sCometShuffleManagerandspark.comet.exec.shuffle.enabled=true. - Kernel tuning knobs are JVM system properties, not Spark confs, because the kernels have no
Spark dependency:
vecruntime.vectorBits,vecruntime.platform,vecruntime.agg.interleave,vecruntime.agg.maskPathMaxGroups,vecruntime.selection.minFraction,vecruntime.agg.plainDictMaxEntries;README.mddescribes each.