Spark SQL was built around rows. Every operator in a classic physical plan consumes and produces InternalRow objects, and whole-stage code generation fuses those operators into one tight loop over rows. Yet the files Spark reads most often, Parquet and ORC, store data by column, and the fastest modern engines process data by column too. Columnar processing in Spark is the set of data structures, readers, plan rules and extension points that let data stay in column form for as long as it pays to.
This article explains how that works in open-source Spark, which parts are columnar by default and which are not, how to read a plan to see where data switches between rows and columns, how plugins replace whole stages with columnar implementations, and what goes wrong in practice. Configuration names and defaults were checked against the Spark 4.2.0 documentation on 2026-10-03. Databricks Photon is a separate proprietary engine and is covered in its own article.
The data path in one picture
Why columns are faster
Consider a table with 80 columns and a query that touches three of them. A row layout stores each record contiguously, so a scan reads all 80 fields and discards 77. A column layout stores each column contiguously, so the scan reads three columns and skips the rest. That saving is about IO and is delivered by the file format alone.
The second saving is about the CPU. When values of one type sit next to each other in memory, a loop over them is simple: load, compare, store, repeat. The branch predictor sees the same pattern every iteration, the cache prefetcher streams the data, and the compiler can use SIMD instructions that apply one operation to several values at once. Row-at-a-time code, by contrast, jumps between fields of different types and calls virtual methods per value. Columnar formats also compress well, because a column of similar values suits dictionary and run-length encoding, and some operations can work on the encoded form directly.
The cost is that some operations are naturally row-shaped. Building a hash table for a join, sorting records, or calling a user function that takes a whole record all want the fields of one row together. Every engine therefore has boundaries where data is pivoted between layouts, and much of the engineering in Spark's columnar support is about putting those boundaries in the right place.
ColumnVector and ColumnarBatch
Two classes in org.apache.spark.sql.vectorized carry columnar data through Spark. A ColumnVector holds one column for a run of rows, with typed accessors such as getInt(rowId), getUTF8String(rowId) and isNullAt(rowId), plus hasNull() and numNulls() so code can skip null checks when a vector has none. A ColumnarBatch is a set of column vectors that share a row count, exposed through numRows() and column(i).
Spark's own implementations are writable vectors backed either by Java arrays on the heap or by off-heap memory, chosen with spark.sql.columnVector.offheap.enabled (default false). Arrow-backed vectors wrap Arrow buffers so that data exchanged with Python or an external engine is not copied. Because ColumnVector is an abstract class, a plugin can supply its own implementation, for example one whose buffers live in GPU memory. Reading a batch directly looks like this:
import org.apache.spark.sql.vectorized.ColumnarBatch
// Sum a double column, skipping the null check when the vector has none.
def sumColumn(batch: ColumnarBatch, ordinal: Int): Double = {
val vec = batch.column(ordinal)
val n = batch.numRows()
var total = 0.0
var i = 0
if (!vec.hasNull()) {
while (i < n) { total += vec.getDouble(i); i += 1 } // tight, branch-free loop
} else {
while (i < n) { if (!vec.isNullAt(i)) total += vec.getDouble(i); i += 1 }
}
total
}The two loops show the point of the layout. When the vector reports no nulls, the loop body is a single load and add over contiguous doubles, which the JIT compiler can unroll and often vectorize. A row-based equivalent would fetch a field from a different offset in each record.
Where stock Spark is columnar
In stock Spark, columnar data appears in a few specific places rather than throughout the plan. The table lists them with the settings that govern them.
| Place | What is columnar | Key settings (Spark 4.2 defaults) |
|---|---|---|
| Parquet scan | Pages decoded straight into column vectors, nested types included | spark.sql.parquet.enableVectorizedReader true; ...enableNestedColumnVectorizedReader true; ...columnarReaderBatchSize 4096 |
| ORC scan | Same idea for ORC stripes | spark.sql.orc.enableVectorizedReader true; nested reader true; batch size 4096 |
| In-memory cache | Cached tables stored as compressed column batches | spark.sql.inMemoryColumnarStorage.compressed true; ...batchSize 10000 |
| Python exchange | Arrow record batches to and from pandas UDFs and toPandas | spark.sql.execution.arrow.pyspark.enabled true |
| Plugins | Whole operators replaced with columnar ones | spark.plugins, spark.sql.extensions |
Everything between the scan and the output, the filters, projections, aggregations, joins and sorts, runs on rows in open-source Spark. The vectorized reader still matters a great deal, because decoding Parquet is often the most expensive part of a scan-heavy query and a batch decoder is much faster than a row-at-a-time one. But it is worth being precise: turning on the vectorized reader does not make Spark a vectorized engine. It makes the scan columnar, and a transition then converts batches to rows for the rest of the stage.
Transitions and how to read them in a plan
Every physical operator declares whether it can produce columnar output through supportsColumnar and whether it produces rows through supportsRowBased. After planning, a rule called ApplyColumnarRulesAndInsertTransitions walks the plan and inserts a ColumnarToRowExec wherever a columnar operator feeds a row-based one, and a RowToColumnarExec wherever the reverse happens. You can see them in any plan:
spark.read.parquet("s3://lake/sales").filter("amount > 100") \
.groupBy("country").sum("amount").explain()
== Physical Plan == (abridged)
*(2) HashAggregate(keys=[country], functions=[sum(amount)])
+- Exchange hashpartitioning(country, 200)
+- *(1) HashAggregate(keys=[country], functions=[partial_sum(amount)])
+- *(1) Filter (isnotnull(amount) AND (amount > 100.0))
+- *(1) ColumnarToRow
+- FileScan parquet [country,amount] Batched: true, PushedFilters: [...]Read it bottom up. Batched: true on the scan means the vectorized reader is active and emits ColumnarBatch objects. ColumnarToRow sits inside whole-stage codegen stage 1, marked by the *(1) prefix, because it supports code generation: the generated loop reads values straight out of the column vectors by row index and feeds them to the filter and partial aggregate without materialising intermediate row objects. That is why the transition is cheap in the common case. If Batched: false appears instead, the reader is producing rows; check whether a setting disabled vectorized reading or the source does not support it.
The expensive case is a transition in the middle of a plan, especially RowToColumnar, which must copy every value of every row into fresh vectors. That appears when a columnar plugin handles some operators but not the one beneath them, and it can erase the plugin's benefit.
The ColumnarRule extension API
The extension point for columnar execution has been part of Spark since 3.0. A ColumnarRule has two hooks, preColumnarTransitions and postColumnarTransitions, each a Rule[SparkPlan]. The first runs before transitions are inserted and is where a plugin swaps operators for columnar versions; the second runs after and is where it cleans up, for example by replacing Spark's generic transitions with its own faster ones. A columnar operator overrides supportsColumnar to return true and implements doExecuteColumnar(), returning an RDD[ColumnarBatch]. The skeleton below shows the wiring; the operator body is where real plugins do their work.
import org.apache.spark.sql.{SparkSession, SparkSessionExtensions}
import org.apache.spark.sql.catalyst.rules.Rule
import org.apache.spark.sql.execution.{ColumnarRule, FilterExec, SparkPlan}
class MyColumnarExtensions extends (SparkSessionExtensions => Unit) {
override def apply(ext: SparkSessionExtensions): Unit =
ext.injectColumnar(session => new MyColumnarRule(session))
}
class MyColumnarRule(session: SparkSession) extends ColumnarRule {
override def preColumnarTransitions: Rule[SparkPlan] = plan => plan.transformUp {
// Replace only what we can run end to end; leave everything else as rows.
case f: FilterExec if MyFilterExec.supports(f.condition, f.child.output) =>
MyFilterExec(f.condition, f.child) // supportsColumnar = true, doExecuteColumnar()
}
}
// Enable with: --conf spark.sql.extensions=com.example.MyColumnarExtensionsTwo design rules follow from how transitions are inserted. First, a replacement only helps if its neighbours are columnar too; replacing one filter between two row operators adds two transitions and gains nothing. Second, an operator must either support every expression in it or decline, because a half-supported operator is a correctness risk. Production plugins therefore carry large tables of supported expressions and types, and fall back to Spark for anything else.
Columnar plugins
Several projects use this machinery to run most of a query outside the row engine. The NVIDIA RAPIDS Accelerator for Apache Spark runs supported operators on GPUs with columnar data in device memory. Apache Gluten offloads operators to native engines, with Velox and ClickHouse backends. Apache DataFusion Comet runs operators natively on the Apache DataFusion query engine. They share a pattern: a plugin registered through spark.plugins or spark.sql.extensions, a planning rule that converts supported subtrees, columnar shuffle so batches cross stage boundaries without becoming rows, and fallback to Spark for unsupported expressions.
Speed-ups are workload-specific, and every project publishes its own benchmarks, so measure on your queries. The two numbers to watch are the fraction of the plan that actually ran in the plugin, which each project reports through its own explain or fallback logging, and the number of transitions left in the final plan. A plugin that converts 60 percent of operators but leaves a transition pair around every join can be slower than stock Spark.
Worked example: a CPU-bound nightly job
A team runs a nightly job over a 4 TB Parquet events table with 80 columns. It reads five columns, filters to one day, applies a Python function that normalises a URL, and aggregates by page. The job takes 70 minutes and the executors show high CPU but modest IO.
Step 1, read the plan. explain() shows FileScan parquet ... Batched: true followed by ColumnarToRow, so the scan is vectorized. Above it sits BatchEvalPython, the operator for an ordinary row-at-a-time Python UDF. Each row is serialised, shipped to a Python worker, processed and shipped back, which is where the CPU is going.
Step 2, remove the row boundary. URL normalisation turns out to be lower-casing, stripping a query string and trimming a trailing slash, all of which built-in functions such as lower and regexp_replace express, so the UDF disappears and the whole stage stays in generated code. Where logic truly needs Python, a pandas UDF moves the data as Arrow batches instead, shown in the plan as ArrowEvalPython.
Step 3, check batch memory. One of the five columns holds raw JSON payloads averaging 200 KB. A Parquet batch of 4,096 rows therefore needs about 800 MB for that column alone, per task. With four concurrent tasks per executor that is over 3 GB of vectors before any other work, and some tasks fail with out-of-memory errors on large row groups. Dropping the payload column from the projection, or lowering spark.sql.parquet.columnarReaderBatchSize to 512 for this job, removes the failures.
Step 4, measure. Compare stage time, task CPU time and peak execution memory before and after each change, one change at a time, on the same input partition. Record the plan so the next person can see which boundary was removed.
Failure modes
- Silent row fallback. A setting, an unsupported source or an older connector turns
Batched: trueintoBatched: falseand the scan slows down with no error. Assert on the plan in a test for critical jobs. - Transition sandwiches. A plugin converts some operators, leaving
RowToColumnarandColumnarToRowpairs around the rest. Count transitions in the final plan. - Batch memory blow-ups. Batch size is in rows, so wide string or binary columns make each batch huge. Size batches from bytes, not rows.
- Python UDF boundaries. Row-at-a-time Python UDFs pivot every row through serialisation. Prefer built-in functions, then pandas UDFs.
- Over-wide schemas. Whole-stage codegen is deactivated above
spark.sql.codegen.maxFields(100) fields including nested ones, which removes the generated loop the cheap transition depends on. Project early. - Cache confusion. The in-memory columnar cache is not Parquet; its batch size, 10,000 rows, and its compression are separate settings, and caching a wide table can cost more memory than expected.
Trade-offs
Columnar execution wins on scans, filters, projections and aggregations over many rows of a few columns, which is most analytical SQL. It wins less, or loses, on narrow lookups, row-oriented user code, and plans that pivot repeatedly between layouts. Stock Spark's choice, columnar scans with row execution fused by code generation, is a reasonable middle ground that needs no extra dependency. Native and GPU plugins buy more speed at the cost of another engine to operate: version compatibility with each Spark release, differences in edge-case semantics such as floating-point and decimal behaviour, extra memory configuration, and fallback paths that need watching. Adopt one when profiling shows CPU-bound operators that it supports, and keep a switch to turn it off per job.
For more depth on the layers around this one, see Spark and Arrow for the interchange format, whole-stage code generation for how row operators are fused, Tungsten for the row memory format, columnar database architecture for encodings and pruning, and Hive vectorization for a comparable design in another engine.
What to do next
- Run
explain()on your three most expensive jobs and recordBatchedon every scan and every transition operator. - Add a test that fails if a critical job's scan stops being batched.
- Replace row-at-a-time Python UDFs with built-in functions or pandas UDFs, and confirm the plan changed.
- Estimate batch bytes per task from average column widths times
columnarReaderBatchSize, and lower the batch size for jobs with wide binary or string columns. - Project only needed columns before anything else, and keep schemas under the code generation field limit.
- If you trial a columnar plugin, record the share of operators it ran and the transitions left, then compare stage CPU time against stock Spark on the same input.