Most Spark jobs spend their time scanning columnar files, decoding them, filtering, joining, aggregating, sorting and shuffling. Those are data-parallel operations over millions of values of the same type, which is exactly what a GPU is built for: thousands of simple cores and memory bandwidth measured in terabytes per second rather than the tens to low hundreds of gigabytes per second a CPU socket gets. The catch is everything around the computation: data has to cross a PCIe link to reach the GPU, GPU memory is small compared with host memory, and any operator the GPU cannot run forces data back to the CPU.
This article explains how GPU acceleration for Spark SQL and DataFrames actually works, using NVIDIA's plugin, now documented as the NVIDIA cuDF plugin for Apache Spark and previously known as the RAPIDS Accelerator for Apache Spark. We follow a task through the hardware, configure a cluster, size memory and concurrency, find and remove CPU fallbacks, handle shuffle and spill, and work through when it saves money and when it does not.
What actually runs on the GPU
The plugin is a jar loaded with spark.plugins=com.nvidia.spark.SQLPlugin. It does not change your code. After Catalyst produces a physical plan, the plugin walks it and replaces each operator it supports, such as a file scan, filter, project, hash aggregate, sort-merge or broadcast join, window or shuffle exchange, with a GPU version implemented on cuDF, the RAPIDS GPU DataFrame library. Where an operator or expression is not supported, it stays on the CPU and the plugin inserts transitions between the columnar GPU format and Spark's row format around it.
The unit of work on the GPU is a columnar batch, one buffer per column plus null bitmaps and string offsets, targeted at spark.rapids.sql.batchSizeBytes, default 1 GiB. That is far larger than a CPU vectorised batch, because a GPU kernel needs millions of values to keep thousands of cores busy. The RDD API, Python UDFs and arbitrary Scala closures do not go through this path; the acceleration is for SQL and the DataFrame API. Operator and expression coverage changes between releases, so check the supported-operators documentation for your version. The current user guide lists Spark 3.3.0 through 3.5.8 on Scala 2.12, and 3.5.x, 4.0.x and 4.1.1 on Scala 2.13.
Following one task through the hardware
File bytes still arrive through CPU threads: the plugin's multithreaded readers fetch and decompress Parquet or ORC pages on the host, and in many cases coalesce small files into one larger batch. The bytes are copied to GPU memory over PCIe, and the GPU does the decoding, filtering and everything downstream. Three hardware facts shape the configuration.
- PCIe is the narrow pipe. A PCIe Gen4 x16 link gives roughly 25 GB/s per direction in practice, and Gen5 roughly double that, against GPU memory bandwidth in the terabytes per second. Every trip across the link must buy a lot of computation to be worth it.
- Pinned memory makes transfers fast. DMA engines can only copy from page-locked host memory. From pageable memory the driver first copies into a pinned staging buffer, which roughly halves the effective rate.
spark.rapids.memory.pinnedPool.sizedefaults to 0, so set it, typically to a few GiB per executor. - GPU memory is the scarce resource. A 16 to 80 GB device is shared by every task running on it. What does not fit must be spilled to host memory or disk, which costs PCIe bandwidth again.
The GPU architecture itself, including why coalesced access to device memory matters, is covered in the GPU memory hierarchy article. For Spark, the practical rule is simple: keep data on the GPU across as many consecutive operators as possible.
Setting up a cluster
Spark's resource scheduling, available since Spark 3.0, lets executors request GPUs and tasks request fractions of a GPU. The plugin documentation recommends one GPU per executor. The fraction per task decides how many tasks the scheduler places on the executor; the plugin then decides separately how many may use the GPU at the same moment.
spark-submit \
--master yarn \
--jars rapids-4-spark_2.12-<version>.jar \
--conf spark.plugins=com.nvidia.spark.SQLPlugin \
--conf spark.executor.cores=16 \
--conf spark.executor.memory=48g \
--conf spark.executor.memoryOverhead=16g \
--conf spark.executor.resource.gpu.amount=1 \
--conf spark.task.resource.gpu.amount=0.0625 \
--conf spark.executor.resource.gpu.discoveryScript=./getGpusResources.sh \
--conf spark.rapids.memory.pinnedPool.size=4g \
--conf spark.sql.files.maxPartitionBytes=512m \
--conf spark.sql.adaptive.enabled=true \
nightly_orders_etl.pyspark.task.resource.gpu.amount is 1/16 here so all 16 cores can run tasks. Only some of those tasks are on the GPU at any instant; the rest do CPU work such as reading files or waiting on shuffle. The discovery script tells YARN and Kubernetes executors which device addresses they own. On Kubernetes the executor pod also requests nvidia.com/gpu through the device plugin. Managed platforms such as Dataproc, EMR and Databricks provide setup paths that install the jar and driver. Pin the plugin version to a CUDA driver version you have tested, because the jar is built against a CUDA major version.
Memory and concurrency: the four knobs that matter
The RMM pool. GPU allocations go through RAPIDS Memory Manager, because cudaMalloc per batch would be slow and would fragment memory. spark.rapids.memory.gpu.pool defaults to ASYNC, CUDA's stream-ordered allocator, and spark.rapids.memory.gpu.allocFraction defaults to 1.0 of free memory, with a small reserve held back. If anything else shares the GPU, such as a Python worker, reduce the fraction.
GPU task concurrency. A semaphore limits how many tasks hold the GPU at once. spark.rapids.sql.concurrentGpuTasks sets the starting value. In recent releases the plugin adjusts it dynamically per stage, and if you do not set it, it estimates a starting point from GPU memory. Too low and the GPU idles while tasks do I/O; too high and tasks fight for memory and spill. Two to four concurrent tasks is a common range on large devices.
Batch size. Larger batches keep kernels efficient but need more memory per task. Keep the default unless the profiling tool shows heavy spill, and raise spark.sql.files.maxPartitionBytes so each task reads enough data to fill a batch.
Spill. When GPU memory runs out, the plugin spills spillable buffers, first to host memory and then to local disk, and reads them back when needed. Spill keeps jobs alive, but spill-heavy stages can be slower than the CPU. In the Spark UI, the GPU operators report spill metrics; treat sustained spill as a sizing problem, not a normal state.
Fallbacks: the tax you must measure
Any unsupported operator or expression runs on the CPU, and data crosses PCIe twice around it, plus a columnar-to-row conversion and a row-to-columnar conversion. One fallback in the middle of a pipeline can erase the benefit for the whole stage. spark.rapids.sql.explain defaults to NOT_ON_GPU, which logs, at planning time, each part of the plan that stays on the CPU and why. Read those logs before tuning anything else.
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.getOrCreate()
spark.conf.set("spark.rapids.sql.explain", "NOT_ON_GPU") # the default, made explicit
orders = spark.read.parquet("gs://lake/orders/dt=2026-09-29/")
custs = spark.read.parquet("gs://lake/customers/")
daily = (orders
.filter(F.col("status") == "PAID")
.join(custs, "customer_id")
.groupBy("region", "segment")
.agg(F.sum("amount").alias("revenue"),
F.countDistinct("customer_id").alias("buyers")))
daily.explain() # Gpu* operators = on the GPU; plain names = CPU fallback
daily.write.mode("overwrite").parquet("gs://lake/marts/daily_revenue/dt=2026-09-29/")The usual causes are Python and Scala UDFs, and expressions or data types with no GPU implementation in your version. Fix them in order of cost: replace UDFs with built-in functions, often also a large CPU speed-up (see Python performance in Spark), and move an unavoidable UDF to the end of the pipeline so it runs once over a small result. Note the opposite case too: spark.rapids.sql.incompatibleOps.enabled defaults to true, so operations whose GPU results are not bit-for-bit identical run on the GPU. Floating-point aggregation is the classic case, because addition order differs. Use decimals for money, or set the flag to false and accept more fallbacks.
To assess a job before moving it, spark.rapids.sql.mode=explainOnly runs the plugin's planner on a CPU cluster and logs what would run on the GPU, without executing anything there.
Shuffle and adaptive execution
Shuffle is often where GPU jobs spend their time, because partitioning is fast on the GPU but the data must still be serialised, written, fetched and deserialised. The plugin ships its own shuffle manager. Its default mode, MULTITHREADED, parallelises shuffle writes and reads with thread pools on the host; a UCX mode can move blocks directly between GPUs over RDMA-capable networks, at the cost of installing and operating UCX. Recent releases configure the plugin's shuffle manager automatically on supported Spark versions unless spark.shuffle.manager is already set, so check the effective value in the environment tab rather than assuming.
Adaptive query execution works with the plugin and matters more than on the CPU. Fewer, larger partitions suit the GPU, and AQE coalescing small post-shuffle partitions and handling skewed joins keeps each GPU task working on a batch big enough to be efficient. Set spark.sql.shuffle.partitions lower than you would on a CPU cluster of the same core count, and let AQE split where needed.
Measure first: the qualification and profiling tools
NVIDIA ships two tools that read Spark event logs. The qualification tool runs over logs from your existing CPU jobs and estimates, per application, how much would run on the GPU and the likely speed-up and savings. The profiling tool runs over logs from a GPU run and recommends configuration changes and flags problems such as fallbacks and spill.
pip install spark-rapids-user-tools
spark_rapids qualification --platform dataproc --eventlogs gs://spark-logs/prod/2026-09/Treat the output as a ranking, not a promise. The jobs that qualify best share a profile: large scans of columnar files, joins and aggregations on wide tables, heavy shuffles, and few UDFs. Jobs dominated by small files, object-store latency, Python UDFs or tiny partitions qualify poorly, and moving them wastes GPU hours.
Worked example: is the nightly ETL worth moving?
A nightly job reads 4 TB of compressed Parquet, filters to 1.2 TB, joins with two dimension tables, aggregates to a few million rows and writes a mart. On the CPU it runs on 40 executors with 16 cores each and takes 95 minutes. The qualification tool reports 97 percent of task time in supported operators, with one Scala UDF that normalises phone numbers.
Step 1: replace the UDF with regexp_replace and substring, which the explain log confirms stay on the GPU. The CPU job alone drops to 88 minutes. Step 2: run on 10 executors, each with one GPU and 16 cores, pinned pool 4 GiB, maxPartitionBytes 512 MiB and 400 shuffle partitions with AQE. First run: 31 minutes, with 180 GB spilled in the join stage. Step 3: the profiling tool suggests reducing concurrent GPU tasks for that stage and adding executor memory overhead for host spill. Second run: 24 minutes, spill under 20 GB.
Cost, with illustrative prices: if a CPU node costs 1 unit per hour and a GPU node 3.5 units, the CPU job costs 40 × 88/60 ≈ 59 units and the GPU job 10 × 24/60 × 3.5 = 14 units, roughly a quarter. The same arithmetic on a job that spends half its time listing and opening small files would come out worse on GPUs, which is why you measure each job rather than migrating a cluster.
Failure modes and how to recognise them
| Symptom | Likely cause | Fix |
|---|---|---|
| Speed-up below 1.5x | Fallbacks mid-pipeline, or I/O-bound scans | Read NOT_ON_GPU logs; fix small files; check GPU utilisation |
| GPU OOM or heavy spill | Too many concurrent GPU tasks, or huge skewed partitions | Lower concurrency; enable AQE skew join; more partitions for that stage |
| Executor killed by the cluster manager | Host memory for pinned pool and spill not in overhead | Raise memoryOverhead to cover pinned pool plus host spill |
| Results differ slightly from CPU | Float aggregation order; incompatible ops are on by default | Use decimals for money, or disable incompatibleOps |
| Executors fail at start | Driver, CUDA and jar mismatch, or discovery script missing | Pin tested versions in the image; validate with a smoke job |
| GPU idle between bursts | Tasks waiting on object-store reads | Multithreaded reader threads, file coalescing, larger partitions |
Trade-offs
GPU clusters cost more per node, are harder to obtain on demand, and add a driver and CUDA stack to your images. Correctness risk is small but not zero, as the incompatible-ops default shows, and operator coverage varies by version, and an upgrade can move an operator onto or off the GPU. Compared with other accelerated engines, such as Databricks Photon on the CPU, a GPU path gives the largest gains on heavy joins, aggregations and sorts, and the smallest on jobs limited by storage. Dynamic allocation works with GPU executors, but scaling slowly in idle periods matters more when each idle executor holds an expensive device; see dynamic allocation for the timeouts.
What to do next
- Run the qualification tool over a month of CPU event logs and pick the top three jobs by estimated savings.
- Remove UDFs from those jobs in favour of built-in functions, and measure the CPU gain on its own.
- Build a pinned image with a tested driver, CUDA and plugin version; set one GPU per executor and a task GPU fraction of 1/cores.
- Set a pinned pool, raise memory overhead to cover it, raise maxPartitionBytes, lower shuffle partitions and enable AQE.
- Run with explain at NOT_ON_GPU, eliminate fallbacks, then run the profiling tool on the GPU event log and act on spill and concurrency advice.
- Compare cost per run, not runtime, and add a regression check that fails the pipeline if a plugin upgrade introduces new fallbacks.