TL;DR: An Apache Spark job is not one thing. It is a stack of nine layers, and almost every layer has
parts you can swap: where executors come from, how the code is written, which engine runs the plan, how joins
and shuffles are decided, where shuffle files live, how memory is divided, how the scheduler reacts to slow and
dying nodes, how bytes sit on disk, and how files are laid out so the planner can skip them. This article maps
all of those options by the problem each one solves, with the--confline or PySpark call for every one of them.โฑ๏ธ 40-minute read. It is the reference companion to
What actually happens when you callspark.read?, which followed one job through the
internals. This one is the parts catalogue. Everything here was checked against Spark 3.5 and Spark 4.0
documentation; where a platform (Databricks, EMR, Dataproc) owns a feature, it says so.๐ฎ Rather break a Spark job than read about one? Every component in this article is a chip you can drag in the
free, open-source Distributed Compute Playground.
Build a job layer by layer and it shows you theexplain()plan, the stages, the task counts and the runtime you
just caused, then hands you 25 on-call tickets to fix. If it clicks, a โญ means the world.
(More about it at the end)
Table of Contents
- How to read this article
- Layer 1: Cluster manager
- Layer 2: Code and APIs
- Layer 3a: Execution engine
- Layer 3b: Planner, joins and AQE
- Layer 4a: Shuffle
- Layer 4b: Memory, caching and serialization
- Layer 5: Scheduling and resilience
- Layer 6a: File and table formats
- Layer 6b: Data layout and the read path
- Combinations that are worth more than their parts
- Six starter profiles
- Cheat sheet: symptom to lever
- Play it: the Distributed Compute Playground ๐ฎ
- Further reading
How to read this article
The previous article told a story: Maya, a new data engineer at RideHub, inherits a report that used to take six minutes and now takes fifty, and we followed one DataFrame through Catalyst, the file index, the scheduler and the shuffle to find out why. If you have not read it, you do not need to; this article stands alone. But the dataset is the same, so when an example says trips, it means 2 TB of Parquet partitioned by dt, 16,000 files, one year of ride-hailing history.
This time the question is different. Not what happens, but what can I change.
The nine layers
Every Spark job, whether it is a notebook cell or a 5 TB nightly join, sits on the same stack:
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ L1 CLUSTER MANAGER local[*] ยท Standalone ยท YARN ยท K8s ยท serverless โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ L2 CODE & APIs DataFrame ยท RDD ยท UDF ladder ยท Connect ยท โ
โ Thrift ยท Structured Streaming โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโฌโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ L3 EXECUTION ENGINE โ L3 PLANNER, JOINS & AQE โ
โ Catalyst+Tungsten (JVM) โ AQE ยท broadcast ยท DPP ยท CBO ยท โ
โ Gluten ยท Comet ยท Photon ยท โ salting ยท bucketing ยท โ
โ RAPIDS โ shuffle.partitions ยท pushdown โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโผโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ L4 SHUFFLE โ L4 MEMORY, CACHING, SERIALIZATION โ
โ ESS ยท tracking ยท Celeborn โ executor shape ยท overhead ยท โ
โ Uniffle ยท Magnet ยท disks โ off-heap ยท cache levels ยท Arrowโ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโดโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ L5 SCHEDULING & RESILIENCE dynamic allocation ยท speculation ยท FAIR ยท โ
โ retries ยท decommission ยท spot ยท history โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโฌโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ L6 FILE / TABLE FORMAT โ L6 DATA LAYOUT & READ PATH โ
โ Parquet ยท ORC ยท Avro ยท โ partitionBy ยท bucketBy ยท โ
โ CSV ยท JSON ยท Delta ยท โ Z-ORDER ยท compaction ยท โ
โ Iceberg ยท Hudi ยท Hive โ schema ยท committers ยท stats โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
Each section below takes one layer, lists every part you can plug into it, and for each part answers the same four questions:
- What problem does it solve?
- What does it cost? (Every lever has a price. Some are hidden.)
- What does it need? (Prerequisites and conflicts. Dynamic allocation without a shuffle service is a config error, not a tuning choice.)
-
How do you turn it on? (The
--confor the PySpark call.)
The eleven dimensions
To compare parts that do completely different things, it helps to score them against the same questions. The playground uses eleven, and this article uses them informally too:
| Dimension | The question it answers |
|---|---|
| Batch throughput | How many bytes per hour can the nightly job chew through? |
| Interactive latency | How fast does a dashboard query come back? |
| Shuffle efficiency | How much data crosses the network, and how painful is that? |
| Memory safety | How close is the driver or an executor to an OutOfMemoryError? |
| Skew resilience | What happens when one key holds 30% of the rows? |
| I/O and small-file efficiency | How much of the time is spent listing, opening and parsing files? |
| Fault tolerance | What happens when a node disappears mid-stage? |
| Cost efficiency | Dollars per terabyte processed. |
| Python friendliness | How much does writing it in Python cost compared with SQL? |
| Streaming fitness | Can it run forever with bounded state? |
| Operational simplicity | How many things can page you at 04:00? |
Most of the dangerous settings in this article look great on one dimension and are catastrophic on another. spark.sql.autoBroadcastJoinThreshold = 2g is a shuffle-efficiency win and a memory-safety disaster. The point of scoring against all eleven is that you notice.
Three kinds of "dangerous"
Throughout, settings are marked three ways:
- โ ๏ธ Warning: a legitimate tool that is wrong more often than it is right. Example:
MEMORY_ONLYpersistence on a dataset bigger than RAM. - ๐งจ Danger: a setting that will eventually take production down, usually on the first busy day after the person who set it has left. Example:
spark.task.maxFailures = 1. - โน๏ธ Note: a prerequisite or platform limit that is easy to miss. Example: Kubernetes has no external shuffle service.
Layer 1: Cluster manager
The cluster manager answers one question: where do executors come from, and what happens to their disks when they go away? That second half decides more than people expect. It is the reason dynamic allocation works out of the box on YARN and needs a workaround on Kubernetes.
driver
โ asks for N executors
โโโโโโโโโโโโโดโโโโโโโโโโโโโ
โผ โผ
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ cluster manager โ โ cluster manager โ
โ (YARN / K8s / โ ... โ decides which โ
โ standalone) โ โ node, how big โ
โโโโโโโโโโฌโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ launches
โโโโโโโโโโผโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ executor (JVM) โ โ executor (JVM) โ โ executor (JVM) โ
โ cores ยท heap ยท โ โ โ โ โ
โ LOCAL DISK โโโโโโผโโโผโโ shuffle files live here, and die with the node
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
local[*]: one JVM
Driver and executors share one process on your laptop. * means one task slot per core.
-
Solves: running tests, reading
explain()output, learning. Everything in this article can be observed onlocal[*]with a few hundred megabytes of Parquet. - Costs: no fault tolerance, one machine of RAM, one machine of disk for shuffle. Scores poorly on throughput and resilience by definition.
- โน๏ธ
local[*]is a learning tool, not a deployment target. If the nightly job "works on my laptop" at 2 TB, it is because the laptop has the data cached and nobody is sharing it.
spark.master=local[*]
Spark Standalone
Spark's own master/worker cluster manager. No Hadoop, no Kubernetes, one process per node and a web UI.
- Solves: a dedicated Spark cluster with the least moving parts. The worker can host the external shuffle service, so dynamic allocation works.
-
Costs: no queues, no multi-tenancy worth the name, no preemption. One application tends to take the whole cluster unless you cap
spark.cores.max.
spark.master=spark://master:7077
Docs: Spark Standalone Mode.
Hadoop YARN
The classic. YARN runs a ResourceManager and a NodeManager per node; the NodeManager can host Spark's external shuffle service as an auxiliary service.
-
Solves: queues, capacity limits, preemption between teams, and shuffle files that survive executor removal. The reason the "Nightly ETL starter" profile at the end of this article uses YARN is that everything simply works:
spark.shuffle.service.enabled=true,spark.dynamicAllocation.enabled=true, done. -
Costs: a Hadoop installation to operate. Executor sizes are quantised by
yarn.scheduler.minimum-allocation-mb. Container start-up is slower than a pod. - Deploy mode matters:
clusterputs the driver inside the YARN application (survives your SSH session dying);clientkeeps it where you ranspark-submit.
spark.master=yarn
spark.submit.deployMode=cluster
Docs: Running Spark on YARN.
Kubernetes
Executors are pods. You bring a container image; Spark asks the API server for pods and tears them down when done.
- Solves: elasticity, one platform for Spark and everything else, custom images with exactly the Python packages a job needs, namespaces and quotas.
-
Costs: โน๏ธ there is no external shuffle service on Kubernetes. When an executor pod is removed, its shuffle files go with it. That means dynamic allocation needs either shuffle tracking (keep executors alive while their shuffle data is needed) or a remote shuffle service (Celeborn, Uniffle), both covered in Layer 4a. Expect to spend time on image builds, service accounts and volume mounts for
spark.local.dir.
spark.master=k8s://https://k8s-api:6443
spark.kubernetes.container.image=registry/spark-py:3.5
spark.kubernetes.namespace=data-eng
spark.kubernetes.authenticate.driver.serviceAccountName=spark
Docs: Running Spark on Kubernetes.
Managed serverless Spark
Dataproc Serverless, EMR Serverless and Databricks serverless compute remove the cluster entirely: you submit a job, the platform sizes and scales it.
- Solves: operational simplicity, start-up time, autoscaling you did not have to configure, and the "who left the 40-node cluster running over the weekend" bill.
-
Costs: โน๏ธ a ceiling on what you may change. Executor shapes are chosen from a menu, some plugins are unavailable, and a few
spark.*settings are silently overridden. Read the platform limits before relying on a knob from this article. - Platform-specific: it is one of the parts the playground marks as not on open-source Spark.
Docs: Dataproc Serverless, EMR Serverless, Databricks serverless compute.
Choosing
| You want | Pick | Then also |
|---|---|---|
| To learn, test, read plans | local[*] |
nothing |
| One team, one cluster, no Hadoop | Standalone | external shuffle service |
| Multi-tenant queues, battle-tested dynamic allocation | YARN | external shuffle service, push-based shuffle |
| Cloud-native, custom images, elasticity | Kubernetes | shuffle tracking or Celeborn / Uniffle |
| Nobody to operate a cluster | Serverless | accept the platform's limits |
Layer 2: Code and APIs
How the job is written decides what the optimizer can see. This is the layer where a one-line change (a @udf decorator) can make a job 50ร slower, and where the fix is usually a different decorator.
DataFrame / SQL API
Declarative plans. You describe what, Catalyst decides how: column pruning, predicate pushdown, join reordering, whole-stage code generation. Python, Scala, Java and SQL all compile to the same logical plan, so a PySpark DataFrame pipeline runs at JVM speed.
df = spark.read.parquet("s3a://ridehub/trips/")
report = (df.filter(F.col("dt") >= "2026-03-01")
.groupBy("zone_id").agg(F.sum("fare")))
This is the default and the baseline. Everything else in this layer is measured against it.
RDD API
The original abstraction: a distributed collection of opaque objects with map, filter, reduceByKey.
-
Solves: full control over partitioning and per-record logic, access to things the DataFrame API hides (custom partitioners,
mapPartitionswith external state). - โ ๏ธ Costs: RDDs skip Catalyst. No column pruning, no pushdown, no AQE, no codegen. In Python, every row is pickled from the JVM to a Python worker process and back. On a 2 TB Parquet scan that reads six columns out of forty, the RDD path reads all forty and deserializes them one object at a time.
- Use it when you genuinely need it, and convert back to a DataFrame as early as possible.
Docs: RDD Programming Guide.
The Python UDF ladder
This is the single most important diagram in this layer. Every rung is "run my Python function on every row"; the difference is how the rows get to Python.
JVM executor Python worker
โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ
row UDF โ row โโบ pickle โโผโโโโ one row โโโโโโโโโบโ unpickle โโบ f()โ ~10โ100ร slower
@udf โ row โโบ pickle โโผโโโโ one row โโโโโโโโโบโ unpickle โโบ f()โ than native
โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ
Arrow UDF โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ
@udf( โ 10,000 rows โโโโผโโโ Arrow batch โโโโโโบโ f() per value โ same code,
useArrow) โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ batched transfer
pandas UDF โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ
@pandas_udf โ 10,000 rows โโโโผโโโ Arrow batch โโโโโโบโ f(pd.Series) โ vectorized:
โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ numpy speed
mapInPandas โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ
โ whole partitionโผโโโ Arrow batches โโโโบโ f(iter[pd.DF]) โ model.predict()
โโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโ on a whole batch
Row-at-a-time Python UDF (@F.udf). The rung everyone starts on.
@F.udf("boolean")
def is_premium(fare):
return fare is not None and fare > 100
df.filter(is_premium("fare")) # BatchEvalPython in the plan
- โ ๏ธ Typically 10โ100ร slower than the equivalent
F.col("fare") > 100. Worse, Catalyst treats it as a black box: nothing above it in the plan can be pushed down into the scan, so a UDF in afiltermeans the scan reads everything. Look forBatchEvalPythoninexplain(); it is the smoking gun.
Arrow-optimized Python UDF (Spark 3.5+, @F.udf(..., useArrow=True)). Same function, same row-by-row semantics, but rows cross the boundary in Arrow batches instead of pickle. A free 2โ4ร for the price of one keyword argument.
@F.udf("double", useArrow=True)
def tip_ratio(tip, fare):
return tip / fare if fare else None
Pandas UDF (@F.pandas_udf). Your function receives a pd.Series (or several) per Arrow batch and returns one. This is where Python logic stops being slow, because pandas and numpy do the inner loop in C.
@F.pandas_udf("double")
def zscore(v: pd.Series) -> pd.Series:
return (v - v.mean()) / v.std()
mapInPandas / applyInPandas. Hand an iterator of pandas DataFrames (a whole partition, or a whole group) to a function. The natural home for model scoring: load the model once per partition, call predict on a batch.
def score(batches):
for pdf in batches:
pdf["p"] = model.predict(pdf[FEATURES])
yield pdf
df.mapInPandas(score, schema="trip_id string, p double")
- All of the Arrow-based rungs still block pushdown for expressions above them. Put filters before the UDF in the pipeline, and make sure the UDF's inputs are the only columns that survive to that point.
- All of them need memory outside the JVM heap for the Python worker and the Arrow buffers. See
memoryOverheadin Layer 4b. The classic failure isContainer killed by YARN for exceeding memory limitswith a heap that is barely used.
Docs: Apache Arrow in PySpark.
pandas API on Spark
import pyspark.pandas as ps gives you pandas syntax on a Spark DataFrame. Good for migrating notebooks.
import pyspark.pandas as ps
psdf = ps.read_parquet("s3a://ridehub/trips/")
psdf.groupby("zone_id")["fare"].sum()
- โน๏ธ pandas has a sequential index; Spark does not. The default index type builds one, and some operations (
iloc,shift) collapse to a single partition to do it. Setps.set_option("compute.default_index_type", "distributed")unless you need pandas-exact ordering.
Docs: pandas API on Spark.
Scala / Java UDF
Runs inside the JVM, so there is no serialization tax. Still opaque to Catalyst: no pushdown through it, no codegen of the expression itself. Cheaper than any Python rung, more expensive than a built-in function.
Built-in functions only
The top of the ladder is not a UDF at all. pyspark.sql.functions has over 400 functions, and Spark 3.5 added F.call_function plus higher-order functions (transform, filter, aggregate) over arrays and maps. Almost every "I need a UDF for this" has a native expression, and the native expression is visible to Catalyst.
df.filter(F.col("fare") > 100).withColumn("tip_ratio", F.col("tip") / F.col("fare"))
Docs: Functions.
Spark Connect
Spark 3.4+ splits the client from the driver. Your notebook holds a thin gRPC client; the driver (and the JVM, and the dependency hell) runs on a server. Clients can be upgraded independently, one bad client cannot crash the driver, and pip install pyspark-connect is all a laptop needs.
spark.remote=sc://spark-connect:15002
- โน๏ธ RDD APIs and
SparkContextcalls are not available over Connect. If your code usessc.parallelize, it is not a Connect job yet.
Docs: Spark Connect overview.
Thrift server (JDBC / BI)
A long-running SparkSession that speaks HiveServer2's protocol, so Tableau, Looker or beeline can run SQL against your lake.
- Solves: dashboards. The session stays warm, caches survive between queries, and the file index for a table is built once.
- Costs: one application shared by everyone. Without FAIR scheduler pools (Layer 5), one analyst's 30-minute query starves every dashboard.
Docs: Distributed SQL Engine.
Structured Streaming and its three companions
The DataFrame API over an unbounded source. Micro-batch by default (continuous mode exists but is rarely the right choice), checkpointed, exactly-once with the right sink.
events = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "rides").load())
Three parts turn a streaming demo into a streaming pipeline. They only make sense together with Structured Streaming, which is why the playground refuses to activate them without it.
RocksDB state store. Streaming aggregations, stream-stream joins and dropDuplicates keep state. The default state store keeps it in the JVM heap, which works until the state is bigger than the heap and the executor dies in a GC storm on Friday evening. RocksDB keeps state off-heap on local disk with changelog checkpointing.
spark.sql.streaming.stateStore.providerClass=org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider
Watermark and trigger tuning. Without a watermark, a windowed aggregation keeps every window forever, because a late event could still arrive. withWatermark("event_ts", "10 minutes") tells Spark when a window can be finalised and its state dropped. The trigger interval decides the latency/throughput trade; pick it on purpose.
(events.withWatermark("event_ts", "10 minutes")
.groupBy(F.window("event_ts", "5 minutes"), "zone_id").count()
.writeStream.trigger(processingTime="1 minute") ...)
foreachBatch with an idempotent sink. File sinks and Delta/Iceberg are exactly-once by design. A JDBC table or a MERGE is not, unless you make each micro-batch idempotent, keyed by batchId, so a replay after a crash does not double-count.
def upsert(batch_df, batch_id):
batch_df.createOrReplaceTempView("b")
spark.sql("""MERGE INTO rides t USING b ON t.id = b.id
WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *""")
agg.writeStream.foreachBatch(upsert).option("checkpointLocation", chk).start()
Docs: Structured Streaming Programming Guide, RocksDB state store.
Choosing
| Situation | Use |
|---|---|
| Anything expressible in SQL | DataFrame API with built-in functions. Always first. |
| Python logic over a column | pandas UDF; Arrow UDF if the logic is inherently scalar |
| Model scoring |
mapInPandas, model loaded once per partition |
| Migrating a pandas notebook | pandas API on Spark, with the distributed index |
| Custom partitioning or per-partition state in the JVM | RDD, then back to DataFrame |
| Many notebooks, one cluster | Spark Connect |
| BI tools | Thrift server + FAIR pools + cache |
| Kafka in, table out, forever | Structured Streaming + RocksDB + watermark + idempotent sink |
Layer 3a: Execution engine
The engine is the one part of the stack you cannot leave empty. It is what turns a physical plan into CPU instructions. For fifteen years that meant one thing: Spark's own JVM engine. Since about 2022 there have been real alternatives, and they all work the same way: a Spark plugin intercepts the physical plan and replaces operators it knows how to run with native (C++, Rust or GPU) versions.
Catalyst physical plan
โ
โผ
โโโโโโโโโโโโโโโโโโโโโโโ supported operator?
โ engine plugin โโโโโโโโโโโโโโโโฌโโโโโโโโโโโโโโโ
โ (Gluten / Comet / โ yes no
โ Photon / RAPIDS) โ โ โ
โโโโโโโโโโโโโโโโโโโโโโโ โผ โผ
native columnar JVM row engine
operator (Tungsten codegen)
โ โฒ
โโโ columnarโrow conversion โโโ
(the hidden cost of every fallback)
Catalyst + Tungsten (JVM): stock Spark
Whole-stage code generation (operators fused into one Java method per stage), a vectorized Parquet/ORC reader, the unified memory manager, and UnsafeRow binary rows that skip Java object overhead.
- Solves: everything, everywhere, with every UDF and every data source. The safest default and the baseline every native engine benchmarks against.
- Costs: CPU-bound SQL on wide tables runs 2โ4ร slower than the native engines. GC pauses on large heaps. That is the whole pitch of the alternatives.
Gluten + Velox
Apache Gluten (incubating) plugs Meta's Velox C++ library under the Spark API. It also ships a columnar shuffle manager so data stays in Arrow-like columnar batches across the exchange.
- Solves: CPU-bound SQL and DataFrame ETL. Typical published numbers are 2โ3ร on TPC-H/DS-style workloads; the playground models it as about 2ร.
-
Costs: โน๏ธ operators and expressions Velox does not support fall back to the JVM, with a columnar-to-row conversion on each side of the fallback. A Python UDF in the middle of a Gluten plan is slower than the same UDF in a plain JVM plan. It also needs off-heap memory configured, because Velox allocates outside the heap; forgetting
spark.memory.offHeap.sizeis the most common first-day failure.
spark.plugins=org.apache.gluten.GlutenPlugin
spark.shuffle.manager=org.apache.spark.shuffle.sort.ColumnarShuffleManager
spark.memory.offHeap.enabled=true
spark.memory.offHeap.size=16g
Apache Comet (DataFusion)
Apache DataFusion Comet is the Rust-based equivalent from the Arrow project: a native Parquet scan, native operators and a native columnar shuffle, for Spark 3.4 through 4.0.
- Solves: the same problem as Gluten with a smaller runtime, Rust memory safety and tight Arrow integration.
-
Costs: โน๏ธ operator coverage is still growing. Check the compatibility guide for your plan, and run
explain(): Comet-executed operators appear with aCometprefix, so you can see exactly where it fell back.
spark.plugins=org.apache.spark.CometPlugin
spark.comet.exec.shuffle.enabled=true
spark.memory.offHeap.enabled=true
spark.memory.offHeap.size=8g
Photon
Databricks' proprietary C++ engine. On Databricks it is a checkbox; off Databricks it does not exist.
- Solves: the same CPU-bound SQL problem, with the best integration of the four (Delta, SQL warehouses, streaming support) and no configuration.
- Costs: platform lock-in and a premium DBU rate. โน๏ธ UDFs and RDDs still fall back to the JVM, so a Photon cluster running a row-at-a-time Python UDF is paying the premium for nothing. Photon on YARN is not a configuration; the playground calls this a platform conflict and tells you the job runs nowhere.
Docs: Photon.
RAPIDS Accelerator (GPU)
NVIDIA's RAPIDS Accelerator for Apache Spark runs scans, joins, aggregates and sorts on GPUs, optionally with a UCX-based GPU-to-GPU shuffle.
- Solves: very wide, very numeric ETL and joins where the GPU's memory bandwidth dwarfs the CPU's. On the right workload it is the fastest thing here by a wide margin.
- โ ๏ธ Costs: GPU memory is small (tens of GB) and unforgiving; wide strings, UDFs and unsupported expressions spill or fall back, and the fallbacks are expensive because the data has to come back off the GPU. GPU nodes cost several times more per hour, so the speed-up has to be large to pay for itself. Benchmark your plan before committing a budget. Pair it with stage-level scheduling (Layer 5) so only the stage that benefits runs on GPU executors.
spark.plugins=com.nvidia.spark.SQLPlugin
spark.rapids.sql.enabled=true
spark.executor.resource.gpu.amount=1
spark.task.resource.gpu.amount=0.25
Choosing an engine
| Workload | First choice | Why |
|---|---|---|
| Anything with Python UDFs or RDDs | JVM | native engines fall back anyway and pay a conversion each way |
| CPU-bound SQL/DataFrame ETL on open-source Spark | Gluten or Comet, with off-heap sized | 2โ3ร for the cost of a plugin and some memory config |
| Same, on Databricks | Photon | zero configuration, best Delta integration |
| Wide numeric ETL, GPU nodes already available | RAPIDS, scoped with stage-level scheduling | memory bandwidth |
| Streaming | JVM (or Photon on Databricks) | the open-source native engines have thin streaming support |
A rule that holds across all of them: run explain() and count the fallbacks. A native engine that executes 95% of the plan is a win; one that executes 40% of it, with conversions at every boundary, is slower than stock Spark.
Layer 3b: Planner, joins and AQE
This is the layer the previous article spent most of its time in, so here is the compressed version with every knob attached. Catalyst produces a physical plan from statistics it can see at planning time; Adaptive Query Execution re-plans at every shuffle boundary using the sizes it actually observed. Most of the parts below either feed the planner better information or tell AQE what it is allowed to change.
Adaptive Query Execution
On by default since Spark 3.2. At each shuffle boundary, AQE looks at the real map output sizes and may coalesce partitions, switch join strategies or split skewed partitions. In explain() you see AdaptiveSparkPlan isFinalPlan=false before the first action and isFinalPlan=true after.
spark.sql.adaptive.enabled=true
Three sub-features sit under it. They do nothing without it.
Coalesce shuffle partitions. The default spark.sql.shuffle.partitions is 200, whatever the data size. AQE merges small post-shuffle partitions up to an advisory size, so a 1 GB aggregate ends up in 16 reducers instead of 200 tiny ones, and the write produces 16 files instead of 200.
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.advisoryPartitionSizeInBytes=64m
Skew join handling. If one post-shuffle partition is more than five times the median and more than 256 MB, AQE splits it into several reader tasks on the big side and duplicates the matching partition on the small side. The three-hour straggler becomes twelve two-minute tasks. This is the single most valuable automatic fix in modern Spark.
spark.sql.adaptive.skewJoin.enabled=true
spark.sql.adaptive.skewJoin.skewedPartitionFactor=5
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256m
Runtime broadcast conversion. The planner chose a sort-merge join because the dimension's file size looked big; after the shuffle, the dimension's actual bytes are 4 MB. AQE converts to a broadcast hash join on the fly, and the local shuffle reader lets the big side read its own map outputs instead of a real exchange.
spark.sql.adaptive.autoBroadcastJoinThreshold=10m
spark.sql.adaptive.localShuffleReader.enabled=true
Docs: Adaptive Query Execution.
How Spark picks a join
Before the knobs, the decision tree, because every knob below moves you to a different leaf:
is one side โค autoBroadcastJoinThreshold (10 MB)
or hinted BROADCAST?
โ
yes โโโโโโโโโดโโโโโโโโ no
โ โ
BroadcastHashJoin equi-join and preferSortMergeJoin=false
ยท build side collected and one side small enough per partition?
on the DRIVER, sent โ
to every executor yes โโโโโโดโโโโโ no
ยท zero Exchange for โ โ
the big side ShuffledHashJoin SortMergeJoin
ยท shuffle both ยท shuffle both, sort both
ยท hash the small ยท spills gracefully
side per task ยท the default for big-to-big
โ
both sides bucketed on the key,
same bucket count?
โ
yes โโโโโโดโโโโโ no
โ โ
SortMergeJoin SortMergeJoin
with NO Exchange with 2 Exchanges
Broadcast hint on small tables
Tell the planner the dimension is small instead of letting it guess from directory size. The directory size estimate is what fooled the planner in the opening story: a 50 MB table that compressed to 12 MB on disk still looked bigger than the 10 MB threshold.
trips.join(F.broadcast(zones), "zone_id")
# or in SQL: SELECT /*+ BROADCAST(z) */ ...
- Solves: a sort-merge join with two exchanges becomes a broadcast join with zero exchanges on the fact side.
- Costs: the build side is collected on the driver and sent to every executor. Hint a table that turns out to be 3 GB and you have found the fastest way to kill a driver.
autoBroadcastJoinThreshold = 100 MB
Raise the automatic threshold from 10 MB to 100 MB, so moderately sized dimensions broadcast without hints.
- Solves: a star schema with a handful of 20โ80 MB dimensions, none of which anyone remembered to hint.
- โน๏ธ Costs: 100 MB of Parquet can be a gigabyte of Java objects once decompressed and hashed on the driver. Raise
spark.driver.memoryat the same time, every time.
spark.sql.autoBroadcastJoinThreshold=100m
spark.driver.memory=8g
๐งจ autoBroadcastJoinThreshold = 2 GB
"The executors have plenty of RAM." They do; the driver does not, and the driver is where every broadcast relation is built first. Broadcast builds are capped at 8 GB and 512 million rows, and spark.sql.broadcastTimeout defaults to 300 seconds. A 2 GB threshold works on the quiet day the setting was tested and dies on the first busy one with Could not execute broadcast in 300 secs or a driver OOM, taking every job in the application with it.
spark.sql.autoBroadcastJoinThreshold=2g # don't
Dynamic partition pruning
On by default since 3.0. In a star-schema join where the dimension is filtered (calendar.quarter = 'Q1'), DPP runs the dimension side first, collects the matching keys, and uses them to prune the fact table's partitions before scanning. In explain() the fact scan shows PartitionFilters: [dynamicpruningexpression(...)].
- Solves: the "the filter is on the calendar table, so the 5 TB fact scan reads two years" problem. Requires the fact table to be partitioned by the join key.
spark.sql.optimizer.dynamicPartitionPruning.enabled=true
Cost-based optimizer and statistics
Catalyst's default size estimates come from file sizes. ANALYZE TABLE stores row counts, distinct counts and min/max per column in the catalog; with CBO on, the planner reorders multi-way joins and estimates join output sizes properly.
- Solves: wrong join order in three-way-plus joins, and dimensions that broadcast only after a filter because the planner now knows how selective the filter is.
-
Costs: statistics go stale, and
ANALYZEon a 2 TB table is itself a job. Table formats (Delta, Iceberg) carry per-file statistics that remove most of the need.
spark.sql.cbo.enabled=true
spark.sql.cbo.joinReorder.enabled=true
spark.sql("ANALYZE TABLE trips COMPUTE STATISTICS FOR ALL COLUMNS")
Prefer shuffled hash join
By default Spark prefers sort-merge over shuffled hash because SMJ spills gracefully. Turning the preference off skips the sort when one side is much smaller.
- โน๏ธ The build side of a shuffled hash join must fit in memory per task. Keep it for clearly asymmetric joins where you have measured the gain.
spark.sql.join.preferSortMergeJoin=false
Key salting for hot keys
The manual skew fix, from before AQE existed and still needed when the skew is in an aggregation rather than a join, or when the hot key is a NULL. Append a random suffix to the hot key on the big side, explode the small side into every suffix, join on both.
big = big.withColumn("salt", (F.rand() * 16).cast("int"))
small = small.withColumn("salt", F.explode(F.array([F.lit(i) for i in range(16)])))
big.join(small, ["key", "salt"])
- Costs: code complexity, the small side grows 16ร, and the next hot key will not announce itself. Prefer AQE skew handling when the shape allows it.
Bucketed joins
Both tables written with bucketBy(N, key) into a metastore table. Rows with the same key hash land in the same bucket number on both sides, so the join needs no Exchange at all. With Iceberg, the same idea is called a storage-partitioned join (Spark 3.3+).
df.write.bucketBy(256, "driver_id").sortBy("driver_id").saveAsTable("trips_b")
- Solves: the recurring 5 TB-to-5 TB join on the same key, every night. The shuffle is paid once at write time instead of every read.
- Costs: both sides must agree on bucket count and column, the write is a shuffle, and plain Parquet directories cannot carry bucketing metadata (hence the Hive metastore or Iceberg prerequisite).
Tune spark.sql.shuffle.partitions
The default 200 reducers is right for roughly 25 GB of shuffle at 128 MB each. For 1.2 TB of shuffle it means 6 GB per reducer and spills everywhere; for 200 MB it means 200 one-megabyte files. Size it to the data, or let AQE coalescing handle the small end.
spark.sql.shuffle.partitions=1200 # โ shuffle bytes / 128 MB
๐งจ spark.sql.shuffle.partitions = 1
"I only want one output file." Now every aggregation and every join in the application runs on a single task, on a single core, with a single task's memory. Use coalesce(1) on the final write if you must have one file, and leave the shuffle alone.
Pushdown-friendly predicates
Not a config: a habit. The planner can only prune partitions and push filters into Parquet when the predicate is a plain comparison on the partition column or a stats-bearing column.
# โ prunes partitions: PartitionFilters: [dt >= 2026-03-01, dt <= 2026-03-31]
trips.filter(F.col("dt").between("2026-03-01", "2026-03-31"))
# โ scans everything: the filter is on a derived expression over a different column
trips.filter(F.month("pickup_ts") == 3)
Check explain() for PartitionFilters: [] and PushedFilters: []. Empty brackets on a partitioned table are the most expensive two characters in Spark.
Choosing
| Symptom | Lever |
|---|---|
| 200 tiny output files | AQE coalesce, or coalesce() before the write |
| One task runs for hours, the rest finish in minutes | AQE skew join; salting if it is an aggregation |
SortMergeJoin against a table you know is small |
broadcast hint; or CBO stats so the planner knows too |
| Star join scans the whole fact table | DPP + partition the fact by the join key |
| Same big join every night | bucket both tables |
PartitionFilters: [] |
rewrite the predicate on the partition column |
| Driver OOM during a join | your broadcast threshold is too high |
Layer 4a: Shuffle
A shuffle is the moment Spark's "every task is independent" model breaks: to group by key, every reducer needs rows from every mapper. Where the intermediate files live decides whether losing an executor costs you a task or a stage, and whether you can release executors at all.
mappers write reducers read
โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ
โ executor A โ shuffle files โ executor C โ
โ task 1 โโโโโโโผโโโบ [p0 p1 p2 โฆ ] โ reducer p0 โโโผโโ fetches p0 from A and B
โโโโโโโโโโโโโโโโ โ โโโโโโโโโโโโโโโโ
โโโโโโโโโโโโโโโโ โ WHERE do these live?
โ executor B โ โ (a) executor's local disk โ die with the executor
โ task 2 โโโโโโโผโโโบ [p0 p1 p2 โฆ ] (b) node's shuffle service โ survive the executor
โโโโโโโโโโโโโโโโ (c) remote shuffle cluster โ survive the NODE
External shuffle service (YARN, Standalone)
A long-running process per node that serves shuffle files on behalf of executors that have exited.
- Solves: executors can be removed mid-job (dynamic allocation) without losing the map outputs they produced. Also survives an executor crash: the stage does not need to be recomputed.
- Needs: YARN (as a NodeManager auxiliary service) or Standalone (in the worker). โน๏ธ Kubernetes has no equivalent.
- Does not solve: losing the node. The files are still on that node's disk.
spark.shuffle.service.enabled=true
Shuffle tracking (Kubernetes)
The Kubernetes workaround for dynamic allocation without a shuffle service: Spark tracks which executors hold shuffle files that are still needed, and refuses to release those.
- Solves: scale-up on Kubernetes, and scale-down for executors whose shuffle data has been consumed.
- โน๏ธ Costs: scale-down is slow and partial. An executor that produced map output for a long stage stays alive until that stage's reducers have all read it. With a 120 s timeout you still pay for idle pods.
spark.dynamicAllocation.shuffleTracking.enabled=true
spark.dynamicAllocation.shuffleTracking.timeout=120s
Apache Celeborn (remote shuffle)
Apache Celeborn runs a separate cluster of shuffle servers. Mappers push their output to Celeborn workers, which merge it per reducer; executors keep nothing on local disk.
- Solves: three problems at once. Executors become stateless, so dynamic allocation can release any of them instantly and spot nodes can vanish without triggering recomputation. Reducers do a few large sequential reads instead of thousands of small random ones. And Kubernetes gets a proper shuffle service.
- Costs: another cluster to run, size and monitor, and an extra network hop for every shuffle byte. It pays for itself from the first spot interruption on.
spark.shuffle.manager=org.apache.spark.shuffle.celeborn.SparkShuffleManager
spark.celeborn.master.endpoints=celeborn-master:9097
Apache Uniffle (remote shuffle)
Apache Uniffle, from Tencent, is the other Apache remote shuffle service with the same architecture (coordinator + shuffle servers) and the same benefits. Pick one; they do not coexist in one application.
spark.shuffle.manager=org.apache.spark.shuffle.RssShuffleManager
spark.rss.coordinator.quorum=rss-coord:19999
Push-based shuffle (Magnet)
LinkedIn's contribution to Spark 3.2: mappers push blocks to the external shuffle service, which merges them per reducer ahead of time. Reducers read one merged file instead of N small ones.
- Solves: the small-random-read problem on very large shuffles, without a separate shuffle cluster.
- Needs: YARN with the external shuffle service. Not available on Kubernetes or Databricks.
spark.shuffle.push.enabled=true
spark.shuffle.push.server.mergedShuffleFileManagerImpl=org.apache.spark.network.shuffle.RemoteBlockPushResolver
zstd shuffle compression
The shuffle codec defaults to lz4. zstd produces 20โ30% smaller shuffle files and spills for slightly more CPU. On a network- or disk-bound shuffle that is a net win.
spark.io.compression.codec=zstd
โ ๏ธ Disabling shuffle compression
"Compression costs CPU." It does: about 1% of the task time with lz4. Turning it off costs 3โ5ร the network and disk traffic. Shuffle stages are almost never CPU-bound; leave it on.
spark.shuffle.compress=false # don't
spark.shuffle.spill.compress=false # really don't
Fast local disks
spark.local.dir is where shuffle files and spills go. On cloud VMs the default is often the boot disk: a network-attached volume with a few thousand IOPS shared by everything. Point it at local NVMe and list several disks so Spark stripes across them.
spark.local.dir=/mnt/nvme0,/mnt/nvme1
- โน๏ธ On Kubernetes and serverless platforms this is a pod spec or platform setting, not a Spark conf.
Dataproc Enhanced Flexibility Mode
Dataproc-specific: shuffle data for jobs on preemptible secondary workers is written to the primary (on-demand) workers, so a preempted secondary loses nothing. The same idea as a remote shuffle service, built into the platform.
Docs: Dataproc Enhanced Flexibility Mode.
Bigger fetch buffers
spark.reducer.maxSizeInFlight (default 48 MB) bounds how much a reducer requests at once. On high-latency networks, doubling it reduces round trips. Pair with more fetch retries for flaky networks.
spark.reducer.maxSizeInFlight=96m
spark.shuffle.io.maxRetries=10
Choosing
| Cluster | For dynamic allocation | For spot nodes |
|---|---|---|
| YARN | external shuffle service (+ Magnet for huge shuffles) | Celeborn / Uniffle, or EFM on Dataproc |
| Standalone | external shuffle service | Celeborn / Uniffle |
| Kubernetes | shuffle tracking (slow scale-down) or Celeborn / Uniffle | Celeborn / Uniffle + decommissioning |
| Serverless | built in | built in |
And everywhere: zstd on, compression never off, local disks that are actually local.
Layer 4b: Memory, caching and serialization
Executor memory is not one number. It is a heap split by the unified memory manager, plus memory outside the heap that Python workers, Arrow buffers and native engines use, plus the container limit that YARN or Kubernetes enforces. Most memory incidents are a mismatch between those three.
container limit (what YARN / K8s kills you for exceeding)
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ spark.executor.memory (JVM heap) โ memoryOverhead โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโฌโโโโโโโโโโโโ โ (off-JVM: โ
โ โ unified region (60% of heap-300MB)โ user + โ โ Python workers, โ
โ โ โโโโโโโโโโโโโโโโโโฌโโโโโโโโโโโโโโ โ reserved โ โ Arrow buffers, โ
โ โ โ execution โ storage โ โ โ โ native engine, โ
โ โ โ (shuffle, sort,โ (cache, โ โ โ โ thread stacks) โ
โ โ โ join, agg) โ broadcast) โ โ โ โ โ
โ โ โ โโโ borrows either way โโโบ โ โ โ โ + offHeap.size โ
โ โ โโโโโโโโโโโโโโโโโโดโโโโโโโโโโโโโโ โ โ โ if enabled โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโดโโโโโโโโโโโโ โ โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
Right-sized executors (4โ5 cores)
The hardware provisioning guide's advice has not changed in a decade: 4โ5 cores and 16โ32 GB per executor. Enough parallelism per JVM to amortise overhead, small enough that GC pauses stay short and a lost executor loses little work. Above 5 cores per executor, HDFS and S3 client throughput stops scaling anyway.
spark.executor.cores=4
spark.executor.memory=20g
spark.executor.instances=10
โ ๏ธ Fat executors (32 cores, 200 GB)
One executor per node sounds efficient. In practice: multi-second GC pauses on a 200 GB heap, 32 tasks lost when one executor dies, and object-store clients saturating at a fraction of the cores. Spark Connect and native engines soften this a little; they do not remove it.
Memory overhead for Python and native code
Python workers, Arrow buffers and Velox/DataFusion/GPU allocations all live outside the JVM heap. memoryOverhead (default 10% of executor memory, minimum 384 MB) is the only room they have. The classic failure is Container killed by YARN for exceeding memory limits. 20.5 GB of 20 GB physical memory used while the heap graph shows 40% free.
spark.executor.memoryOverhead=4g
spark.executor.pyspark.memory=4g # caps each Python worker; counts toward the container
Off-heap memory
Lets Tungsten allocate execution memory with sun.misc.Unsafe outside the heap, so sort and hash buffers do not churn the garbage collector. Required by Gluten and Comet. On a plain JVM job it is a modest GC win on large executors.
spark.memory.offHeap.enabled=true
spark.memory.offHeap.size=8g # counts toward the container limit too
cache() a reused DataFrame (MEMORY_AND_DISK)
If an intermediate is read more than once, materialise it. cache() uses MEMORY_AND_DISK: blocks that do not fit in storage memory spill to local disk instead of being recomputed.
march = trips.filter(...).cache()
march.count() # materialises
march.groupBy("zone_id").count().show() # reads the cache
- Costs: storage memory competes with execution memory; a big cache can push shuffles into spilling. Cache what is reused, not what is big. And remember that cached data is per-application: a nightly batch job caching its only pass gains nothing.
โ ๏ธ persist(MEMORY_ONLY) on a huge dataset
2 TB into 160 GB of RAM. MEMORY_ONLY never spills; when storage memory is short it evicts whole partitions and recomputes them from lineage the next time they are read. The job gets slower than with no cache at all, and the Storage tab shows "27% cached" forever. Use MEMORY_AND_DISK, or do not cache.
MEMORY_ONLY_SER and Kryo
Serialized cache blocks are 2โ5ร smaller than deserialized Java objects, at the cost of CPU on every read. With Kryo registered for your classes, smaller still. This matters for RDDs and Datasets of case classes; DataFrames are already stored as compact UnsafeRows, so the gain there is small.
spark.serializer=org.apache.spark.serializer.KryoSerializer
spark.kryo.registrationRequired=false
Arrow for toPandas() and createDataFrame()
Without Arrow, toPandas() pickles rows one at a time through the driver. With it, the JVM ships columnar batches that pandas ingests directly. Minutes become seconds. There is no reason to leave this off in Spark 3.x.
spark.sql.execution.arrow.pyspark.enabled=true
spark.sql.execution.arrow.maxRecordsPerBatch=10000
spark.driver.maxResultSize sized on purpose
A guard rail. The default 1 GB limit makes a careless collect() fail fast with a clear error instead of filling the driver heap and taking the application down. Raise it only as far as the driver can actually hold.
spark.driver.maxResultSize=2g
๐งจ collect() / toPandas() the whole table
"I just want it in pandas." Every row lands in driver memory, after being serialized through it. trips.toPandas() on 3.2 billion rows is a driver OOM with extra steps. Aggregate first, sample first, or write to storage and read the output with pandas.
pdf = trips.toPandas() # 3.2 billion rows โ one JVM. Don't.
pdf = trips.filter(...).groupBy(...).agg(...).toPandas() # a few thousand rows. Fine.
Bigger driver (16 GB)
The driver holds the file index (16,000 FileStatus objects for trips; millions for a small-files disaster), every broadcast build, every collect() result and the plans themselves. 1 GB is the default; 8โ16 GB is normal for a production job touching large tables.
spark.driver.memory=16g
spark.driver.cores=4
checkpoint() long lineages
Iterative jobs (graph algorithms, 200-round feature pipelines) build lineages hundreds of stages deep. If a partition is lost, Spark replays the lineage to rebuild it; two hundred stages later you wish it had not. checkpoint() writes the DataFrame to reliable storage and truncates the lineage there. Structured Streaming checkpoints for the same reason, automatically.
spark.sparkContext.setCheckpointDir("s3a://ridehub/chk/")
df = df.checkpoint()
G1GC with region tuning
G1 is the default collector on JDK 11+ and is right for almost everyone. Tune it only when the Spark UI's GC Time column exceeds roughly 10% of task time; the first knob is starting concurrent collection earlier.
spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:InitiatingHeapOccupancyPercent=35
Choosing
| Symptom | Lever |
|---|---|
Container killed by YARN for exceeding memory limits, heap mostly free |
memoryOverhead (Python / native / Arrow) |
| Long GC pauses, executors "lost" with no error | smaller executors, off-heap, G1 tuning |
| Driver OOM |
maxResultSize guard, no whole-table collect(), bigger driver, lower broadcast threshold |
| Cache shows partial, job slower than uncached |
MEMORY_ONLY โ MEMORY_AND_DISK, or stop caching |
toPandas() takes minutes on small data |
Arrow |
| Iterative job's recovery replays everything | checkpoint() |
Docs: Tuning Spark: memory management, Hardware provisioning, RDD persistence.
Layer 5: Scheduling and resilience
Layers 1โ4 decide how fast the job runs when nothing goes wrong. This layer decides what happens at 03:40 when a spot node gets its two-minute notice, when one task is ten times slower than its siblings, and when S3 returns its fifth 503 of the night.
Dynamic allocation
Add executors when tasks queue, release them when idle. The difference between paying for 50 executors all night and paying for 50 during the join and 4 during the write.
- Needs somewhere safe for shuffle files, or releasing an executor would destroy map output that a later stage needs. On YARN or Standalone: the external shuffle service. On Kubernetes: shuffle tracking or a remote shuffle service. On serverless: built in. โน๏ธ Spark refuses to start with dynamic allocation enabled and none of these present; the playground models it as an inactive component with an error.
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.minExecutors=2
spark.dynamicAllocation.maxExecutors=50
spark.dynamicAllocation.executorIdleTimeout=60s
Docs: Dynamic resource allocation.
Speculative execution
When 90% of a stage's tasks are done and a straggler is running 3ร longer than the median, launch a copy on another executor and take whichever finishes first.
- Solves: slow nodes (a bad disk, a noisy neighbour). It does not solve skew: a copy of the task that holds 30% of the keys is just as slow.
- โน๏ธ Costs: duplicated side effects. Two copies of a task writing to a non-idempotent sink (a JDBC insert, a
foreachthat calls an API) means duplicate rows. Safe with file commit protocols and table formats; dangerous with everything else.
spark.speculation=true
spark.speculation.multiplier=3
spark.speculation.quantile=0.9
FAIR scheduler pools
Inside one application, jobs run FIFO by default: the first notebook cell to submit takes every slot until it finishes. FAIR mode shares slots between pools by weight, so a dashboard query in the dashboards pool gets its 20% even while a 30-minute backfill runs in batch.
spark.scheduler.mode=FAIR
spark.scheduler.allocation.file=fairscheduler.xml
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "dashboards")
- This is scheduling within an application. Between applications, that is the cluster manager's queue system (YARN queues, Kubernetes namespaces).
Docs: Scheduling within an application.
spark.locality.wait = 0 on object storage
Spark waits up to 3 seconds per locality level (process, node, rack) for a slot near the data before scheduling a task elsewhere. On HDFS that wait finds a node-local replica and saves a network read. On S3 or GCS there are no local replicas; every byte is a network read; the wait is pure delay. Three seconds times a few locality levels times a few hundred scheduling rounds is a noticeable chunk of a short job.
spark.locality.wait=0s # object storage only; keep the default on HDFS
Stage-level scheduling
Spark 3.1+ lets a job request a different executor profile for one stage. The feature pipeline runs on cheap 4-core executors until the model-scoring stage, which runs on 8-core, 64 GB executors with a GPU, and then the job goes back to cheap ones.
from pyspark.resource import ResourceProfileBuilder, ExecutorResourceRequests
rp = (ResourceProfileBuilder()
.require(ExecutorResourceRequests().cores(8).memory("64g").resource("gpu", 1))
.build())
scored = features.rdd.withResources(rp).mapPartitions(score)
- Needs: YARN or Kubernetes with dynamic allocation (the new executors have to come from somewhere). Scoped to the RDD API until recent versions; check your Spark version for DataFrame support.
Docs: Stage level scheduling overview.
Generous task retries
spark.task.maxFailures (default 4) is how many times one task may fail before the stage, and the job, fails. On cloud storage and spot nodes, transient failures are normal: an S3 503 SlowDown, a preempted node mid-task, a shuffle fetch from an executor that just died. Eight attempts with the default back-off costs seconds when it matters and nothing when it does not.
spark.task.maxFailures=8
spark.stage.maxConsecutiveAttempts=8
๐งจ spark.task.maxFailures = 1
"Fail fast." One throttled S3 request, three hours into the job, and the job is gone. This is not fail-fast; it is fail-often. If the goal is to stop retrying a genuinely broken task, the default 4 already does that in under a minute.
Graceful decommissioning
Spark 3.1+. When a node is about to be taken away (spot interruption notice, autoscaler scale-in, Kubernetes drain), Spark can stop scheduling new tasks on it and migrate its shuffle blocks and cached RDD blocks to other executors before it dies. Without this, a preempted node means recomputing every stage that produced data on it.
spark.decommission.enabled=true
spark.storage.decommission.enabled=true
spark.storage.decommission.shuffleBlocks.enabled=true
spark.storage.decommission.rddBlocks.enabled=true
- Works best when there is time (the 2-minute spot notice is enough for most shuffle outputs) and somewhere to put the blocks. With a remote shuffle service there is nothing to migrate, which is even better.
Docs: Decommissioning.
Executor exclusion on failures
Formerly "blacklisting". A node with a failing disk makes every task that lands on it fail; after a threshold, Spark stops scheduling there for a while. Cheap insurance.
spark.excludeOnFailure.enabled=true
Spot / preemptible workers
60โ90% cheaper compute that can be reclaimed with two minutes' notice (AWS, GCP) or less (Azure). The single biggest cost lever for batch, and the single biggest source of 04:00 pages when used naked.
- โ ๏ธ Costs: losing a node loses its shuffle files and cached blocks. Every stage that produced data there must be recomputed. Three interruptions in a 2-hour job can double its runtime, and the job may never finish on a bad day.
- Make it safe with: a remote shuffle service (nothing on the node to lose) or decommissioning (move it before the node dies) or EFM on Dataproc, plus generous retries, plus the driver on an on-demand node. The playground gives the combination spot + Celeborn + dynamic allocation a synergy bonus for exactly this reason.
Event logs and the History Server
The Spark UI disappears when the application ends. Event logging writes every event to storage, and the History Server replays them into the same UI, days later. It is the only way to debug yesterday's failure, and the only way to compare tonight's plan with last week's.
spark.eventLog.enabled=true
spark.eventLog.dir=s3a://ridehub/spark-events/
spark.history.fs.logDirectory=s3a://ridehub/spark-events/
Docs: Monitoring and instrumentation.
Prometheus metrics sink
Executor memory, GC time, task counts and shuffle bytes as time series, with alerts. Spark 3.0+ exposes a Prometheus servlet on the driver and executors natively.
spark.ui.prometheus.enabled=true
spark.metrics.conf.*.sink.prometheusServlet.class=org.apache.spark.metrics.sink.PrometheusServlet
Choosing
| Goal | Levers |
|---|---|
| Pay only for what the stage needs | dynamic allocation + a shuffle home |
| Survive spot interruptions | remote shuffle or decommissioning, retries โฅ 8, on-demand driver |
| One slow node should not define the job | speculation (idempotent sinks only), executor exclusion |
| Dashboards and backfills on one session | FAIR pools |
| GPU for one stage only | stage-level scheduling |
| Short jobs on S3 feel sluggish | locality.wait = 0 |
| Debug last night | event logs + History Server; Prometheus for trends |
Layer 6a: File and table formats
Everything above this line decides how fast Spark can process a byte. This layer decides how many bytes there are to process, and whether Spark can find out which files it does not need to open before opening them.
can skip columns? can skip row groups? one file = many tasks? file list from
Parquet / ORC yes yes yes directory listing
Avro no no yes directory listing
CSV / JSON no no yes directory listing
gzip CSV / JSON no no NO directory listing
Delta / Iceberg / yes yes + per-FILE stats yes transaction log /
Hudi (on Parquet) (skip whole files) manifests
Parquet
Columnar, splittable, compressed per column, with a footer holding min/max statistics per row group. The default output format and the right one for anything that will be read more than once.
- Solves: column pruning (read 6 columns of 40 and you read ~15% of the bytes), predicate pushdown via row-group statistics, parallel reads of one file.
- Spark's vectorized Parquet reader is the fastest path in the JVM engine; the native engines all have their own.
df.write.parquet("s3a://ridehub/trips/")
Docs: Parquet files.
ORC
The other columnar format, from the Hive world: stripes instead of row groups, built-in indexes and optional bloom filters per column. Functionally equivalent to Parquet for Spark; pick whichever your ecosystem already speaks.
spark.sql.orc.impl=native
spark.sql.orc.filterPushdown=true
Docs: ORC files.
Avro
Row-oriented, schema embedded, splittable. The lingua franca of Kafka and schema registries.
- Solves: schema evolution and record-at-a-time producers.
- โน๏ธ Costs: a
SELECT one_columnreads every column of every row. Fine as a landing format for a streaming pipeline; convert to Parquet before anything analytical touches it.
Docs: Avro files.
CSV (uncompressed)
Text. Every read parses every byte, types are guesses unless you provide a schema, and quoting rules are a daily surprise.
- Splittable (Spark can start a task at any newline), so parallelism is fine. Throughput is not: parsing dominates.
- โน๏ธ Fine as a landing format. Convert on arrival.
โ ๏ธ gzip-compressed CSV or JSON
gzip has no block index, so a file cannot be split: one file is one task on one core, however large. A 50 GB events.csv.gz is a single task that runs for an hour while 399 cores sit idle. Use bzip2 or seekable zstd if you must compress text, split into many ~128 MB files, or convert to Parquet on landing.
JSON lines
Self-describing and verbose. The Jackson parser is the bottleneck; schema inference is a full extra pass over the data. Same advice as CSV, with more emphasis on always passing a schema.
Docs: JSON files, CSV files.
Delta Lake
Parquet files plus a transaction log (_delta_log/). The log lists exactly which files make up the current version, with per-file min/max statistics.
-
Solves: ACID writes (no half-written tables after a failed job), time travel (
VERSION AS OF),MERGE INTOfor upserts and GDPR deletes, data skipping from per-file statistics,OPTIMIZEandZORDER, and no directory listing: the file list comes from the log, which is why a Delta table with two million files still plans quickly. - Streaming source and sink with exactly-once semantics, which makes it the natural landing table for Structured Streaming.
spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension
spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog
df.write.format("delta").partitionBy("dt").save("s3a://ridehub/trips_delta/")
Docs: Delta Lake.
Apache Iceberg
An open table format designed for engine neutrality (Spark, Trino, Flink, Snowflake and BigQuery all read it natively).
-
Solves: the same ACID/time-travel/upsert problems as Delta, plus hidden partitioning (partition by
days(pickup_ts)without adtcolumn that users have to remember to filter on), manifest-level pruning (the planner skips whole manifests of files from their statistics, without listing or opening anything), and storage-partitioned joins (3.3+). - A REST catalog or Hive metastore holds the table pointer.
spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.catalog.lake=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.lake.type=rest
df.writeTo("lake.ridehub.trips").partitionedBy(F.days("pickup_ts")).create()
Docs: Iceberg with Spark.
Apache Hudi
The upsert-first table format, from Uber. Copy-on-write tables rewrite files on update; merge-on-read tables append deltas and compact later. Incremental queries ("give me everything that changed since commit X") are a first-class feature.
- Solves: CDC-style ingestion with record-level updates and incremental downstream pulls.
df.write.format("hudi").option("hoodie.datasource.write.recordkey.field", "trip_id").save(path)
Docs: Apache Hudi.
Hive metastore table (Parquet)
Not a format, but a different read path: plain Parquet files registered as a partitioned table in a metastore. Partition locations come from the catalog, not from walking S3, and bucketing metadata lives here too.
-
Solves: listing cost on partitioned tables, and it is the prerequisite for
bucketBy. -
Costs:
MSCK REPAIR TABLEorALTER TABLE ADD PARTITIONafter every write that creates a new partition, and a metastore to run.
df.write.partitionBy("dt").saveAsTable("ridehub.trips")
spark.sql("MSCK REPAIR TABLE ridehub.trips")
Choosing
| Need | Format |
|---|---|
| Landing zone from Kafka | Avro (or JSON if you must), converted downstream |
| Analytical tables, single engine, no updates | Parquet, partitioned by date |
| Analytical tables with updates, deletes, time travel | Delta (Databricks-centric) or Iceberg (multi-engine) |
| CDC ingestion with incremental consumers | Hudi |
| Millions of files and you cannot change the writer | a table format, purely for the manifest-based file list |
| Anything compressed with gzip | re-compress, or split into many files |
Layer 6b: Data layout and the read path
Same bytes, same format, radically different job: this layer is about where rows sit inside and across files so the planner can skip them, and about the handful of read-side settings that decide whether the first stage is a scan or a surprise.
partitionBy(date) on write
Hive-style directories, one per value: dt=2026-03-01/part-0000.parquet. A filter on dt never opens the other 334 days.
df.write.partitionBy("dt").parquet(path)
trips/
โโโ dt=2026-03-01/ โโโ
โโโ dt=2026-03-02/ โโโ filter dt BETWEEN '2026-03-01' AND '2026-03-31'
โ โฆ โ opens these 31 directories (1,400 files)
โโโ dt=2026-03-31/ โโโ
โโโ dt=2026-04-01/ and never lists or opens the other 334 (14,600 files)
โ โฆ
- Solves: the single biggest I/O win available, for the single most common filter (time).
- Pairs with pushdown-friendly predicates (Layer 3b); the playground gives the pair a synergy bonus because each is nearly useless without the other.
๐งจ partitionBy(user_id) and other high-cardinality keys
"Queries filter by user." 40 million users means 40 million directories with one tiny file each. Listing them takes longer than the old full scan did, the driver's file index runs out of memory, and every write creates thousands of files. Partition on low-cardinality columns only (date, region, tenant). For point lookups on high-cardinality keys, use clustering or bucketing below.
bucketBy(key) on write
Pre-hash rows into N buckets by the join or aggregation key. Later jobs that join or group by that key skip the shuffle entirely. Metastore tables only.
df.write.bucketBy(256, "driver_id").sortBy("driver_id").saveAsTable("ridehub.trips_b")
OPTIMIZE โฆ ZORDER BY (Delta)
Z-ordering rewrites files so that rows close in a multi-dimensional space (say, driver_id ร zone_id) land in the same file. Per-file min/max statistics then become selective for either column, and a point lookup opens a handful of files out of thousands.
spark.sql("OPTIMIZE ridehub.trips ZORDER BY (driver_id, zone_id)")
- Solves: the high-cardinality lookup that partitioning cannot. Costs: a rewrite of the table; schedule it.
Liquid clustering (Databricks)
Databricks' replacement for partitioning plus Z-ORDER: declare clustering columns once, and OPTIMIZE incrementally clusters new data without rewriting old data or picking partition boundaries up front.
CREATE TABLE trips CLUSTER BY (dt, driver_id) AS SELECT ...
Docs: Liquid clustering.
sortWithinPartitions before write
The cheapest clustering there is. Sorting each output partition by a column means each file (and each row group inside it) covers a narrow range of that column, so Parquet's min/max statistics finally skip data on it. The previous article's fare > 300 filter went from "every row group might match" to "3% of row groups might match" with one extra method call.
df.sortWithinPartitions("fare").write.partitionBy("dt").parquet(path)
Compaction to 128 MBโ1 GB files
Streaming sinks, hourly jobs and partitionBy(user_id) all produce small files. Every file costs a listing entry, an S3 GET for the footer, a task scheduling decision and a few milliseconds of fixed overhead. Two million 400 KB files is 800 GB of data with the overhead of 2,000,000 opens. A compaction job rewrites them into a few thousand ~256 MB files.
# Delta / Databricks
spark.sql("OPTIMIZE ridehub.trips")
# Iceberg
spark.sql("CALL lake.system.rewrite_data_files(table => 'ridehub.trips')")
# plain Parquet
spark.read.parquet(src).repartition(200).write.mode("overwrite").parquet(dst)
repartition / coalesce before write
Output file count equals the number of partitions in the final stage, which after a shuffle is spark.sql.shuffle.partitions: 200 files, however small the result. Decide on purpose.
result.coalesce(8).write.parquet(out) # fewer files, no shuffle
result.repartition(64, "dt").write.partitionBy("dt").parquet(out) # one well-sized file per dt
-
coalesce(n)merges partitions without a shuffle but can leave them uneven;repartition(n, col)shuffles and balances. ForpartitionBywrites,repartitionby the partition column avoids every task writing a sliver into every directory.
Delta optimized writes and auto compaction
Databricks' automatic version of the two parts above: an adaptive shuffle before the write to produce ~128 MB files, then a compaction of small files after each commit.
spark.databricks.delta.optimizeWrite.enabled=true
spark.databricks.delta.autoCompact.enabled=true
Explicit schema on read
For CSV and JSON, pass a schema. No inference pass, deterministic types, and no "the column was all nulls on Tuesday so it became a string" incident.
schema = "trip_id string, driver_id long, zone_id int, pickup_ts timestamp, fare double"
spark.read.schema(schema).json("s3a://ridehub/raw/")
โ ๏ธ inferSchema = True
For CSV and JSON this is a full extra pass over every byte before the real job starts, and the inferred types can change between runs. Parquet and ORC carry their schema in the footer, so the option does nothing there.
maxPartitionBytes = 512 MB
spark.sql.files.maxPartitionBytes (default 128 MB) is the target size of a scan task; openCostInBytes (default 4 MB) is the fixed cost charged per file when bin-packing. Raising the target halves the number of scan tasks, which helps when files are many and the work per byte is small; lowering it helps when a few huge files leave cores idle.
spark.sql.files.maxPartitionBytes=512m
spark.sql.files.openCostInBytes=8m
Cloud-native output committer
Spark's default commit protocol writes to a temporary directory and renames on commit. On HDFS a rename is a metadata operation; on S3 it is a copy and delete of every byte, and it is not atomic. The S3A magic committer (and EMR's S3-optimized committer, and GCS's equivalent) use multipart uploads that are only completed at commit time: no copy, no window of half-visible output.
spark.hadoop.fs.s3a.committer.name=magic
spark.sql.sources.commitProtocolClass=org.apache.spark.internal.io.cloud.PathOutputCommitProtocol
spark.sql.parquet.output.committer.class=org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter
- โน๏ธ Table formats solve this one level up: Delta, Iceberg and Hudi make a write visible by committing to their log, so the files underneath never need renaming.
Docs: Committing work to S3 with the S3A committers.
mergeSchema on read
When Parquet files in one directory have evolved schemas, mergeSchema=true reads every footer and reconciles them. Correct, and on 16,000 files it is 16,000 footer reads on the driver before the job starts. Table formats track schema evolution in their log and avoid this.
๐งจ ignoreCorruptFiles = true
"The job keeps failing on one bad file." Now it silently skips any file it cannot read, including the one that failed because S3 throttled the request and the one that was half-written by a job that died. Missing rows, no error, and nobody finds out until the finance numbers do not add up. Quarantine bad files explicitly instead.
zstd Parquet compression
Snappy is the default Parquet codec: fast, moderate ratio. zstd gives noticeably smaller files (often 20โ40%) at similar read speed on modern CPUs. Less S3 storage, less network, slightly more CPU on write.
spark.sql.parquet.compression.codec=zstd
Reuse one DataFrame per source
spark.read.parquet(path) builds an InMemoryFileIndex by listing the directory tree. Calling it again lists again. Read once, filter many times, and the listing (which on 16,000 S3 objects is several seconds of API calls) is paid once per session.
trips = spark.read.parquet(path) # lists once
march = trips.filter(...)
april = trips.filter(...) # no second listing
Parquet column index and bloom filters
Parquet 1.11+ writes page-level min/max (the column index), so pushdown can skip pages inside a row group, not just whole row groups. Bloom filters, enabled per column at write time, answer "is this trip_id definitely not in this row group?" for point lookups on high-cardinality keys where min/max statistics are useless.
spark.hadoop.parquet.bloom.filter.enabled#trip_id=true
spark.hadoop.parquet.bloom.filter.expected.ndv#trip_id=50000000
Choosing
| Symptom | Lever |
|---|---|
| Date filter still scans the whole table |
partitionBy(dt) + a plain predicate on dt
|
| Point lookups on a high-cardinality key | Z-ORDER / liquid clustering / bloom filters, not partitioning |
| Listing takes minutes | compaction, table format, reuse the DataFrame |
| 200 tiny output files |
coalesce / repartition before write, optimized writes, AQE coalesce |
| Range filter on a non-partition column reads everything |
sortWithinPartitions on that column at write time |
| Job spends minutes before stage 0 |
inferSchema on text, mergeSchema, a cold listing: pass a schema, use a table format |
| Half-written output visible after a failure | cloud committer or a table format |
Combinations that are worth more than their parts
A few of these parts only make sense together, and the playground gives them explicit synergy bonuses. They are worth listing because each is a small design pattern:
| Combination | Why it is more than the sum |
|---|---|
| AQE + coalesce + skew join + runtime broadcast | Self-tuning shuffles. The planner's three biggest mistakes (too many reducers, one hot reducer, wrong join strategy) all fixed after the fact. |
partitionBy(dt) + pushdown-friendly predicates |
Partition pruning on every read. Either one alone does nothing. |
pandas UDF + Arrow + memoryOverhead
|
Production-grade PySpark. Vectorized logic, fast transfer, and room for the worker to live. |
| Celeborn + dynamic allocation + spot | Elastic, spot-tolerant shuffle. Stateless executors can be added, removed or reclaimed at any moment. |
| Spot + graceful decommissioning | Nodes drain instead of crash. The two-minute notice becomes enough time. |
| Kubernetes + Celeborn + dynamic allocation | Stateless autoscaling on Kubernetes. The missing shuffle service, supplied. |
Structured Streaming + RocksDB + watermark + idempotent foreachBatch
|
A streaming pipeline that runs for months. Bounded state, off-heap, exactly-once. |
| Gluten / Comet + off-heap | A native engine with the memory it needs. Without it, you get the fallbacks and none of the speed. |
| Iceberg + bucketed joins | Storage-partitioned joins. The big-to-big join with no Exchange. |
Thrift server + cache() + FAIR pools |
A shared BI endpoint that stays responsive. Warm data, fair slots. |
| Delta + Z-ORDER + compaction | Point lookups on a lakehouse table. Few files, each covering a narrow key range. |
| Event logs + Prometheus | An observable cluster. Post-mortems and trends. |
Six starter profiles
The playground ships these as one-click presets. Each is a defensible starting point for its workload, not a tuned endpoint; the point is to see which layers a given problem actually touches.
๐ Nightly ETL starter
YARN ยท DataFrame API ยท built-in functions ยท JVM engine ยท AQE (coalesce, skew) ยท broadcast hint ยท pushdown-friendly predicates ยท external shuffle service ยท zstd shuffle ยท 4-core executors ยท dynamic allocation ยท History Server ยท Parquet ยท partitionBy(dt) ยท coalesce before write
The RideHub job after Maya fixed it. Nothing exotic: every component is in open-source Spark, and the whole configuration fits in a dozen lines.
spark.master=yarn
spark.submit.deployMode=cluster
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.skewJoin.enabled=true
spark.shuffle.service.enabled=true
spark.io.compression.codec=zstd
spark.executor.cores=4
spark.executor.memory=20g
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.maxExecutors=50
spark.eventLog.enabled=true
๐ Interactive SQL endpoint
Kubernetes ยท DataFrame API ยท Thrift server ยท JVM ยท AQE (runtime broadcast) ยท CBO statistics ยท Celeborn ยท cache() the hot week ยท FAIR pools ยท locality.wait = 0 ยท Delta ยท Z-ORDER ยท compaction ยท one DataFrame per source
Dashboards at nine o'clock. The levers are all about latency: a warm long-running session, fair sharing, few well-clustered files, and no three-second locality waits on object storage.
๐ PySpark ML features
Kubernetes ยท DataFrame API ยท pandas UDF ยท mapInPandas ยท JVM ยท AQE (coalesce) ยท shuffle tracking ยท memoryOverhead ยท Arrow ยท 4-core executors ยท dynamic allocation ยท stage-level scheduling ยท Parquet ยท partitionBy(dt) ยท pushdown-friendly predicates
Python at scale without the 40ร penalty. Note what is absent: no row-at-a-time UDF, and no native engine (it would fall back on the UDF anyway).
๐ Streaming Kafka โ Delta
Kubernetes ยท DataFrame API ยท Structured Streaming ยท RocksDB ยท watermark ยท idempotent foreachBatch ยท JVM ยท AQE ยท 4-core executors ยท memoryOverhead ยท generous retries ยท History Server ยท Delta ยท optimized writes
Runs for months. Bounded state off the heap, a sink that tolerates replay, retries that tolerate Kafka hiccups, and a table format that compacts its own small files.
๐๏ธ Cost-optimised spot batch
Kubernetes ยท DataFrame API ยท built-in functions ยท JVM ยท AQE (coalesce, skew) ยท Celeborn ยท zstd shuffle ยท 4-core executors ยท dynamic allocation ยท spot workers ยท graceful decommissioning ยท generous retries ยท Parquet ยท partitionBy(dt) ยท pushdown-friendly predicates ยท zstd Parquet ยท coalesce before write
The nightly ETL at 30% of the price. Every resilience lever is on because every node is disposable: shuffle lives off the executors, nodes drain when reclaimed, and eight retries absorb the churn.
๐ Native engine ETL (Gluten)
YARN ยท DataFrame API ยท built-in functions ยท Gluten + Velox ยท AQE (coalesce) ยท broadcast hint ยท pushdown-friendly predicates ยท external shuffle service ยท off-heap memory ยท 4-core executors ยท dynamic allocation ยท Parquet ยท partitionBy(dt) ยท compaction
The CPU-bound ETL, roughly twice as fast. Built-in functions only (no fallbacks), off-heap sized for Velox, and compacted inputs so the native scan reads big files.
Cheat sheet: symptom to lever
| You see | It usually means | Reach for |
|---|---|---|
PartitionFilters: [] on a partitioned table |
predicate is not on the partition column, or wraps it in a function | rewrite the predicate; DPP for star joins |
| 16,000 scan tasks for one month of data | no pruning, or tiny files | pruning above; compaction; maxPartitionBytes
|
SortMergeJoin with two Exchanges against a small table |
planner sized the dimension from files | broadcast hint; CBO stats; AQE runtime broadcast |
Could not execute broadcast in 300 secs / driver OOM in a join |
broadcast threshold too high, or a hint on a big table | lower the threshold; bigger driver; let it SMJ |
| 199 tasks in 1 minute, 1 task in 3 hours | hot key | AQE skew join; salting for aggregations |
BatchEvalPython in the plan, 40ร slower |
row-at-a-time Python UDF | built-in functions; pandas UDF; Arrow UDF |
Container killed by YARN for exceeding memory limits, heap free |
off-JVM memory (Python / Arrow / native) |
memoryOverhead; pyspark.memory; off-heap |
| 200 output files of 2 MB | default shuffle partitions | AQE coalesce; coalesce(); optimized writes |
One task per .csv.gz, hours long |
unsplittable compression | re-compress; split; convert to Parquet |
| Job spends 5 minutes before stage 0 | listing, inferSchema, mergeSchema
|
schema; table format; reuse DataFrame |
FetchFailedException after a node vanished |
shuffle files died with the node | ESS; remote shuffle; decommissioning |
| Dynamic allocation never scales down on Kubernetes | shuffle tracking holding executors | Celeborn / Uniffle |
| Dashboards freeze while a backfill runs | FIFO within one session | FAIR pools |
| Streaming executor OOM after days | unbounded state on heap | watermark; RocksDB |
| Duplicate rows after a retry | non-idempotent sink with speculation or restarts |
foreachBatch keyed by batchId; table format |
| Native engine slower than JVM | fallbacks on every operator | count fallbacks in explain(); remove UDFs; size off-heap |
| Job silently produced fewer rows | ignoreCorruptFiles |
turn it off; quarantine explicitly |
Play it: the Distributed Compute Playground
Reading a parts catalogue is one thing. Bolting the wrong parts together and watching the task count explode is how it sticks.
So I turned this whole article into a game. ๐ฎ
Distributed Compute Playground is a drag-and-drop sandbox where you assemble a Spark job from exactly the nine layers above. Every component in this article (all 104 of them) is a chip. Drop it in and the playground tells you, instantly, what Spark would do with the job you just built:
- the
explain()-style physical plan, withPartitionFilters, the join strategy andAQEShuffleReadfilled in from your choices - the stage DAG with task counts and bytes per stage
- how many output files, how many Spark jobs, and roughly how long it takes on your cluster
Drop in month(pickup_ts) = 3 and watch PartitionFilters: [] come back with 16,000 tasks. Swap it for a range on dt and watch it fall to 1,400. Add a broadcast hint and watch the SortMergeJoin turn into a BroadcastHashJoin. Raise the threshold to 2 GB and watch the memory-safety score fall off a cliff.
What's inside
- ๐งฑ Build it: 104 real Spark components across 9 colour-coded layers. Prerequisites and conflicts are enforced: dynamic allocation without a shuffle home is an error, Photon on YARN "runs nowhere", a Gluten plan with a Python UDF gets a warning about fallbacks.
- ๐ฌ Read the plan it would produce: a live job simulator for 10 workloads, from the nightly ETL to two million tiny files, a hot-key join, a 5 TB-to-5 TB join and a Kafka stream.
- ๐ Get scored: the eleven dimensions from this article, a grade from S to D, and loud warnings when you do something reckless (yes,
spark.task.maxFailures = 1, I'm looking at you ๐). - ๐ต๏ธ Crack 25 missions: on-call tickets written like cloud certification exam questions. Each is a ~300-word story about a company in trouble, with a hidden question and goals that check the plan, not just the score. Mission 1 is FIN-2231, the ticket from the previous article โ except the company is now Capsule Corp and the engineer is Sakura, because every company and person in the playground borrows a name from the Naruto or Dragon Ball Z universe (no real company ever gets quoted). Then there is the three-hour straggler, the job that only fails on Fridays, the GDPR delete by Friday, and the Photon bill.
- ๐ Light and dark themes, zero dependencies, zero network calls. Nothing leaves your browser.
Up and running in 10 seconds
git clone https://github.com/AnikethSD/distributed-computing-playground.git
cd distributed-computing-playground
open index.html # macOS; use xdg-open on Linux, or just double-click it
No npm install. No Docker. No cluster. Just a browser.
One small favour โญ
If this article or the game saved you even one evening of staring at a Spark UI, please drop a โญ on the repo.
It takes two seconds and costs nothing. It's the single biggest signal that tells me to keep building (the .write article and a second engine family are next ๐), and it helps other engineers find the project.
Got a production incident that would make a brutal mission? A straggler story? A plan that fooled you? Issues and PRs are very welcome. There's a mission authoring guide to get you started, and if databases are more your thing, its sibling, the Database Playground, is waiting too.
Further reading
The companion article
-
What Actually Happens When You Call
spark.read? One Line of Python, a Thousand Tasks: the story version, following one job through Catalyst, the file index, the scheduler and the shuffle.
Official Spark docs, by layer
- L1: Cluster mode overview, Standalone, YARN, Kubernetes
- L2: Arrow and pandas UDFs, Spark Connect, Thrift server, Structured Streaming
- L3: SQL performance tuning (AQE, join hints, CBO), Apache Gluten, Apache DataFusion Comet, Photon, RAPIDS Accelerator
- L4: Shuffle behavior configuration, Push-based shuffle, Apache Celeborn, Apache Uniffle, Memory management, Hardware provisioning
- L5: Job scheduling (dynamic allocation, FAIR pools), Stage-level scheduling, Monitoring
- L6: Parquet, ORC, Avro, CSV, JSON, Delta Lake, Apache Iceberg, Apache Hudi, S3A committers
- Everything: Spark configuration reference
Longer reads
- The Internals of Spark SQL by Jacek Laskowski, operator by operator
- Spark: The Definitive Guide (Chambers & Zaharia), Part II, for the execution model in long form
This is the second article in a series on distributed computing, told through one dataset and one team. The first followed a job through Spark's internals; this one catalogued the parts. Next: what happens when you call .write, and why your output has 200 files. Spark moves fast and platforms move faster: if a default or a feature here does not match what you see, tell me in the comments and I will fix it.


Top comments (1)
Really well explained, Aniketh! ๐ฅ Love how you broke down the different knobs in a Spark job and made distributed compute feel much easier to understand. Great read! ๐