ETL CompareETL Compare

University · Glossary

Spark ETL performance glossary

In brief

Thirty-eight terms cover most of what Spark accelerator vendors and projects talk about, from shuffle, disk spill and data skew to fallback, off-heap memory, adaptive query execution and TPC-DS. Each entry gives a short definition, a link to a longer page on ETL Compare, and a source link where the term is defined in project or cloud provider documentation.

Adaptive query execution (AQE)
A Spark SQL feature that re-optimizes a query plan during execution using runtime statistics. It is on by default (spark.sql.adaptive.enabled) and can coalesce small shuffle partitions and split skewed join partitions. Read more Source: Spark SQL tuning
Apache Arrow
A columnar in-memory data format. Apache DataFusion describes itself as a query engine "that uses Apache Arrow as its in-memory format", which is why Comet, Auron and Flarion share an Arrow-based execution layer. Read more Source: datafusion.apache.org
Apache Iceberg
An open table format for large analytic tables. Iceberg ships Spark procedures for table maintenance, including rewrite_data_files for compaction, rewrite_manifests and expire_snapshots. Read more Source: Iceberg Spark procedures
Apache Parquet
A columnar file format widely used for Spark tables. Several accelerators document native Parquet handling: Comet documents native Parquet scans, Photon scans Parquet natively, and the RAPIDS Accelerator documents Parquet writing on GPUs. Read more
Bloom filter join
A join optimization that builds a compact Bloom filter from the join keys on one side and uses it to drop non-matching rows on the other side before they are shuffled. AWS lists Bloom filter joins among the Spark performance optimizations in the Amazon EMR runtime. Read more Source: Amazon EMR Spark performance
Broadcast join
A join in which Spark copies the smaller table to every executor instead of shuffling both sides. Tables below spark.sql.autoBroadcastJoinThreshold (default 10 MB) are broadcast automatically. Read more Source: Spark SQL tuning
Columnar execution
Processing data column by column in memory rather than row by row. Read more
Compaction
Rewriting many small data files into fewer, larger ones. In Apache Iceberg the rewrite_data_files procedure does this with a binpack, sort or zorder strategy and a default target file size of 536870912 bytes (512 MB). Read more Source: Iceberg Spark procedures
Data skew
Uneven distribution of data across partitions, so a few tasks process far more data than the rest and the stage waits on them. Read more
Databricks Unit (DBU)
The unit in which Databricks meters compute. Databricks documents that Photon compute consumes DBUs at a different rate from non-Photon compute, so the cost effect of turning Photon on has to be worked out per job. Read more Source: Photon documentation
DataFusion
Apache DataFusion, an "extensible query engine written in Rust that uses Apache Arrow as its in-memory format". Comet, Auron and Flarion build on it. Read more Source: datafusion.apache.org
DataFusion Comet
Apache DataFusion Comet, an open-source Spark accelerator that translates Spark physical plans to DataFusion physical plans. Its 1.0.0 release came out on 7 August 2026. Read more Source: Comet 1.0.0 release notes
Disk spill
When a Spark task runs out of execution memory and writes intermediate data to local disk, then reads it back. Read more Source: Spark web UI
Dynamic partition pruning
A runtime optimization that uses the filter values from one side of a join to skip partitions of a partitioned table on the other side, so Spark reads less data. AWS lists it among the Spark optimizations in the Amazon EMR runtime. Read more Source: Amazon EMR Spark performance
Event log
The record of a Spark application's events that lets the Spark UI be rebuilt after the job ends. It is turned on with spark.eventLog.enabled and written to the directory set in spark.eventLog.dir. Read more Source: Spark monitoring
Executor
A process launched for a Spark application on a worker node that runs the application's tasks and holds its execution memory. Tasks running on the same executor share that memory, which is one reason spill appears even on large executors. Read more
Fallback
What an accelerator does with an operation it cannot run natively: it hands that part of the plan back to standard Spark execution. The RAPIDS Accelerator, Photon, Gluten and Comet all document automatic fallback. Read more Source: RAPIDS Accelerator FAQ
GC time
A task metric in the Spark UI: the time the executor's JVM spent in garbage collection while the task ran. High GC time relative to task duration often means tasks are short of memory. Read more Source: Spark web UI
Gluten
Apache Gluten, a top-level Apache project that offloads Spark SQL execution to native engines such as Velox. Read more Source: gluten.apache.org
History Server
The Spark component that rebuilds the web UI of finished applications from their event logs. It listens on port 18080 by default and serves the same data as JSON under /api/v1. Read more Source: Spark monitoring
Native engine
An execution engine compiled to machine code that runs outside the JVM, such as Velox (C++), Photon (C++) or DataFusion (Rust). Spark accelerators hand supported parts of the physical plan to a native engine and fall back to standard Spark for the rest. Read more
Off-heap memory
Memory that Spark manages outside the JVM heap, turned on with spark.memory.offHeap.enabled and sized with spark.memory.offHeap.size. Gluten's Velox backend and Comet both run in this pool, so it has to be sized for them. Read more Source: Comet tuning guide
Partition
A slice of a dataset processed by one task. File reads pack up to spark.sql.files.maxPartitionBytes (default 128 MB) into one partition, and shuffles produce spark.sql.shuffle.partitions (default 200) partitions unless adaptive query execution coalesces them. Read more Source: Spark SQL tuning
Photon
Databricks' native vectorized query engine, part of the Databricks Runtime. Read more Source: Photon documentation
Physical plan
The concrete set of operators Spark executes for a query, chosen after optimization. Accelerators work at this level: the RAPIDS Accelerator FAQ describes replacing "parts of the physical plan that Apache Spark considers internal". Read more Source: RAPIDS Accelerator FAQ
Qualification Tool
An NVIDIA tool for the RAPIDS Accelerator that reads event logs from existing CPU Spark runs and estimates which applications are good candidates for GPU acceleration. Read more Source: RAPIDS Accelerator user guide
Shuffle
The redistribution of data across executors between stages, needed by joins, aggregations and repartitioning. Read more Source: Spark web UI
Shuffle manager
The pluggable Spark component that writes and reads shuffle data. Accelerators often replace it: Gluten with ColumnarShuffleManager, Comet with CometShuffleManager, the RAPIDS Accelerator with the RAPIDS Shuffle Manager and Auron with its own shuffle manager. Read more Source: Gluten Velox getting started
Shuffle Read Fetch Wait Time
A task metric in the Spark UI: the time a task spent blocked waiting for shuffle data to arrive from other executors. High values on a slow stage point to a shuffle bottleneck rather than slow computation. Read more Source: Spark web UI
Skew join optimization
The part of adaptive query execution that splits oversized partitions in sort-merge joins. A partition counts as skewed when it is larger than spark.sql.adaptive.skewJoin.skewedPartitionFactor (default 5.0) times the median partition size and also larger than spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes (default 256MB). Read more Source: Spark SQL tuning
Small files problem
Too many small input or output files, so Spark spends its time opening files and planning tasks rather than processing data. Read more
Spark plugin
Code loaded into a Spark application through configuration, such as the spark.plugins setting, without changing application code. Gluten, for example, is enabled with spark.plugins=org.apache.gluten.GlutenPlugin. Read more Source: Gluten Velox getting started
Stage
A set of tasks that run the same work on different partitions, bounded by shuffles. The Spark UI Stages tab reports per-stage metrics such as shuffle read and write, spill and GC time. Read more Source: Spark web UI
Straggler task
A task that runs much longer than the median task in its stage, usually because its partition holds far more data. The stage cannot finish until the straggler does. Read more
Task
The unit of work Spark runs on one partition within one stage. A stage has as many tasks as it has partitions. Read more Source: Spark web UI
TPC-DS and TPC-H
Decision-support benchmark specifications from the Transaction Processing Performance Council. Spark projects usually publish results from query suites derived from them; Comet's benchmarking guide says it runs "benchmarks derived from TPC-H and TPC-DS". Read more Source: Comet benchmarking guide
Vectorized execution
Processing a batch of values with one operation instead of one value at a time. Read more
Velox
An open-source C++ execution engine created at Meta, used by Gluten as a Spark backend. Read more Source: velox-lib.io

New to Spark internals? Start with the lessons in the Basics track.