DEV Community

Cover image for The Distributed Compute Playground: Everything You Can Plug In and Tune in an Apache Spark Job
Aniketh Deshpande
Aniketh Deshpande

Posted on

The Distributed Compute Playground: Everything You Can Plug In and Tune in an Apache Spark Job

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 --conf line or PySpark call for every one of them.

โฑ๏ธ 40-minute read. It is the reference companion to
What actually happens when you call spark.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 the explain() 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

  1. How to read this article
  2. Layer 1: Cluster manager
  3. Layer 2: Code and APIs
  4. Layer 3a: Execution engine
  5. Layer 3b: Planner, joins and AQE
  6. Layer 4a: Shuffle
  7. Layer 4b: Memory, caching and serialization
  8. Layer 5: Scheduling and resilience
  9. Layer 6a: File and table formats
  10. Layer 6b: Data layout and the read path
  11. Combinations that are worth more than their parts
  12. Six starter profiles
  13. Cheat sheet: symptom to lever
  14. Play it: the Distributed Compute Playground ๐ŸŽฎ
  15. 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    โ”‚
 โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜
Enter fullscreen mode Exit fullscreen mode

Each section below takes one layer, lists every part you can plug into it, and for each part answers the same four questions:

  1. What problem does it solve?
  2. What does it cost? (Every lever has a price. Some are hidden.)
  3. What does it need? (Prerequisites and conflicts. Dynamic allocation without a shuffle service is a config error, not a tuning choice.)
  4. How do you turn it on? (The --conf or 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_ONLY persistence 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
   โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜
Enter fullscreen mode Exit fullscreen mode

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 on local[*] 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[*]
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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: cluster puts the driver inside the YARN application (survives your SSH session dying); client keeps it where you ran spark-submit.
spark.master=yarn
spark.submit.deployMode=cluster
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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")))
Enter fullscreen mode Exit fullscreen mode

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, mapPartitions with 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
  • โš ๏ธ 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 a filter means the scan reads everything. Look for BatchEvalPython in explain(); 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
Enter fullscreen mode Exit fullscreen mode

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()
Enter fullscreen mode Exit fullscreen mode

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")
Enter fullscreen mode Exit fullscreen mode
  • 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 memoryOverhead in Layer 4b. The classic failure is Container killed by YARN for exceeding memory limits with 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()
Enter fullscreen mode Exit fullscreen mode
  • โ„น๏ธ 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. Set ps.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"))
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
  • โ„น๏ธ RDD APIs and SparkContext calls are not available over Connect. If your code uses sc.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())
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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") ...)
Enter fullscreen mode Exit fullscreen mode

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()
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode

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.size is 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
Enter fullscreen mode Exit fullscreen mode

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 a Comet prefix, 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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) */ ...
Enter fullscreen mode Exit fullscreen mode
  • 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.memory at the same time, every time.
spark.sql.autoBroadcastJoinThreshold=100m
spark.driver.memory=8g
Enter fullscreen mode Exit fullscreen mode

๐Ÿงจ 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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 ANALYZE on 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
Enter fullscreen mode Exit fullscreen mode
spark.sql("ANALYZE TABLE trips COMPUTE STATISTICS FOR ALL COLUMNS")
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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"])
Enter fullscreen mode Exit fullscreen mode
  • 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")
Enter fullscreen mode Exit fullscreen mode
  • 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
Enter fullscreen mode Exit fullscreen mode

๐Ÿงจ 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)
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

โš ๏ธ 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
  • โ„น๏ธ 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
Enter fullscreen mode Exit fullscreen mode

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     โ”‚
  โ”‚ โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ดโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜ โ”‚                  โ”‚
  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

โš ๏ธ 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
  • 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

๐Ÿงจ 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.
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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()
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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 foreach that 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "dashboards")
Enter fullscreen mode Exit fullscreen mode
  • 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
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode
  • 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
Enter fullscreen mode Exit fullscreen mode

๐Ÿงจ 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
Enter fullscreen mode Exit fullscreen mode
  • 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
Enter fullscreen mode Exit fullscreen mode

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/
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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/")
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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_column reads 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 INTO for upserts and GDPR deletes, data skipping from per-file statistics, OPTIMIZE and ZORDER, 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
Enter fullscreen mode Exit fullscreen mode
df.write.format("delta").partitionBy("dt").save("s3a://ridehub/trips_delta/")
Enter fullscreen mode Exit fullscreen mode

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 a dt column 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
Enter fullscreen mode Exit fullscreen mode
df.writeTo("lake.ridehub.trips").partitionedBy(F.days("pickup_ts")).create()
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode

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 TABLE or ALTER TABLE ADD PARTITION after 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")
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode
   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)
   โ”‚   โ€ฆ
Enter fullscreen mode Exit fullscreen mode
  • 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")
Enter fullscreen mode Exit fullscreen mode

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)")
Enter fullscreen mode Exit fullscreen mode
  • 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 ...
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
  • coalesce(n) merges partitions without a shuffle but can leave them uneven; repartition(n, col) shuffles and balances. For partitionBy writes, repartition by 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
Enter fullscreen mode Exit fullscreen mode

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/")
Enter fullscreen mode Exit fullscreen mode

โš ๏ธ 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode
  • โ„น๏ธ 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

๐Ÿ“Š 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, with PartitionFilters, the join strategy and AQEShuffleRead filled 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.

The Job plan tab: explain() output, stage table and estimate for the nightly ETL

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.

The mission board: 25 on-call tickets across 8 categories

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
Enter fullscreen mode Exit fullscreen mode

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

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)

Collapse
 
pranav_nandan_eb5ba9a76e8 profile image
Pranav Nandan •

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! ๐Ÿ‘