Every Apache Spark release, and the problem each one was built to solve
A release-by-release walk through Spark's history, from the research releases to the current line: the theory behind what each version added, which defaults it changed under you, what it deleted, and a runnable example for the features that matter. Every example run against a live cluster rather than recalled.
- How to run the examples
- Architecture: the four eras, and what moved between them
- Before 1.0: the research releases
- The 1.x line: structure arrives
- The 2.x line: unification and codegen
- The 3.x line: the planner stops guessing
- The 4.x line: SQL becomes the surface again
- Which defaults changed underneath you?
- What was removed, and when?
- Common misconceptions
- Which version should you be on?
- Frequently asked questions
- The mental model
- References
- Trademarks
TL;DR
- Spark’s centre of gravity has moved four times: from the RDD, to the DataFrame, to the query plan, to a protocol between client and driver. Every release makes more sense once you know which of those four eras it belongs to.
- The biggest upgrades are rarely the headline features. They are the default changes: sort shuffle, a 10 MB broadcast threshold, adaptive execution, bloom-filter joins, ANSI mode and Arrow-backed Python UDFs all flipped on in different releases and silently changed what your existing jobs do.
- Some of the oldest defaults have never moved. The broadcast-join threshold is still
10485760bandspark.sql.shuffle.partitionsis still200, eleven years after they were set.- Most features belong to a handful of long threads — Python performance, streaming state, SQL completeness, runtime re-planning — and each thread took several releases to finish. Knowing where a feature sits on its thread tells you whether it is experimental, usable or default.
- The same query can return a number, a
NULL, a wrapped-around negative or an exception depending on the release. Integer overflow is the cleanest example, and it is shown running on both demo lines below.- The newest release ships a SQL surface ahead of its engine in places. The change-data-capture clause parses and analyses, then stops at the catalog, because the built-in one does not implement it.
You have inherited a Spark job. It is pinned to some version because that is what the platform team installed in 2021, and someone is now asking whether moving it forward is worth a quarter of engineering time. The release notes are no help: each one is a list of between 1,100 and 5,100 resolved tickets, sorted by component, with no indication of which three entries will change your job’s behaviour and which 3,000 you can ignore.
That is the question this post answers. Not “what is in each release”, which the project already documents, but which changes altered the execution model, which changed a default under you, and which deleted something you might still be using.
Everything below was checked against a running cluster rather than recalled. The version under test is 4.1.3 unless the text says otherwise, and the examples filed under the 3.x releases were run again on 3.5.9 to confirm they are not 4.x-only. Before-and-after contrasts with older behaviour use 3.5.9, the oldest line this blog demos on. The one exception is the 4.2 section: features that do not exist in 4.1.3 were run on 4.2.0, and are labelled as such. A good number of things I was confident about turned out to be wrong when I ran them, and I have kept those corrections in the text where they are useful.
Every release section lists its features one by one. Each feature starts with its name in bold, followed by what it is and the theory behind it, a runnable example where one makes sense, and a Good to know note for the detail that trips people up. The examples are self-contained, so any one of them can be pasted into a file and run on its own.
How to run the examples
Every Python example is a standalone script. The fastest way to run one is the official image, with the script’s directory mounted:
docker run --rm -v "$PWD":/w -w /w apache/spark:4.1.3-python3 \
/opt/spark/bin/spark-submit --master "local[4]" example.py
The apache/spark images do not include pandas or PyArrow. The pandas UDF,
pandas-on-Spark and stateful-processor examples need them, and the Spark
Connect and Declarative Pipelines examples also need the Connect client
packages:
pip install pandas pyarrow # pandas UDFs, applyInPandas, pyspark.pandas
pip install grpcio grpcio-status protobuf \
googleapis-common-protos zstandard pyyaml # Spark Connect, spark-pipelines
One trap from my own setup: if the image you use carries a spark-defaults.conf
pointing at object storage or a Hive metastore, Python’s tempfile.mkdtemp()
paths get resolved against that default file system, and the streaming examples
fail with an UnknownHostException. Pointing SPARK_CONF_DIR at an empty
directory runs everything as plain local Spark.
Architecture: the four eras, and what moved between them
Spark has had one durable architectural habit: each era moves the thing you program against one level further from the machine, and one level closer to a description of intent.
flowchart LR
E1["<b>1.x</b><br/>the RDD era<br/><i>you describe<br/>the computation</i>"]
E2["<b>2.x</b><br/>the DataFrame era<br/><i>you describe<br/>the data</i>"]
E3["<b>3.x</b><br/>the plan era<br/><i>the engine revises<br/>its own decisions</i>"]
E4["<b>4.x</b><br/>the protocol era<br/><i>the client is<br/>detached from the driver</i>"]
E1 --> E2 --> E3 --> E4
Each move buys something and costs something:
The 1.x RDD era: you describe the computation
You program against an RDD of ordinary JVM or Python objects, and you pass
functions to map, filter and reduceByKey. That buys total control: any
type, any logic. The cost is that the engine cannot see inside your closures.
It does not know which fields a lambda reads or what a filter keeps, so it cannot
reorder, prune or push anything down. Every optimisation is yours to do by hand.
The 2.x DataFrame era: you describe the data
You program against a DataFrame over a schema, and your operations become an expression tree the engine can read. That buys whole-stage code generation, columnar scans, predicate pushdown and one API for batch and streaming. The cost is that anything the engine cannot express, above all a Python UDF, becomes the slow path, because the optimizer has to treat it as a black box.
The 3.x plan era: the engine revises its own decisions
You write the same DataFrame, but the plan is revised at runtime using measured sizes rather than estimates. That buys decisions that fit the data: partitions coalesced, skew split, joins switched to broadcast. The cost is that a plan you read before execution is no longer the plan that ran, so debugging means reading the final adaptive plan, not the initial one.
The 4.x protocol era: the client is detached from the driver
Your program sends a logical plan over the wire to a Spark Connect server. That buys thin clients in many languages, isolation from the driver JVM, and several users sharing one server. The cost is two API surfaces, classic and Connect, and features that sometimes arrive on one before the other.
The version-by-version sections follow that arc. If you only read one section, read Defaults that changed underneath you, because default changes cause more upgrade surprises than features do.
Every release at a glance
Each line below gives the month of the .0 release, its headline features, and
the final maintenance release in that line according to the project’s release
archive, or the newest one where the line is still maintained. Every feature
named here has its own entry, with a description and an example, further down.
- Spark 0.5 (June 2012): Mesos 0.9 support,
sortByKey, faster shuffles. Final patch 0.5.2. - Spark 0.6 (October 2012): standalone deploy mode and the Java API. Final patch 0.6.2.
- Spark 0.7 (February 2013): PySpark and Spark Streaming as an alpha. Final patch 0.7.3.
- Spark 0.8 (September 2013): MLlib, the web UI on port 4040, YARN in mainline. Final patch 0.8.1.
- Spark 0.9 (February 2014): GraphX alpha,
SparkConf, Scala 2.10. Final patch 0.9.2. - Spark 1.0 (May 2014):
spark-submit, Spark SQL alpha, API stability. Final patch 1.0.2. - Spark 1.1 (September 2014): sort shuffle available, JDBC server. Final patch 1.1.1.
- Spark 1.2 (December 2014): sort shuffle and Netty by default, data source API. Final patch 1.2.2.
- Spark 1.3 (March 2015):
DataFrame, direct Kafka stream. Final patch 1.3.1. - Spark 1.4 (June 2015): SparkR, window functions, sort-merge join. Final patch 1.4.1.
- Spark 1.5 (September 2015): Tungsten: codegen, binary rows, managed memory. Final patch 1.5.2.
- Spark 1.6 (January 2016):
Dataset, unified memory, the first adaptive execution. Final patch 1.6.3. - Spark 2.0 (July 2016):
SparkSession, whole-stage codegen, Structured Streaming. Final patch 2.0.2. - Spark 2.1 (December 2016): event-time watermarks. Final patch 2.1.3.
- Spark 2.2 (July 2017): Structured Streaming GA, cost-based optimizer. Final patch 2.2.3.
- Spark 2.3 (February 2018): Kubernetes, pandas UDFs, stream-stream joins. Final patch 2.3.4.
- Spark 2.4 (November 2018): higher-order functions, Avro, barrier mode. Final patch 2.4.8.
- Spark 3.0 (June 2020): adaptive execution rebuilt, dynamic partition pruning. Final patch 3.0.3.
- Spark 3.1 (March 2021): Kubernetes GA, ANSI mode errors. The line starts at 3.1.1; there is no 3.1.0. Final patch 3.1.3.
- Spark 3.2 (October 2021): AQE by default, pandas API, RocksDB state store. Final patch 3.2.4.
- Spark 3.3 (June 2022): bloom-filter runtime filters, error classes. Final patch 3.3.4.
- Spark 3.4 (April 2023): Spark Connect, parameterized SQL,
DEFAULTcolumns. Final patch 3.4.4. - Spark 3.5 (September 2023): Python UDTFs, Arrow Python UDFs, Connect clients. Latest patch 3.5.9, still maintained.
- Spark 4.0 (May 2025): ANSI by default,
VARIANT, collations, Scala 2.13 only. Latest patch 4.0.4. - Spark 4.1 (December 2025): Declarative Pipelines, SQL scripting GA, recursive CTEs. Latest patch 4.1.3.
- Spark 4.2 (July 2026): geospatial types, CDC syntax, Arrow Python by default. 4.2.0 is current.
The threads that run through the releases
Reading releases one at a time hides the fact that most features arrive in stages. A capability is introduced as experimental, becomes usable a release or two later, and becomes a default a release or two after that. Five of those threads run across the whole history, and placing a feature on its thread tells you whether it is experimental, usable or finished.
Python performance
It started with row-at-a-time Python lambdas on RDDs in 0.7 and Python SQL UDFs
in 1.1, where every row is pickled and sent to a Python worker. It became usable
at scale with pandas UDFs over Arrow in 2.3, which ship whole batches, and
type-hinted pandas UDFs in 3.0. It is finishing with Arrow-optimized plain UDFs
in 3.5 (opt-in), the arrow_udf decorator in 4.1, and Arrow on by default in
4.2, so ordinary UDFs get the fast path without being rewritten.
Streaming
It started with DStreams in 0.7, a stream as a sequence of small RDDs. It became usable with Structured Streaming in 2.0, a stream as an ever-growing table, watermarks in 2.1 to bound state, and GA in 2.2. It matured with the RocksDB state store in 3.2, the State API v2 and a state reader in 4.0, and real-time mode in 4.1 for sub-second latency.
Runtime re-planning
It started with reducer-count selection in 1.6. It became usable with the AQE rewrite and dynamic partition pruning in 3.0, and finished when AQE was switched on by default in 3.2 and bloom-filter runtime filters in 3.4.
SQL correctness
It started with ANSI mode as an opt-in in 3.0 and raising errors under it in 3.1. It became usable with ANSI GA in 3.2, error classes in 3.3 and SQLSTATE codes in 3.4, and finished when ANSI mode became the default in 4.0.
Client and driver separation
It started with Spark Connect for Python in 3.4. It became usable with the Scala and Go clients, streaming and pandas on Connect in 3.5, and matured in 4.1 with a Connect JDBC driver, ML on Connect GA, and Declarative Pipelines built as a Connect client.
The pattern holds well enough to use as a rule. When a feature first appears in a release note, expect roughly two releases before it is safe to depend on and two more before it is the default. Kubernetes support took from 2.3 to 3.1; adaptive execution took from 3.0 to 3.2; ANSI mode took from 3.0 to 4.0.
Before 1.0: the research releases
Spark was written at UC Berkeley’s AMPLab and open-sourced in 2010. The public releases before 1.0 are worth knowing about for one reason: they set down every primitive the 1.x line then built on. The RDD, lineage-based recovery, the shell, caching and the first three language bindings were all in place before anything was declared stable.
In this section and every release section below, each feature is written the same way: the feature name in bold, what it is and why it matters, a runnable example where one makes sense, and a Good to know note for anything that trips people up.
Resilient distributed datasets and lineage (0.5 and earlier)
The theory that everything rests on. An RDD is an immutable, partitioned collection plus a record of how to compute each partition from its parents. That record, the lineage, is the fault-tolerance mechanism: if an executor dies, Spark does not restore a checkpoint, it recomputes the lost partitions from their parents. Transformations only extend the lineage; nothing runs until an action asks for a result.
You can print a lineage, and it still looks the way it did in these releases.
Each indentation step with +- is a shuffle boundary, which is also where the
scheduler cuts the job into stages:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("rdd-wordcount").getOrCreate()
sc = spark.sparkContext
lines = sc.parallelize(["spark is fast", "spark is lazy", "rdds are immutable"], 2)
counts = (lines.flatMap(lambda l: l.split())
.map(lambda w: (w, 1))
.reduceByKey(lambda a, b: a + b))
print(counts.toDebugString().decode())
print(sorted(counts.collect()))
spark.stop()
On 4.1.3:
(2) PythonRDD[5] at RDD at PythonRDD.scala:59 []
| MapPartitionsRDD[4] at mapPartitions at PythonRDD.scala:171 []
| ShuffledRDD[3] at partitionBy at NativeMethodAccessorImpl.java:0 []
+-(2) PairwiseRDD[2] at reduceByKey at /w/10_rdd_lineage.py:9 []
| PythonRDD[1] at reduceByKey at /w/10_rdd_lineage.py:9 []
| ParallelCollectionRDD[0] at readRDDFromFile at PythonRDD.scala:300 []
[('are', 1), ('fast', 1), ('immutable', 1), ('is', 2), ('lazy', 1), ('rdds', 1), ('spark', 2)]
The (2) is the partition count, and the single +- marks the one shuffle that
reduceByKey introduces. Two stages, two tasks each.
Good to know: flatMap and map do not appear as separate RDDs. PySpark
pipelines consecutive narrow Python transformations into one PythonRDD, here
the one labelled reduceByKey, which also holds the map-side combine. All of it
runs in a single pass of the Python worker over each partition.
Pair-RDD operators: sortByKey and takeSample (0.5)
0.5 filled out the key-value operators on RDDs of pairs. sortByKey does a
range-partitioned global sort, sampling the keys first to pick partition
boundaries, and takeSample returns a fixed-size random sample to the driver:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("pairs05").getOrCreate()
sc = spark.sparkContext
sales = sc.parallelize([("pune", 30), ("goa", 5), ("hyderabad", 12), ("pune", 7)])
print(sales.sortByKey().collect())
# [('goa', 5), ('hyderabad', 12), ('pune', 30), ('pune', 7)]
print(sales.reduceByKey(lambda a, b: a + b).sortBy(lambda kv: -kv[1]).first())
# ('pune', 37)
print(len(sales.takeSample(False, 2, seed=42)))
# 2
spark.stop()
Good to know: sortByKey runs an extra job before the sort to sample key
boundaries, which is why a sort shows up as two jobs in the UI.
Mesos support (0.5)
Spark was originally built as a demonstration framework for Mesos, and 0.5 tracked Mesos 0.9. Mesos stayed a supported cluster manager until it was removed in 4.0.
Standalone deploy mode (0.6)
A Spark-only cluster manager: a master and a set of workers, started with the
scripts in sbin/, needing nothing but Java. It is still shipped and is the
simplest way to run a small multi-node cluster:
$SPARK_HOME/sbin/start-master.sh # UI on :8080, master on spark://<host>:7077
$SPARK_HOME/sbin/start-worker.sh spark://<host>:7077 # run on each worker node
spark-submit --master spark://<host>:7077 app.py
Java API (0.6)
The second language binding after Scala. Every later API, including DataFrames, kept a Java surface from here on.
Per-RDD storage levels with persist() (0.6)
Before 0.6 caching meant memory only. persist() let each RDD choose where its
cached partitions live: memory, disk, both, serialized or not, and replicated
or not:
from pyspark import StorageLevel
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("persist06").getOrCreate()
rdd = spark.sparkContext.parallelize(range(1_000_000), 4).map(lambda x: x * 2)
rdd.persist(StorageLevel.MEMORY_AND_DISK)
print(rdd.count(), rdd.getStorageLevel())
# 1000000 Disk Memory Serialized 1x Replicated
rdd.unpersist()
spark.stop()
Good to know: the output says Serialized even though the level is
MEMORY_AND_DISK. From Python, cached data is always stored as pickled bytes,
so PySpark’s storage levels are the serialized ones whatever you name.
addFile and addJar (0.6)
Ship a file or a JAR to every executor at runtime. The --files, --jars and
--py-files flags of spark-submit are the command-line form of the same
mechanism.
PySpark (0.7)
The Python API. It arrived with the architecture it still has for RDD code: the driver runs in Python and talks to a JVM, and each executor starts Python worker processes and pipes pickled rows to them. That serialization cost is the thread that runs through every Python feature in this post, up to Arrow by default in 4.2.
Spark Streaming (DStreams), alpha (0.7)
The first streaming API modelled a stream as a sequence of small RDDs, one per batch interval. It was superseded by Structured Streaming in 2.0 and is now deprecated; the 2.x section explains why.
MLlib (0.8)
The machine-learning library, which shipped with seven algorithms: SVMs,
logistic regression, linear regression variants, k-means and collaborative
filtering. The DataFrame-based pyspark.ml API that replaced the RDD-based one
still has k-means:
from pyspark.ml.clustering import KMeans
from pyspark.ml.feature import VectorAssembler
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("kmeans08").getOrCreate()
pts = spark.createDataFrame([(1.0, 1.0), (1.2, 0.8), (9.0, 9.5), (8.8, 9.1)], "x DOUBLE, y DOUBLE")
vec = VectorAssembler(inputCols=["x", "y"], outputCol="features").transform(pts)
model = KMeans(k=2, seed=1).fit(vec)
model.transform(vec).select("x", "y", "prediction").orderBy("x").show()
# +---+---+----------+
# | x| y|prediction|
# +---+---+----------+
# |1.0|1.0| 0|
# |1.2|0.8| 0|
# |8.8|9.1| 1|
# |9.0|9.5| 1|
# +---+---+----------+
spark.stop()
Web UI on port 4040 and the metrics system (0.8)
Every running application has served its UI on port 4040 since this release, with the next free port used when 4040 is taken. The metrics system sends the same numbers to sinks such as JMX, Ganglia, CSV files or Graphite; a Prometheus endpoint was added in 3.0.
YARN in mainline, and the fair scheduler (0.8)
YARN support moved from experimental into mainline, including secured
clusters. The fair scheduler let several jobs inside one application share
executors instead of queueing behind each other; turn it on with
spark.scheduler.mode=FAIR.
Good to know: 0.8 was Spark’s first release inside the Apache incubator. It became a top-level Apache project in February 2014.
SparkConf (0.9)
A configuration object for building a SparkContext. Every spark.* setting
in this post goes through it, whether it is set in code, in
spark-defaults.conf or with --conf:
from pyspark import SparkConf
from pyspark.sql import SparkSession
conf = (SparkConf().setAppName("conf09").setMaster("local[2]")
.set("spark.sql.shuffle.partitions", "8"))
spark = SparkSession.builder.config(conf=conf).getOrCreate()
print(spark.sparkContext.appName, spark.sparkContext.master, spark.conf.get("spark.sql.shuffle.partitions"))
# conf09 local[2] 8
spark.stop()
Good to know: precedence, from strongest to weakest, is values set in code,
then --conf flags, then spark-defaults.conf.
GraphX, alpha (0.9)
A graph-processing library on top of RDDs, with PageRank, connected components and triangle counting built in. It is Scala-only and has had little development since 2.x; GraphFrames, a separate project, is the DataFrame-based alternative most people use from Python.
Spark Streaming out of alpha, and Scala 2.10 (0.9)
Streaming gained driver high availability and faster windowed operators, and the build moved from Scala 2.9 to 2.10.
Broadcast variables and accumulators (0.x)
The two shared-variable primitives from this period. A broadcast variable ships a read-only value to each executor once instead of once per task, and an accumulator lets tasks add to a value that only the driver can read:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("shared-vars").getOrCreate()
sc = spark.sparkContext
city_codes = sc.broadcast({"pune": "PNQ", "hyderabad": "HYD"})
unknown = sc.accumulator(0)
def code(city):
if city not in city_codes.value:
unknown.add(1)
return city_codes.value.get(city, "???")
print(sc.parallelize(["pune", "hyderabad", "goa"]).map(code).collect())
# ['PNQ', 'HYD', '???']
print("unknown cities:", unknown.value)
# unknown cities: 1
spark.stop()
Good to know: an accumulator updated inside a transformation can count twice if a task is retried or a stage is recomputed. Only updates made inside actions are guaranteed to be applied exactly once. Treat accumulators in transformations as debugging counters, not as business metrics.
The 1.x line: structure arrives
The first line’s job was to turn a fast research engine into something an operations team would accept, and then to discover that the RDD was the wrong abstraction for most of the work people were doing with it. Read the 1.x features in two groups: 1.0 to 1.2 make Spark operable (submission, security, monitoring, a shuffle that scales), and 1.3 to 1.6 replace the RDD with the DataFrame and rebuild execution underneath it (Tungsten).
Spark 1.0, May 2014
The release that made Spark deployable rather than merely usable.
spark-submit
One submission path for local, Standalone, Mesos and YARN; before this, each cluster manager had its own launch story.
How it works. The theory behind it is the split between the driver,
which holds the SparkContext, builds the DAG of stages and schedules tasks,
and the executors, which run those tasks and hold cached data.
spark-submit is a thin launcher: it resolves the master URL, builds the
classpath from your application, --jars and --packages (resolved from Maven
through Ivy), merges configuration, and then starts the driver.
--deploy-mode decides where that driver lives:
clientkeeps the driver in the process that ran the command. Its output arrives in your terminal, which is convenient for interactive work, but the job dies when your session does, and the driver’s network traffic to every executor crosses from your machine into the cluster.clusterasks the cluster manager to start the driver on a cluster node. The job survives your disconnection and its logs live with the cluster manager, which is what you want for anything scheduled.
The same command shape has held from 1.0 to today:
spark-submit \
--master yarn \
--deploy-mode cluster \
--name nightly-orders \
--num-executors 10 \
--executor-cores 4 \
--executor-memory 8g \
--conf spark.sql.shuffle.partitions=400 \
--py-files deps.zip \
orders_job.py --run-date 2026-09-26
The master URL is the only part that changes between environments:
local[*]runs everything in one JVM, one thread per core;local[4]uses four threads.spark://host:7077targets a Standalone master.yarnreads the ResourceManager address fromHADOOP_CONF_DIR.k8s://https://api-server:6443targets a Kubernetes API server (from 2.3).
Good to know: configuration is merged in a fixed order. Values set in code on
the SparkConf or session builder win, then --conf flags, then
spark-defaults.conf. When a setting “does not take effect”, it is usually
being overridden by a stronger source; spark-submit --verbose prints every
resolved value and where it came from. For a tool that assembles these commands
from a form, see the
Spark submit formatter.
History server
The web UI disappears when an application ends. The history server rebuilds it afterwards, which is what makes post-mortem debugging of a finished or failed job possible.
How it works. With event logging on, the driver attaches a listener to its internal event bus and writes every scheduler event, one JSON object per line, to a file named after the application ID. The history server scans the log directory and, when you open an application, replays its file through the same listener code the live UI uses, so the pages look identical. You can read an event log yourself:
import json, os, tempfile
from pyspark.sql import SparkSession
logdir = tempfile.mkdtemp()
spark = (SparkSession.builder.appName("eventlog10")
.config("spark.eventLog.enabled", "true")
.config("spark.eventLog.dir", "file://" + logdir)
.config("spark.eventLog.compress", "false") # 4.x default is zstd-compressed
.getOrCreate())
spark.range(1000).selectExpr("id % 3 AS k").groupBy("k").count().collect()
spark.stop()
for root, _, files in os.walk(logdir):
for f in sorted(files):
print(os.path.relpath(os.path.join(root, f), logdir))
events = []
for root, _, files in os.walk(logdir):
for f in files:
if f.startswith("events_"):
events += [json.loads(l)["Event"] for l in open(os.path.join(root, f))]
print(len(events), "events")
for name in ["SparkListenerLogStart", "SparkListenerApplicationStart", "SparkListenerJobStart",
"SparkListenerStageCompleted", "SparkListenerTaskEnd", "SparkListenerApplicationEnd"]:
print(f"{name:32s}{events.count(name)}")
# eventlog_v2_local-1790409099524/.appstatus_local-1790409099524.crc
# eventlog_v2_local-1790409099524/appstatus_local-1790409099524
# eventlog_v2_local-1790409099524/events_1_local-1790409099524
# 32 events
# SparkListenerLogStart 1
# SparkListenerApplicationStart 1
# SparkListenerJobStart 2
# SparkListenerStageCompleted 2
# SparkListenerTaskEnd 5
# SparkListenerApplicationEnd 1
The log is not a single file on 4.x. It is a directory, eventlog_v2_<app-id>,
holding numbered events_N_ files and an appstatus marker that disappears
when the application finishes. That is the rolling format, and in 4.1.3 it is
the default, together with compression, as Spark’s own configuration entries
show:
spark.eventLog.rolling.enabled = true
spark.eventLog.compress = true
spark.eventLog.compression.codec = zstd
I switched compression off above so the file could be read as plain JSON; left
at the default, the events file is events_1_<app-id>.zstd. Every task produces
a SparkListenerTaskEnd event with its full metrics, which is why event logs of
long jobs with many tasks grow large. On a cluster the configuration is:
# in spark-defaults.conf, for every application
spark.eventLog.enabled true
spark.eventLog.dir hdfs:///spark-events
spark.eventLog.compress true
# for the history server itself
spark.history.fs.logDirectory hdfs:///spark-events
spark.history.fs.cleaner.enabled true
spark.history.fs.cleaner.maxAge 14d
$SPARK_HOME/sbin/start-history-server.sh # UI on :18080
Good to know: event logs are never cleaned up unless you enable the cleaner, so on a busy cluster the directory grows without limit. Rolling logs arrived in 3.0 for streaming jobs that run for weeks, so the history server does not have to replay one enormous file; tools that parse event logs themselves need to handle both the old single-file format and the rolling, compressed one.
Hadoop and YARN security
Kerberos authentication and delegation-token handling, the blocker for regulated Hadoop shops.
How it works. Kerberos proves identity with a ticket obtained from a keytab
or a kinit. Executors cannot each talk to Kerberos, so the driver uses its
ticket to obtain delegation tokens for HDFS, Hive and HBase and ships them
to the executors, which present the tokens instead. Tokens expire, typically
after a day, and can only be renewed up to a maximum lifetime, typically seven
days. A job that runs longer than that needs to obtain fresh tokens, which is
why long-running jobs pass a principal and keytab:
spark-submit --master yarn --deploy-mode cluster \
--principal etl@EXAMPLE.COM --keytab /etc/security/etl.keytab \
orders_job.py
Kerberos protects access to Hadoop services. Spark’s own traffic between driver and executors is protected separately:
--conf spark.authenticate=true # shared-secret auth for internal RPC
--conf spark.network.crypto.enabled=true # encrypt RPC and block transfers
--conf spark.io.encryption.enabled=true # encrypt shuffle and spill files on local disk
Good to know: a streaming job that fails after exactly seven days with an
HDFS token ... can't be found in cache error has hit the token maximum
lifetime. The keytab is the fix, not a restart schedule.
Spark SQL, alpha, with the Catalyst optimizer
Spark SQL arrived with SchemaRDD, not DataFrame; the name survived exactly
three releases. Catalyst, introduced in the same release, is still the
optimizer in 4.2.
How it works. Catalyst represents a query as a tree and transforms it through four phases, each a set of rules that pattern-match on the tree and rewrite it:
- Parsing turns SQL text into an unresolved logical plan, in which table and column names are just strings.
- Analysis resolves those names against the catalog, assigns each column a unique ID and a type, and fails if anything does not exist.
- Optimization applies rule batches repeatedly until the plan stops changing: constant folding, predicate pushdown, column pruning, filter combining and dozens more.
- Physical planning picks an operator for each logical step, such as which join algorithm, and code generation compiles the result.
Every phase is visible from Python:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("catalyst10").getOrCreate()
spark.range(0, 100).selectExpr("id", "id % 5 AS city_id").createOrReplaceTempView("orders")
df = spark.sql("SELECT id * 2 AS doubled FROM orders WHERE city_id = 3 AND 1 = 1")
qe = df._jdf.queryExecution()
for label, plan in [("parsed", qe.logical()), ("analyzed", qe.analyzed()),
("optimized", qe.optimizedPlan()), ("physical", qe.executedPlan())]:
print(f"== {label} ==")
print(plan.toString().strip())
# == parsed ==
# 'Project [('id * 2) AS doubled#9]
# +- 'Filter (('city_id = 3) AND (1 = 1))
# +- 'UnresolvedRelation [orders], [], false
# == analyzed ==
# Project [(id#7L * cast(2 as bigint)) AS doubled#9L]
# +- Filter ((city_id#8L = cast(3 as bigint)) AND (1 = 1))
# +- SubqueryAlias orders
# +- View (`orders`, [id#7L, city_id#8L])
# +- Project [id#7L, (id#7L % cast(5 as bigint)) AS city_id#8L]
# +- Range (0, 100, step=1, splits=Some(4))
# == optimized ==
# Project [(id#7L * 2) AS doubled#9L]
# +- Filter ((id#7L % 5) = 3)
# +- Range (0, 100, step=1, splits=Some(4))
# == physical ==
# *(1) Project [(id#7L * 2) AS doubled#9L]
# +- *(1) Filter ((id#7L % 5) = 3)
# +- *(1) Range (0, 100, step=1, splits=4)
spark.stop()
Reading the four plans top to bottom shows each phase’s job. The parsed plan
has 'UnresolvedRelation [orders] and quoted names: nothing is known yet. The
analyzed plan has found the view, given every column an ID (id#7L), and
inserted cast(2 as bigint) because id is a bigint. The optimized plan has
folded the casts into literals, deleted 1 = 1, inlined the view, and
substituted city_id with its definition, so the filter reads
(id#7L % 5) = 3 directly against the Range. The physical plan picks
operators, and the *(1) markers show that all three were fused into one
generated method.
The pushdown Catalyst was built for is visible in any Parquet scan:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("sql10").getOrCreate()
path = tempfile.mkdtemp()
spark.range(0, 1000).selectExpr("id", "id % 5 AS city_id").write.mode("overwrite").parquet(path)
spark.read.parquet(path).createOrReplaceTempView("orders")
df = spark.sql("SELECT id FROM orders WHERE city_id = 3 AND id > 990")
print(df.collect())
scan = [l for l in df._jdf.queryExecution().executedPlan().toString().splitlines() if "FileScan" in l][0]
print(scan[scan.index("PushedFilters"):scan.index("ReadSchema")].strip().rstrip(","))
# [Row(id=993), Row(id=998)]
# PushedFilters: [IsNotNull(city_id), IsNotNull(id), EqualTo(city_id,3), GreaterThan(id,990)]
spark.stop()
Both conditions were handed to the Parquet reader, which uses each row group’s min and max statistics to skip groups that cannot match.
Good to know: a pushed filter is a hint to the source, not a guarantee. Spark still re-applies every filter after the scan, so a source that can only use some filters stays correct. The optimizer itself is covered in depth in Inside Spark’s Catalyst optimizer.
API stability for the 1.x line
A promise that code written against 1.0 would keep compiling and running on every 1.x release.
How it works. Public APIs are stable unless marked otherwise. Classes and
methods annotated @Experimental or @DeveloperApi may change in any release,
and the build runs the MiMa tool on every change to catch accidental binary
incompatibilities in the stable surface.
Good to know: the guarantee is per major line. Each major release, 2.0, 3.0 and 4.0, is where Spark spends its compatibility budget, which is why upgrades across a major version need a migration pass and upgrades within one usually do not.
Spark 1.1, September 2014
Sort-based shuffle, available
A new shuffle writer, opt-in here and the default from 1.2.
How it works. A shuffle moves data so that all records with the same key end up in the same partition. The map side of a shuffle writes each task’s output grouped by destination partition; the reduce side fetches its slice from every map task. The sort writer buffers records in memory tagged with their destination partition ID, sorts the buffer by that ID (and by key, when there is a map-side combine), spills sorted runs to disk when memory runs out, and finally merges everything into one data file plus one index file per map task. The index holds the byte offset where each reduce partition’s data starts, so a reducer fetches exactly its byte range.
You can see the files it leaves on disk. Four map tasks, fifty reduce partitions:
import os, tempfile
from pyspark.sql import SparkSession
local = tempfile.mkdtemp()
spark = (SparkSession.builder.appName("sortshuffle11")
.config("spark.local.dir", local)
.config("spark.sql.adaptive.enabled", "false")
.config("spark.sql.shuffle.partitions", "50")
.getOrCreate())
spark.range(0, 100_000, numPartitions=4).selectExpr("id % 1000 AS k").groupBy("k").count().collect()
files = sorted(f for _, _, fs in os.walk(local) for f in fs if f.startswith("shuffle_"))
print(len(files), "shuffle files for 4 map tasks x 50 reduce partitions")
print(files)
# 12 shuffle files for 4 map tasks x 50 reduce partitions
# ['shuffle_0_0_0.checksum.ADLER32', 'shuffle_0_0_0.data', 'shuffle_0_0_0.index',
# 'shuffle_0_1_0.checksum.ADLER32', 'shuffle_0_1_0.data', 'shuffle_0_1_0.index',
# 'shuffle_0_2_0.checksum.ADLER32', 'shuffle_0_2_0.data', 'shuffle_0_2_0.index',
# 'shuffle_0_3_0.checksum.ADLER32', 'shuffle_0_3_0.data', 'shuffle_0_3_0.index']
spark.stop()
Twelve files: for each of the four map tasks, one .data file, one .index
file and one checksum file, named shuffle_<shuffleId>_<mapId>_0. There are
fifty reduce partitions, but no per-partition files: each reducer reads its byte
range out of every map task’s data file. The checksum file is newer than the
sort writer; since 3.2 Spark records a checksum per partition, so that when a
reducer reads a corrupt block it can tell whether the disk or the network is to
blame.
Good to know: when there are few reduce partitions (at most
spark.shuffle.sort.bypassMergeThreshold, default 200) and no map-side
aggregation, Spark skips the sort and writes one temporary file per partition,
then concatenates them. Either way the result on disk is the same data-plus-index
pair. The shuffle files are what the external shuffle service, from 1.2, serves
after an executor is gone.
JDBC/ODBC server
The Thrift server lets BI tools and anything with a JDBC or ODBC driver run SQL against Spark:
$SPARK_HOME/sbin/start-thriftserver.sh --master yarn \
--conf spark.scheduler.mode=FAIR
beeline -u jdbc:hive2://localhost:10000 -e "SHOW TABLES"
How it works. It speaks the HiveServer2 protocol, so any Hive client can
connect. Behind it is one Spark application: every connection gets its own
session, with its own SQL settings, temporary views and current database, but
all sessions share the same SparkContext, executors and cache.
Good to know: because everyone shares one application, one heavy query can
starve the rest. Run it with spark.scheduler.mode=FAIR, as above, and assign
sessions to pools so that interactive users are not queued behind a batch
report. Spark Connect, from 3.4, is the newer alternative for programmatic
clients.
JSON source with schema inference
Point Spark at JSON lines and it scans the data to infer a schema, including nested structs and arrays:
import os, tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("json11").getOrCreate()
d = tempfile.mkdtemp()
with open(os.path.join(d, "e.json"), "w") as f:
f.write('{"id": 1, "user": {"name": "ravi", "tags": ["a", "b"]}}\n')
f.write('{"id": 2, "user": {"name": "sita"}, "vip": true}\n')
df = spark.read.json(d)
df.printSchema()
df.select("id", "user.name", "vip").orderBy("id").show()
# root
# |-- id: long (nullable = true)
# |-- user: struct (nullable = true)
# | |-- name: string (nullable = true)
# | |-- tags: array (nullable = true)
# | | |-- element: string (containsNull = true)
# |-- vip: boolean (nullable = true)
# +---+----+----+
# | id|name| vip|
# +---+----+----+
# | 1|ravi|NULL|
# | 2|sita|true|
# +---+----+----+
spark.stop()
How it works. Spark infers a type for every record, then merges them pairwise:
a field present in only some records becomes nullable (vip above), numeric
types widen (long and double merge to double), and types that cannot be
reconciled fall back to string. Fields are sorted alphabetically in the
result, which is why the schema does not follow the order in the file.
Good to know: inference reads the whole input, which is slow on large data.
In production, pass an explicit schema with .schema(...), or set
samplingRatio to infer from a fraction of the records. The default expects one
JSON object per line; a file holding a single pretty-printed document needs
multiLine=true, available from 2.2.
Dynamic bytecode generation for expressions
How it works. Without code generation, evaluating a + b * 2 means walking an
expression tree: each node is an object whose eval method calls its children’s
eval, boxing every intermediate value into a Java object. Code generation turns
the tree into Java source for one specialised function, compiles it at runtime
with the Janino compiler, and caches the compiled class, so each row costs a few
primitive operations instead of a chain of virtual calls.
Good to know: this release compiled single expressions. 1.5 turned it on for almost every function, and 2.0’s whole-stage codegen compiled whole chains of operators. If generated code is too large to compile, Spark silently falls back to interpreted evaluation, which is one reason a query with hundreds of columns can be unexpectedly slow.
UDF registration from Python, Scala and Java
Registering a Python function makes it callable from SQL:
from pyspark.sql import SparkSession
from pyspark.sql.types import StringType
spark = SparkSession.builder.appName("udf11").getOrCreate()
spark.udf.register("initcap_city", lambda s: s.title(), StringType())
spark.sql("SELECT initcap_city('new delhi') AS city").show()
# +---------+
# | city|
# +---------+
# |New Delhi|
# +---------+
spark.stop()
It works, and it is the slowest way to do this.
How it works. A Python UDF cannot run inside the JVM. Spark inserts a separate operator into the physical plan that sends input rows to a Python worker process on the same executor, pickled in batches, waits for the results and joins them back to the rows. The plan shows the difference directly:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StringType
spark = SparkSession.builder.appName("udfplan11").getOrCreate()
title = F.udf(lambda s: s.title(), StringType())
df = spark.createDataFrame([("new delhi",)], "city STRING")
for label, q in [("python udf", df.select(title("city"))), ("built-in", df.select(F.initcap("city")))]:
plan = q._jdf.queryExecution().executedPlan().toString()
print(f"{label:11s}", [n for n in ["BatchEvalPython", "ArrowEvalPython", "Project"] if n in plan])
# python udf ['BatchEvalPython', 'Project']
# built-in ['Project']
spark.stop()
BatchEvalPython is the round trip to Python. The built-in initcap is a plain
Project that runs inside generated JVM code. The UDF is also opaque to the
optimizer: Catalyst cannot push a filter through it, reorder it or know what it
returns for a null.
Good to know: the rule that follows from this release has never changed:
reach for a built-in function first, a vectorized UDF second, and a
row-at-a-time UDF last. The rest of this post follows the project’s attempts to
shrink that last option’s cost, from pandas UDFs in 2.3 to Arrow by default in
4.2, where the operator becomes ArrowEvalPython.
Disk spilling for large cached blocks
How it works. Caching a partition means materialising all of its records. Before this change, a partition too large for memory could exhaust the heap while it was being built. Spark now “unrolls” a partition incrementally, checking memory as it goes; if the partition will not fit and the storage level allows disk, it is written to disk instead of failing.
Good to know: with MEMORY_ONLY, a partition that does not fit is simply not
cached, and is recomputed every time it is used. The Storage tab’s “Fraction
Cached” column is where that shows up.
New defaults: Snappy compression and torrent broadcast
spark.io.compression.codec became snappy, though it has moved on since: the
default in the 4.1.3 build I ran is lz4. The codec applies to shuffle files,
spills and broadcast data, and trades a little CPU for much less disk and
network I/O.
How it works (torrent broadcast). spark.broadcast.factory became
TorrentBroadcastFactory. The driver splits a broadcast value into chunks of
spark.broadcast.blockSize (4 MB) and each executor fetches the chunks it lacks
from the driver or from other executors that already have them, then serves
them in turn, like BitTorrent. The driver no longer has to send the whole value
to every executor itself.
Good to know: this is also how broadcast joins distribute the small side, so the same mechanism, and the same limits on driver memory, apply to both. The driver still has to collect the whole value first.
Spark 1.2, December 2014
The release that changed the most defaults in the shortest note.
Sort shuffle becomes the default
spark.shuffle.manager changed from hash to sort, the single most
consequential default flip in Spark’s history.
How it works. Hash shuffle had every map task open one file per reduce partition. With M map tasks and R reduce partitions that is M × R files: 2,000 by 2,000 is four million files, each needing a file handle and a write buffer, which exhausted file descriptors and memory on large jobs. Sort shuffle, as the 1.1 entry shows, writes two files per map task, 2 × M in total, whatever R is.
Good to know: the setting is gone from the 4.x configuration reference; sort shuffle is simply how Spark shuffles now. Everything about shuffle tuning after 1.2 assumes the sort writer, which is described in Apache Spark architecture.
Netty block transfer by default
spark.shuffle.blockTransferService changed from nio to netty.
How it works. Reducers fetch shuffle blocks over the network from the
executors, or shuffle services, that wrote them. The Netty implementation sends
file regions with zero-copy sendfile, so block data goes from the page cache to
the socket without being copied through the JVM heap, and it keeps its buffers
off-heap.
Good to know: the setting no longer exists; Netty is the only implementation.
The knobs that matter now are the fetch retries. A FetchFailedException means
a reducer could not get a block, usually because the executor holding it died;
spark.shuffle.io.maxRetries (default 3) and spark.shuffle.io.retryWait
(default 5s) decide how hard it tries before Spark reruns the map stage.
Dynamic allocation on YARN
Lets an application give executors back to the cluster when it is idle and ask for more when tasks queue up.
How it works. When tasks have been waiting for longer than
spark.dynamicAllocation.schedulerBacklogTimeout (1 second), Spark requests
more executors, and it keeps requesting every
sustainedSchedulerBacklogTimeout while the backlog lasts, doubling each round:
1, then 2, then 4, then 8, up to maxExecutors. An executor with no running
tasks for executorIdleTimeout (60 seconds) is released.
The catch that explains its configuration is shuffle data. An executor that wrote shuffle files cannot simply disappear, because a later stage still needs to read them. On YARN the answer was the external shuffle service, a long-running process on each node that serves shuffle files on behalf of executors that no longer exist:
spark-submit \
--master yarn \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.dynamicAllocation.minExecutors=2 \
--conf spark.dynamicAllocation.maxExecutors=50 \
--conf spark.dynamicAllocation.executorIdleTimeout=60s \
--conf spark.shuffle.service.enabled=true \
orders_job.py
Good to know: executors holding cached data are never released by default,
because spark.dynamicAllocation.cachedExecutorIdleTimeout is infinite, so a
cached DataFrame can pin executors for the application’s whole life. On
Kubernetes, where there is no node-level shuffle service by default, 3.0 added
spark.dynamicAllocation.shuffleTracking.enabled, which keeps an executor alive
while it holds shuffle data anyone still needs.
External data source API
A plug-in interface for reading and writing data that Spark does not ship a
reader for, and the ancestor of every connector you use. It is where the uniform
format(...).option(...).load() shape comes from, the same for every source:
import os, tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ds12").getOrCreate()
d = tempfile.mkdtemp()
with open(os.path.join(d, "c.csv"), "w") as f:
f.write("city,amount\npune,10\ngoa,5\n")
df = spark.read.format("csv").option("header", "true").option("inferSchema", "true").load(d)
df.printSchema()
out = tempfile.mkdtemp() + "/out"
df.write.format("parquet").mode("overwrite").save(out)
print(spark.read.format("parquet").load(out).count())
# root
# |-- city: string (nullable = true)
# |-- amount: integer (nullable = true)
# 2
spark.stop()
How it works. A source implements a small set of Scala traits. A
RelationProvider turns the options into a BaseRelation that reports a schema,
and the relation chooses how much work it can take off Spark’s hands:
TableScan returns every row, PrunedScan receives the list of columns the
query needs, and PrunedFilteredScan also receives the filters. This is the
mechanism behind the pushdown shown in the 1.0 Catalyst entry.
Good to know: this first API, now called V1, was replaced by Data Source V2, experimental in 2.3 and the basis of the 3.0 catalog API. Iceberg, Delta and Hudi are V2 sources. From 4.0 you can write a source in Python.
Streaming write-ahead log, and a Python streaming API
How it works. A receiver-based DStream holds received data in executor memory
until a batch processes it. If the driver fails, that data is lost. With
spark.streaming.receiver.writeAheadLog.enable=true, a receiver writes each
block to a log in the checkpoint directory, on HDFS or another fault-tolerant
file system, before acknowledging it to the source. After a failure the new
driver replays the log.
Good to know: the log costs throughput, since every record is written twice, and it only gives at-least-once delivery. 1.3’s direct Kafka stream and then Structured Streaming removed the need for it by treating the source’s own offsets as the log. DStreams also became usable from Python in this release.
GraphX graduates from alpha
GraphX’s API was declared stable.
How it works. GraphX represents a graph as two RDDs, one of vertices and one of edges, each carrying arbitrary properties, and exposes a Pregel-style API in which vertices exchange messages along edges in rounds until they stop changing. PageRank and connected components are written that way. It is Scala-only:
import org.apache.spark.graphx.{Edge, Graph}
val edges = sc.parallelize(Seq(Edge(1L, 2L, 1), Edge(2L, 3L, 1), Edge(3L, 1L, 1), Edge(4L, 1L, 1)))
val graph = Graph.fromEdges(edges, defaultValue = 0)
graph.pageRank(0.0001).vertices.collect().sortBy(_._1)
.foreach { case (id, rank) => println(f"$id%d $rank%.3f") }
// 1 1.330
// 2 1.281
// 3 1.239
// 4 0.150
Good to know: GraphX has had little development since 2.x. GraphFrames, a separate project with a DataFrame API and Python support, is what most people use now.
The broadcast-join threshold, still unchanged
This release set spark.sql.autoBroadcastJoinThreshold to 10485760, raising
it from 10000. That value has not changed since. On 4.1.3:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("defaults").getOrCreate()
print(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
# 10485760b
print(spark.conf.get("spark.sql.shuffle.partitions"))
# 200
spark.stop()
How it works. At planning time Spark compares each join side’s estimated
size with the threshold and broadcasts a side below it. For a file-based table
the estimate is the size of its files on disk; for a table with ANALYZE TABLE
statistics it is the recorded size; and for a relation Spark knows nothing about
it is effectively infinite, so that side is never broadcast automatically.
Good to know: the estimate is on-disk size, and compressed columnar files can grow ten times or more when decoded into memory, so a “9 MB” table can be a 100 MB broadcast. Eleven years and three major versions later, the threshold a 2014 release chose is still deciding whether your join shuffles; since 3.0, AQE can also switch a join to broadcast at runtime once it has measured the real size.
Spark 1.3, March 2015
The release where Spark stopped being an RDD engine.
The DataFrame API
Named, typed columns in Python, Scala and Java. This is the oldest API in the post that still runs unchanged today:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("df13").getOrCreate()
orders = spark.createDataFrame(
[(1, "pune", 1200), (2, "hyderabad", 800), (3, "pune", 450)],
"order_id INT, city STRING, amount INT",
)
orders.groupBy("city").sum("amount").orderBy("city").show()
# +---------+-----------+
# | city|sum(amount)|
# +---------+-----------+
# |hyderabad| 800|
# | pune| 1650|
# +---------+-----------+
spark.stop()
How it works. A DataFrame is not data. It is an immutable, lazily built
logical plan: every transformation (select, where, groupBy) returns a new
DataFrame whose plan has one more node, and nothing runs until an action
(show, collect, count, write) hands the whole plan to Catalyst. Because
the engine sees the entire query before running it, it can do things no RDD
program gets for free. Here the query mentions four columns, but only asks for
one after filtering on another:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("pruning13").getOrCreate()
path = tempfile.mkdtemp()
spark.range(0, 1000).selectExpr("id AS order_id", "id % 7 AS city_id", "id * 10 AS amount",
"repeat('x', 200) AS notes").write.mode("overwrite").parquet(path)
df = spark.read.parquet(path).where("city_id = 3").select("amount")
df.collect()
scan = [l for l in df._jdf.queryExecution().executedPlan().toString().splitlines() if "FileScan" in l][0]
print(scan[scan.index("PushedFilters"):].strip())
# PushedFilters: [IsNotNull(city_id), EqualTo(city_id,3)], ReadSchema: struct<city_id:bigint,amount:bigint>
spark.stop()
ReadSchema lists only city_id and amount. The wide notes column and
order_id are never read from disk, and the filter was pushed into the scan.
The same logic written as RDD code reads every byte of every row.
Good to know: converting to an RDD with df.rdd ends the optimization at that
point; everything after it runs as opaque functions again. Stay in the
DataFrame API for as long as the logic allows.
SchemaRDD renamed to DataFrame, and Spark SQL out of alpha
The rename looks cosmetic and was not. A SchemaRDD was an RDD that happened to
carry a schema, and it inherited RDD semantics: its operations were RDD
operations. A DataFrame is a description of a computation over named, typed
columns, which the engine is free to execute however it likes. That freedom is
what every later optimisation spends: codegen, columnar scans, adaptive
execution.
JDBC data source
Read from and write to MySQL, Postgres, Oracle and anything else with a JDBC driver. Derby ships inside the Spark distribution, so this example needs no external database:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("jdbc13").getOrCreate()
url = "jdbc:derby:memory:shop;create=true"
driver = "org.apache.derby.jdbc.EmbeddedDriver"
(spark.createDataFrame([(1, "pune", 1200), (2, "goa", 300), (3, "pune", 450)],
"order_id INT, city STRING, amount INT")
.write.format("jdbc")
.option("url", url).option("driver", driver).option("dbtable", "orders")
.option("createTableColumnTypes", "city VARCHAR(32)") # STRING would become a CLOB
.mode("overwrite").save())
df = (spark.read.format("jdbc")
.option("url", url).option("driver", driver)
.option("query", 'SELECT "city", SUM("amount") AS total FROM orders GROUP BY "city"')
.load())
df.orderBy("city").show()
# +----+-----+
# |city|TOTAL|
# +----+-----+
# | goa| 300|
# |pune| 1650|
# +----+-----+
spark.stop()
How it works. When you read with dbtable, Spark first runs the query with
WHERE 1=0 to learn the schema, then issues one SELECT per partition, with the
query’s filters translated into a SQL WHERE clause so the database does the
filtering. On write, each partition opens its own connection and inserts in
batches of batchsize rows (default 1,000).
A parallel read splits a numeric, date or timestamp column into ranges:
df = (spark.read.format("jdbc")
.option("url", url).option("dbtable", "orders")
.option("partitionColumn", "order_id")
.option("lowerBound", 1).option("upperBound", 1_000_000)
.option("numPartitions", 8) # eight queries, eight connections
.option("fetchsize", 10_000) # rows per round trip
.load())
Good to know: three things this small example ran into, plus one that bites in production.
- A JDBC read is one partition and one connection by default, however big
the table. Use the partitioning options above, or
queryto push a whole aggregation into the database, as the first example does. lowerBoundandupperBounddo not filter. They only decide the stride of the ranges; rows outside them all land in the first or last partition.- Spark’s type mapping may not suit the database. Written without
createTableColumnTypes, theSTRINGcolumn became a DerbyCLOB, and theGROUP BYfailed withColumns of type 'CLOB' may not be used in ... GROUP BY. - Spark quotes the column names it creates, so they keep their lowercase
spelling. In a case-folding database such as Derby, Oracle or DB2, a
pushed-down query then has to quote them too: unquoted
cityfailed withColumn 'CITY' is either not in any table in the FROM list.
Parquet schema merging
Files written at different times with compatible but different schemas can be read as one table whose schema is the union of theirs:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("merge13").getOrCreate()
path = tempfile.mkdtemp()
spark.createDataFrame([(1, "pune")], "id INT, city STRING").write.parquet(path + "/day=1")
spark.createDataFrame([(2, "goa", 99.5)], "id INT, city STRING, amount DOUBLE").write.parquet(path + "/day=2")
print(spark.read.parquet(path).columns)
print(spark.read.option("mergeSchema", "true").parquet(path).columns)
spark.read.option("mergeSchema", "true").parquet(path).orderBy("id").show()
# ['id', 'city', 'day']
# ['id', 'city', 'amount', 'day']
# +---+----+------+---+
# | id|city|amount|day|
# +---+----+------+---+
# | 1|pune| NULL| 1|
# | 2| goa| 99.5| 2|
# +---+----+------+---+
spark.stop()
How it works. Every Parquet file stores its schema in its footer. With merging
on, Spark reads the footers of all files and unions the fields; columns missing
from a file read as NULL for its rows. The day column comes from the
directory names, which Spark’s partition discovery turns into a column.
Good to know: merging has been off by default since 1.5, because reading
every footer is slow on large tables. Without it Spark takes the schema from one
file, and columns that only exist in other files silently disappear, which is
what the first print shows. Merging also fails if two files disagree on a
column’s type. Table formats such as Iceberg and Delta solve this properly by
keeping the schema in table metadata.
Direct Kafka stream
How it works. The receiver-based approach ran a long-lived receiver task that consumed from Kafka, stored blocks in executor memory, and wrote them to a write-ahead log for durability. The direct approach has no receivers. At the start of each batch the driver asks Kafka for the latest offsets and defines the batch as an offset range per Kafka partition. Each range becomes one Spark partition, read directly from the broker by a normal task. The ranges, not the data, are what get checkpointed, so after a failure Spark simply re-reads the same ranges.
Good to know: that makes the reading exactly-once. End-to-end exactly-once also needs a sink that is idempotent, or that stores the offsets in the same transaction as the output. This is the same reasoning Structured Streaming later generalises to every source; its Kafka source, from 2.1, is the one to use now.
Spark 1.4, June 2015
SparkR
An R binding built on the DataFrame API:
library(SparkR)
sparkR.session()
df <- createDataFrame(faithful)
head(summarize(groupBy(df, df$waiting), count = n(df$waiting)))
How it works. Like PySpark, SparkR runs your R code in an R process that
drives a JVM through a socket; DataFrame operations are translated into JVM
calls, so they run at the same speed as from Scala, while R functions passed to
dapply or gapply run in R worker processes on the executors.
Good to know: SparkR was deprecated in 4.0 and is still present in 4.2. The
community sparklyr package is the more common way to use Spark from R.
Window functions (SPARK-1442)
A window function computes a value for each row from a set of related rows, the
frame, without collapsing them the way groupBy does.
How it works. A window specification has three parts. PARTITION BY picks
the related rows, ORDER BY orders them within each partition, and the frame
decides which of the ordered rows feed the calculation for the current row. When
an ordering is given, the default frame is “start of partition to current row”,
which is why a sum over an ordered window is a running total. Physically,
Spark shuffles by the partition keys, sorts each partition by the order keys,
and streams through the sorted rows, keeping only the frame in memory:
from pyspark.sql import SparkSession, Window
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("window14").getOrCreate()
sales = spark.createDataFrame(
[("pune", "2026-01-01", 100), ("pune", "2026-01-02", 300),
("pune", "2026-01-03", 200), ("goa", "2026-01-01", 50),
("goa", "2026-01-02", 70)],
"city STRING, day STRING, amount INT",
)
w = Window.partitionBy("city").orderBy("day")
(sales.withColumn("running_total", F.sum("amount").over(w))
.withColumn("rank_in_city", F.rank().over(Window.partitionBy("city").orderBy(F.desc("amount"))))
.withColumn("prev_amount", F.lag("amount").over(w))
.orderBy("city", "day")
.show())
# +----+----------+------+-------------+------------+-----------+
# |city| day|amount|running_total|rank_in_city|prev_amount|
# +----+----------+------+-------------+------------+-----------+
# | goa|2026-01-01| 50| 50| 2| NULL|
# | goa|2026-01-02| 70| 120| 1| 50|
# |pune|2026-01-01| 100| 100| 3| NULL|
# |pune|2026-01-02| 300| 400| 1| 100|
# |pune|2026-01-03| 200| 600| 2| 300|
# +----+----------+------+-------------+------------+-----------+
spark.stop()
The frame can be set explicitly. rowsBetween(-2, 0) is a three-row moving
average; an empty partitionBy() makes the whole table one window, which is how
you compute a share of the total:
from pyspark.sql import SparkSession, Window
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("frames14").getOrCreate()
daily = spark.createDataFrame([(d, a) for d, a in enumerate([10, 20, 30, 40, 50], start=1)], "day INT, amount INT")
by_day = Window.orderBy("day")
(daily.withColumn("moving_avg_3", F.avg("amount").over(by_day.rowsBetween(-2, 0)))
.withColumn("pct_of_total", F.round(F.col("amount") / F.sum("amount").over(Window.partitionBy()), 3))
.withColumn("next_amount", F.lead("amount").over(by_day))
.withColumn("quartile", F.ntile(4).over(by_day))
.orderBy("day").show())
# +---+------+------------+------------+-----------+--------+
# |day|amount|moving_avg_3|pct_of_total|next_amount|quartile|
# +---+------+------------+------------+-----------+--------+
# | 1| 10| 10.0| 0.067| 20| 1|
# | 2| 20| 15.0| 0.133| 30| 1|
# | 3| 30| 20.0| 0.2| 40| 2|
# | 4| 40| 30.0| 0.267| 50| 3|
# | 5| 50| 40.0| 0.333| NULL| 4|
# +---+------+------------+------------+-----------+--------+
spark.stop()
Good to know: rowsBetween counts physical rows; rangeBetween counts values
of the ordering column, so rangeBetween(-6, 0) over a day number means “the
last seven days” even when some days have no rows. Each distinct
PARTITION BY / ORDER BY pair costs its own shuffle and sort. And a window
with no PARTITION BY, like both windows in the second example, moves every row
to a single task; Spark logs a warning, and on real data it is a bottleneck.
Sort-merge join (SPARK-2213)
Until it existed, a shuffle join that did not fit in memory failed. After it, join size is bounded by disk instead of memory.
How it works. Three steps. Both sides are shuffled by the join key, so matching keys land in the same partition number on both sides. Each partition is sorted by the key, spilling to disk if needed. Then the two sorted streams are merged in one pass, like merging two sorted lists, holding only the rows for the current key in memory. The physical plan shows all three:
from pyspark.sql import SparkSession
spark = (SparkSession.builder.appName("smj14")
.config("spark.sql.autoBroadcastJoinThreshold", "-1")
.config("spark.sql.adaptive.enabled", "false").getOrCreate())
a = spark.range(0, 100_000).selectExpr("id AS k", "id AS v")
b = spark.range(0, 100_000).selectExpr("id AS k", "id * 2 AS w")
plan = a.join(b, "k")._jdf.queryExecution().executedPlan().toString()
print("\n".join(l.split(", [plan_id")[0] for l in plan.splitlines()))
# *(5) Project [k#70L, v#71L, w#74L]
# +- *(5) SortMergeJoin [k#70L], [k#73L], Inner
# :- *(2) Sort [k#70L ASC NULLS FIRST], false, 0
# : +- Exchange hashpartitioning(k#70L, 200), ENSURE_REQUIREMENTS
# : +- *(1) Project [id#69L AS k#70L, id#69L AS v#71L]
# : +- *(1) Range (0, 100000, step=1, splits=4)
# +- *(4) Sort [k#73L ASC NULLS FIRST], false, 0
# +- Exchange hashpartitioning(k#73L, 200), ENSURE_REQUIREMENTS
# +- *(3) Project [id#72L AS k#73L, (id#72L * 2) AS w#74L]
# +- *(3) Range (0, 100000, step=1, splits=4)
spark.stop()
The trade-off is explicit: sort-merge join pays a shuffle and a sort on both sides to remove the memory ceiling.
Good to know: the expensive part is the shuffle, not the merge, which is why avoiding it matters so much: a broadcast join (1.5) avoids shuffling the large side, and two tables bucketed on the join key (2.0) avoid both shuffles. The full decision procedure is in Spark joins in depth.
ORC support (SPARK-2883)
Read and write ORC, the columnar format from the Hive world:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("orc14").getOrCreate()
path = tempfile.mkdtemp() + "/orders_orc"
spark.range(0, 100).selectExpr("id", "id % 3 AS city_id").write.orc(path)
df = spark.read.orc(path).where("city_id = 1")
print(df.count(), spark.conf.get("spark.sql.orc.impl"))
# 33 native
spark.stop()
How it works. Like Parquet, ORC stores each column separately in stripes, with min/max statistics per stripe and per 10,000-row group, so filters can skip data and queries read only the columns they need.
Good to know: this first ORC reader went through Hive’s libraries. 2.3 added a
native vectorized reader and 2.4 made it the default, which is what the native
in the output is. Outside Hive-centric shops, Parquet is the more common choice
and gets new Spark features first.
Project Tungsten begins (SPARK-7081)
The first performance work under the Tungsten name.
How it works. Tungsten had three goals: manage memory explicitly instead of leaving it to the JVM’s objects and garbage collector, lay data out so the CPU cache is used well, and generate code instead of interpreting. The first delivery, in this release, was a faster shuffle write path. Instead of sorting Java objects, it serialises records into binary pages and sorts an array of 8-byte entries, each holding a record’s partition ID and a pointer to it. Sorting compact integers is far faster and more cache-friendly than sorting objects.
Good to know: the 1.5 section is where Tungsten reshaped execution as a whole.
DAG visualisation in the UI (SPARK-6942)
How it works. Each job and stage page draws its DAG: boxes for stages, dots for the RDD operations inside them, and edges where data flows. A stage boundary is always a shuffle, so the number of boxes is the number of shuffles plus one. Cached RDDs are drawn in green. The SQL tab later added the same idea for DataFrame queries, with each whole-stage-codegen cluster and its metrics.
Good to know: it is the fastest way to find an unexpected shuffle. A job you expected to be one stage that draws as three has two shuffles you did not plan for.
Python 3 support (SPARK-4897)
PySpark ran on Python 3 for the first time. Python 2 support was removed in 3.1.
Good to know: the driver and the executors must run the same Python minor
version, or tasks fail with Python in worker has different version. Choose the
interpreter explicitly with PYSPARK_PYTHON (executors) and
PYSPARK_DRIVER_PYTHON (driver), or spark.pyspark.python.
REST API for application information (SPARK-3644)
Everything the UI shows is also JSON under /api/v1, on a running
application’s UI port or on the history server. It is how monitoring tools
scrape Spark:
import json, urllib.request
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("rest14").config("spark.ui.enabled", "true").getOrCreate()
spark.range(100).count()
base = spark.sparkContext.uiWebUrl + "/api/v1"
apps = json.load(urllib.request.urlopen(base + "/applications"))
app_id = apps[0]["id"]
jobs = json.load(urllib.request.urlopen(f"{base}/applications/{app_id}/jobs"))
print(apps[0]["name"], [(j["jobId"], j["status"]) for j in jobs])
# rest14 [(1, 'SUCCEEDED'), (0, 'SUCCEEDED')]
spark.stop()
Under /applications/<app-id>/ the useful endpoints are jobs, stages
(with per-task metrics under stages/<id>/<attempt>/taskList), executors
(memory, GC time, failed tasks), storage/rdd (what is cached) and
environment (every resolved setting).
Good to know: one count() produced two jobs here, a small reminder that jobs
in the UI do not map one-to-one to the actions in your code. The same paths work
against the history server on port 18080, for applications that have already
finished.
Spark 1.5, September 2015
Tungsten’s main release. The theme is that the JVM’s object model is the bottleneck, so Spark should stop using it for data.
Code generation on by default
Almost all DataFrame and SQL functions were compiled to bytecode instead of interpreted.
How it works. See the 1.1 entry on dynamic bytecode generation: the expression tree becomes one generated function per projection or filter, with primitive types instead of boxed objects. 1.5 extended it from a handful of expressions to nearly all of them, so a typical query ran almost entirely in generated code.
Good to know: 2.0’s whole-stage codegen extends this from single expressions to whole operator chains, which removes the remaining per-row cost of passing rows between operators.
Cache-friendly hash maps for aggregation
How it works. A groupBy builds a hash map from each group key to its running
aggregates. Before Tungsten that map held Java objects, scattered across the
heap, so each lookup chased pointers and missed the CPU cache. The new map
(BytesToBytesMap) stores binary keys and values in large contiguous memory
pages and uses open addressing, with part of each key’s hash stored alongside the
pointer, so most lookups are resolved without touching the key itself.
Good to know: you see this as HashAggregate in plans. The alternatives,
ObjectHashAggregate and SortAggregate, appear when an aggregate function
cannot work on binary buffers, such as collect_list or some UDAFs, and they
are slower.
External sort-based aggregation
How it works. When the aggregation hash map cannot grow any further, Spark sorts its current contents by key, spills them to disk, and starts a fresh map. At the end it merges all the sorted runs and combines rows with equal keys in a single streaming pass. The aggregation completes, only slower, instead of failing with an out-of-memory error.
Good to know: this is what the “Spill (Memory)” and “Spill (Disk)” columns in
the stage UI count. Steady spilling in an aggregation usually means too few
shuffle partitions, so each task’s hash map is too large; raising
spark.sql.shuffle.partitions is the first fix to try.
Binary row format and explicitly managed memory
“Execution memory is explicitly accounted for, without relying on JVM GC” is the sentence that matters.
How it works. Rows became compact binary records, UnsafeRow. A row starts
with a null bitmap, followed by one 8-byte slot per column: fixed-width values
such as integers and doubles sit in their slot directly, and variable-width
values such as strings store an offset and length pointing into a region at the
end of the row. A row with three integers is one small byte array, where as Java
objects it was a row object, an array and three boxed integers. Equality and
hashing work on the raw bytes, and rows can be sorted, shuffled and spilled
without being turned back into objects.
The memory those rows live in is allocated by Spark’s own memory manager in large pages, on or off the heap, and every operator must acquire memory before using it and spill when it cannot.
Good to know: before Tungsten, a Spark executor’s memory behaviour was the JVM’s behaviour, and tuning meant tuning a garbage collector. After it, Spark tracks its own execution memory and spills deliberately. The memory model that grows out of this is covered in Spark memory management.
Sort-merge join preferred for shuffle joins
Sort-merge join became the default strategy for joins too big to broadcast, replacing shuffled hash join, which needed each partition’s build side to fit in memory.
Good to know: shuffled hash join is faster when it fits, because it skips the
sort. It never went away: spark.sql.join.preferSortMergeJoin=false lets Spark
choose it when one side is much smaller, and 3.0’s SHUFFLE_HASH hint requests
it directly.
Around 100 new built-in functions, and a UDAF interface (SPARK-3947)
Date, string, math and collection functions, each one a reason not to write a UDF:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("fns15").getOrCreate()
spark.sql("""
SELECT date_format(DATE'2026-09-26', 'EEEE') AS weekday,
datediff(DATE'2026-12-31', DATE'2026-09-26') AS days_left,
regexp_extract('order-4821-pune', 'order-(\\\\d+)', 1) AS order_no,
concat_ws('|', 'a', NULL, 'c') AS joined,
round(stddev(x), 2) AS sd
FROM VALUES (1), (2), (3), (4) AS t(x)
""").show()
# +--------+---------+--------+------+----+
# | weekday|days_left|order_no|joined| sd|
# +--------+---------+--------+------+----+
# |Saturday| 96| 4821| a|c|1.29|
# +--------+---------+--------+------+----+
spark.stop()
How it works. Built-in functions are Catalyst expressions: they take part in
code generation, constant folding and pushdown, and they handle nulls by the SQL
rules. concat_ws skipping the NULL above is an example; a hand-written UDF
would have had to remember to.
Good to know: the UDAF interface let JVM code define custom aggregates. From
Python, the equivalent today is a grouped-aggregate pandas UDF, a function from
pd.Series to a scalar. The full function list for your version is
SHOW FUNCTIONS.
The broadcast hint (SPARK-8300)
The tool you reach for when the optimizer’s size estimate is wrong. With automatic broadcasting turned off, the same join produces two different physical operators depending only on the hint:
from pyspark.sql import SparkSession
from pyspark.sql.functions import broadcast
spark = SparkSession.builder.appName("bhj15").getOrCreate()
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1") # disable automatic broadcast
orders = spark.range(0, 1_000_000).selectExpr("id AS order_id", "id % 3 AS city_id")
cities = spark.createDataFrame([(0, "pune"), (1, "goa"), (2, "hyderabad")], "city_id LONG, city STRING")
def join_node(df):
plan = df._jdf.queryExecution().executedPlan().toString()
return next(l.strip(" +-*()0123456789").split(" ")[0] for l in plan.splitlines() if "Join" in l)
print(join_node(orders.join(cities, "city_id")))
# SortMergeJoin
print(join_node(orders.join(broadcast(cities), "city_id")))
# BroadcastHashJoin
spark.stop()
How it works. A broadcast hash join collects the small side to the driver, builds a hash table from it, and ships that table to every executor with the torrent mechanism from 1.1. Each task of the large side then probes the table locally. The large side is never shuffled, which is the whole point.
Good to know: the hint overrides the threshold, so it also lets you broadcast
something larger than 10 MB when you know it fits, or force it when the table’s
statistics are missing and Spark assumes it is huge. The cost lands on the
driver, which must hold the collected table, and on every executor’s memory;
broadcasts above spark.sql.broadcastTimeout (300 seconds) to build fail.
Spark 1.6, January 2016
The Dataset API (SPARK-9999)
The idea that 1.6 added and 2.0 made central: a collection of typed JVM
objects, checked at compile time, but executed by the same engine as a
DataFrame. It is a Scala and Java API only; Python and R have no compile-time
types to check. Run in spark-shell on 4.1.3:
case class Order(id: Long, city: String, amount: Double)
import spark.implicits._
val orders = Seq(Order(1, "pune", 120.0), Order(2, "goa", 80.0), Order(3, "pune", 45.0)).toDS()
val big = orders.filter(o => o.amount > 50) // typed lambda, checked at compile time
val byCity = big.groupByKey(_.city).mapValues(_.amount).reduceGroups(_ + _)
byCity.orderBy("key").show()
// +----+------------------------+
// | key|ReduceAggregator(double)|
// +----+------------------------+
// | goa| 80.0|
// |pune| 120.0|
// +----+------------------------+
How it works. The bridge between objects and Tungsten’s binary rows is an
encoder. For a case class, Spark generates code that converts an Order
object to an UnsafeRow and back, field by field, far faster than Java or Kryo
serialisation. Data stays in binary form inside the engine; objects are created
only where your code needs them. The price of the typing is that a lambda over a
typed object is opaque to Catalyst, just as it was in the RDD, and the encoder
has to deserialise each row into an object before your function sees it.
Good to know: a typo in o.amount is a compile error here and a runtime
AnalysisException with a DataFrame. That is the whole trade: earlier errors in
exchange for filters the optimizer cannot see into. Writing
orders.filter($"amount" > 50) gets both the Dataset type and a
pushdown-friendly expression. Types without a built-in encoder fall back to
Encoders.kryo, which stores each object as an opaque binary blob and loses
columnar benefits entirely.
Unified memory management (SPARK-10000)
How it works. Before 1.6, execution memory (shuffles, joins, sorts, aggregations) and storage memory (cached data, broadcasts) had fixed, separate fractions of the heap, so one could run out while the other sat idle. The unified manager gives them one pool with a borrowing rule: either side may use free memory from the other. The eviction is deliberately asymmetric. Execution can evict cached blocks, down to a protected storage floor; storage can never evict execution. A query that needs memory to finish beats a cache that is only an optimisation.
The defaults are still the 1.6 design, with one retune: spark.memory.fraction
started at 0.75 and was lowered to 0.6 in 2.0 to leave more room for user
objects. Reading the defaults from Spark’s own configuration entries, rather
than from a session where something may have overridden them, on 4.1.3:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("mem16").getOrCreate()
jvm = spark.sparkContext._jvm
cfg = getattr(getattr(jvm.org.apache.spark.internal.config, "package$"), "MODULE$")
for entry in [cfg.MEMORY_FRACTION(), cfg.MEMORY_STORAGE_FRACTION(), cfg.MEMORY_OFFHEAP_ENABLED()]:
print(entry.key(), "=", entry.defaultValueString())
# spark.memory.fraction = 0.6
# spark.memory.storageFraction = 0.5
# spark.memory.offHeap.enabled = false
spark.stop()
Worked through for an executor with an 8 GB heap: 300 MB is reserved, 60% of the remaining 7.7 GB, about 4.6 GB, is the unified pool, and half of that, about 2.3 GB, is storage’s protected floor. The other 40%, about 3.1 GB, is “user memory” for your own objects, UDF data structures and Spark’s internal metadata.
Good to know: PySpark’s Python workers live outside this heap entirely, so
their memory must be budgeted separately, with spark.executor.memoryOverhead
or spark.executor.pyspark.memory. An executor killed by YARN or Kubernetes for
exceeding its container limit is usually over-using that off-heap part, not the
heap.
Off-heap memory for query execution (SPARK-11389)
Tungsten’s binary rows can live outside the JVM heap, where the garbage collector never scans them.
How it works. With off-heap enabled, the memory manager allocates its pages
with Unsafe.allocateMemory instead of as Java byte arrays, and the unified
pool gains a second, off-heap region of the size you give. Execution and storage
can both use it.
--conf spark.memory.offHeap.enabled=true \
--conf spark.memory.offHeap.size=4g
Good to know: the off-heap size is not part of spark.executor.memory, so the
container has to be big enough for heap, off-heap and overhead together. It
helps most for executors with very large heaps, where full GC pauses were the
problem.
The first adaptive query execution (SPARK-9858)
This is dated 1.6, not 3.0.
How it works. The 1.6 version inserted a coordinator between the two sides of
a shuffle. After the map stage ran, it looked at the actual size of each
shuffle partition and combined small neighbouring partitions so that each
reducer read roughly spark.sql.adaptive.shuffle.targetPostShuffleInputSize
bytes. That is all it did: it chose the number of reducers for joins and
aggregations. It was rebuilt from scratch for 3.0, with a far wider set of
runtime changes, and enabled by default in 3.2.
Good to know: anyone who says AQE is new in 3.0 is describing the rewrite, not the idea, and tuning advice from 2016 that mentions adaptive execution is about this much smaller feature.
mapWithState for DStreams (SPARK-2629)
How it works. updateStateByKey, the earlier API, called your function for
every key in the state on every batch, so a job tracking a million users paid
for a million calls every few seconds, even if ten users were active. mapWithState
only calls the function for keys that received data in the batch, stores state
incrementally, and lets you set a timeout so idle keys are removed. The speed-up
on large state was often an order of magnitude.
Good to know: its ideas reappear in Structured Streaming’s
mapGroupsWithState in 2.2 and the State API v2 in 4.0, which is where new
stateful code should go.
Per-operator SQL metrics (SPARK-10412)
The SQL tab of the UI shows row counts, sizes and times on each operator of the physical plan. The same numbers are on the plan objects after a query runs:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("metrics16").config("spark.sql.adaptive.enabled", "false").getOrCreate()
df = spark.range(0, 1_000_000).where("id % 10 = 0").selectExpr("id % 7 AS k").groupBy("k").count()
df.collect()
def walk(node, depth=0):
m = node.metrics()
rows = m.get("numOutputRows")
label = node.nodeName()
if rows.isDefined():
print(" " * depth + f"{label}: {rows.get().value():,} rows")
else:
print(" " * depth + label)
ch = node.children()
for i in range(ch.size()):
walk(ch.apply(i), depth + 1)
walk(df._jdf.queryExecution().executedPlan())
# WholeStageCodegen (2)
# HashAggregate: 7 rows
# InputAdapter
# Exchange
# WholeStageCodegen (1)
# HashAggregate: 28 rows
# Project
# Filter: 100,000 rows
# Range: 1,000,000 rows
spark.stop()
Read bottom-up, the numbers tell the story of the query. Range produced a
million rows, Filter kept 100,000, the partial HashAggregate reduced them to
28 rows before the shuffle (7 keys in each of 4 tasks), and the final aggregate
produced the 7 groups. That partial aggregation before the Exchange is why a
groupBy shuffles 28 rows here instead of 100,000. The WholeStageCodegen
nodes carry no row count because they are wrappers around the fused operators.
Good to know: most performance debugging starts by reading these numbers top-down. The operator where the row count drops sharply is doing the useful filtering, and it should happen as early, as far down the plan, as possible; an operator whose output is far larger than its input, such as a join producing more rows than either side, is where to look for a missing join condition.
The 2.x line: unification and codegen
If 1.x discovered that the DataFrame was the right abstraction, 2.x committed to it: one entry point, one type, one execution strategy, and streaming redefined as a query over an unbounded table. Read the 2.x features in three groups: 2.0 unifies the APIs and compiles whole queries; 2.1 to 2.3 make Structured Streaming production-ready and bring Python up to speed with Arrow; 2.4 fills in SQL and data-source gaps and becomes the long-term 2.x release.
Spark 2.0, July 2016
SparkSession
One entry point replacing SQLContext and HiveContext, with the
SparkContext still available underneath.
How it works. A Spark application has exactly one SparkContext, which owns
the connection to the cluster manager, the executors and the scheduler. A
SparkSession sits on top of it and owns the SQL side: the SQL configuration,
the catalog, temporary views and registered functions. The builder returns the
existing session if there is one:
from pyspark.sql import SparkSession
spark = (SparkSession.builder
.appName("session20")
.master("local[2]")
.config("spark.sql.shuffle.partitions", "8")
.enableHiveSupport()
.getOrCreate())
print(spark.sparkContext is not None, spark.catalog.currentDatabase(), spark.conf.get("spark.sql.catalogImplementation"))
# True default hive
spark.stop()
The session also brought a catalog API, so code can inspect tables and views
without parsing SHOW output:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("catalog20").getOrCreate()
spark.range(5).createOrReplaceTempView("recent_orders")
spark.range(5).write.mode("overwrite").saveAsTable("orders_t")
print([(t.name, t.tableType, t.isTemporary) for t in spark.catalog.listTables()])
print([c.name for c in spark.catalog.listColumns("orders_t")])
print(spark.catalog.tableExists("orders_t"), spark.catalog.currentCatalog())
# [('orders_t', 'MANAGED', False), ('recent_orders', 'TEMPORARY', True)]
# ['id']
# True spark_catalog
spark.sql("DROP TABLE orders_t")
spark.stop()
Good to know: getOrCreate() returns an existing session if there is one, and
then silently ignores most of the builder’s config calls. Settings that shape
the JVM or the executors, such as memory and cores, must be set before the first
session is created, usually on the spark-submit command line.
spark.newSession() creates a second session with its own SQL settings and
temporary views but the same SparkContext, which is how a server can isolate
users cheaply.
DataFrame becomes Dataset[Row]
In Scala and Java, DataFrame became a type alias for a Dataset of generic
rows, so the two APIs share one implementation.
How it works. A Row is an untyped record whose fields are looked up by name
or position at runtime. Dataset[Order] and Dataset[Row] hold the same binary
data; the only difference is which encoder converts it when your code asks for
objects. Untyped operations (select("amount"), where("amount > 50"))
return DataFrames and are fully optimizable; typed operations (map,
filter(o => ...)) return Datasets and cost a deserialisation per row.
Good to know: Python has only DataFrames. Python objects have no compile-time types, so there is nothing for a typed Dataset to check.
SQL 2003 support, and all 99 TPC-DS queries
A native SQL parser replaced the Hive-derived one, with native DDL and support
for correlated and uncorrelated subqueries in SELECT and WHERE:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("subq20").getOrCreate()
spark.createDataFrame([(1, "pune"), (2, "goa"), (3, "delhi")], "id INT, city STRING").createOrReplaceTempView("customers")
spark.createDataFrame([(1, 500), (1, 200), (2, 50)], "cust_id INT, amount INT").createOrReplaceTempView("orders")
spark.sql("""
SELECT c.city,
(SELECT sum(amount) FROM orders o WHERE o.cust_id = c.id) AS spent
FROM customers c
WHERE EXISTS (SELECT 1 FROM orders o WHERE o.cust_id = c.id AND o.amount > 100)
""").show()
# +----+-----+
# |city|spent|
# +----+-----+
# |pune| 700|
# +----+-----+
spark.stop()
How it works. Spark never runs a correlated subquery once per outer row. The
optimizer decorrelates it: EXISTS becomes a left semi join, NOT EXISTS
a left anti join, and a scalar subquery such as spent becomes an aggregation
joined back to the outer query. The plan for the query above contains joins,
not subqueries.
Good to know: because subqueries become joins, they cost what the equivalent
join costs, and they are subject to the same join-strategy choices. NOT IN
with a subquery is the exception to watch: SQL’s null semantics force it into a
null-aware anti join, which is much slower than NOT EXISTS. Prefer
NOT EXISTS when the column can be null.
Whole-stage code generation
The defining change of 2.0.
How it works. Before it, Spark used the Volcano execution model: each
operator was an object with a next() method that pulled one row from its child,
so a scan, a filter and a projection meant three virtual calls per row, and rows
were materialised between every pair of operators. Whole-stage codegen fuses a
chain of operators into one generated Java method with a single loop: read a
value, test the filter, compute the projection, all in local variables, with no
row objects in between. The JIT then compiles it like hand-written code. The
release notes claim speedups of two to ten times for common operators.
You can read the generated code. A filter and a projection become a few lines
inside one processNext() loop:
import io, contextlib
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("codegen20").getOrCreate()
df = spark.range(0, 1000).where("id % 2 = 0").selectExpr("id * 3 AS tripled")
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
df.explain(mode="codegen")
print(buf.getvalue())
spark.stop()
The output is about 150 lines of Java. The inner loop, trimmed of setup and metrics code, is the whole query:
/* 093 */ for (int range_localIdx_0 = 0; range_localIdx_0 < range_localEnd_0; range_localIdx_0++) {
/* 094 */ long range_value_0 = ((long)range_localIdx_0 * 1L) + range_nextIndex_0;
/* 095 */
/* 096 */ do {
/* ... */
/* 104 */ filter_value_1 = (long)(range_value_0 % 2L);
/* ... */
/* 108 */ filter_value_0 = filter_value_1 == 0L;
/* ... */
/* 111 */ if (filter_isNull_0 || !filter_value_0) continue;
/* ... */
/* 119 */ project_value_0 = org.apache.spark.sql.catalyst.util.MathUtils.multiplyExact(range_value_0, 3L, ...);
/* ... */
/* 122 */ range_mutableStateArray_0[2].write(0, project_value_0);
/* 123 */ append((range_mutableStateArray_0[2].getRow()));
/* 124 */
/* 125 */ } while (false);
The range, the filter and the projection are one for loop over primitive
long values: generate range_value_0, continue if it is odd, multiply, write
the result. No row object exists until the final append. Line 119 also shows a
later release at work: id * 3 compiles to MathUtils.multiplyExact, which
throws on overflow, because 4.0 turned ANSI mode on. On 3.5.9, where ANSI mode
is still off, the same line 119 is a plain multiplication:
/* 119 */ project_value_0 = range_value_0 * 3L;
Every plan today shows where the fused stages are. The * markers and the
number after them are codegen stage identifiers:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("codegen").getOrCreate()
df = (spark.range(0, 1_000_000)
.selectExpr("id", "id % 7 AS bucket")
.groupBy("bucket").count())
df.collect()
print(df._jdf.queryExecution().executedPlan().toString())
spark.stop()
On 4.1.3 that prints:
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
ResultQueryStage 1
+- *(2) HashAggregate(keys=[bucket#1L], functions=[count(1)], output=[bucket#1L, count#2L])
+- AQEShuffleRead coalesced
+- ShuffleQueryStage 0
+- Exchange hashpartitioning(bucket#1L, 200), ENSURE_REQUIREMENTS, [plan_id=25]
+- *(1) HashAggregate(keys=[bucket#1L], functions=[partial_count(1)], output=[bucket#1L, count#7L])
+- *(1) Project [(id#0L % 7) AS bucket#1L]
+- *(1) Range (0, 1000000, step=1, splits=11)
+- == Initial Plan ==
HashAggregate(keys=[bucket#1L], functions=[count(1)], output=[bucket#1L, count#2L])
+- Exchange hashpartitioning(bucket#1L, 200), ENSURE_REQUIREMENTS, [plan_id=15]
+- HashAggregate(keys=[bucket#1L], functions=[partial_count(1)], output=[bucket#1L, count#7L])
+- Project [(id#0L % 7) AS bucket#1L]
+- Range (0, 1000000, step=1, splits=11)
Three eras are visible in that one output. *(1) and *(2) are 2.0’s
whole-stage codegen: the range scan, the projection and the partial aggregate
were compiled into a single generated method, and the final aggregate into
another. The Exchange is where one stage ends and the next begins, because a
shuffle cannot be fused. AdaptiveSparkPlan with its Initial Plan and
Final Plan sections is 3.0. AQEShuffleRead coalesced is 3.2 deciding at
runtime that 200 shuffle partitions were too many for seven groups and reading
them as fewer.
Switching codegen off shows what it changes. The same query, planned with
spark.sql.codegen.wholeStage on and off:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("wscg").getOrCreate()
spark.conf.set("spark.sql.adaptive.enabled", "false")
for flag in ["true", "false"]:
spark.conf.set("spark.sql.codegen.wholeStage", flag)
q = spark.range(0, 1_000_000).where("id % 7 = 3").selectExpr("sum(id) AS s")
plan = q._jdf.queryExecution().executedPlan().toString()
print(f"-- wholeStage={flag}")
print("\n".join(l.split(", [plan_id")[0] for l in plan.splitlines()))
spark.stop()
-- wholeStage=true
*(2) HashAggregate(keys=[], functions=[sum(id#9L)], output=[s#13L])
+- Exchange SinglePartition, ENSURE_REQUIREMENTS
+- *(1) HashAggregate(keys=[], functions=[partial_sum(id#9L)], output=[sum#16L])
+- *(1) Filter ((id#9L % 7) = 3)
+- *(1) Range (0, 1000000, step=1, splits=12)
-- wholeStage=false
HashAggregate(keys=[], functions=[sum(id#18L)], output=[s#22L])
+- Exchange SinglePartition, ENSURE_REQUIREMENTS
+- HashAggregate(keys=[], functions=[partial_sum(id#18L)], output=[sum#25L])
+- Filter ((id#18L % 7) = 3)
+- Range (0, 1000000, step=1, splits=12)
The operators are identical; only the *(n) markers are gone. With them, the
Range, Filter and partial HashAggregate run as one generated loop, stage 1.
Without them, each operator is a separate object passing rows to the next one
at a time, the Volcano model. That per-row, per-operator call overhead is what
whole-stage codegen removes, and it grows with the number of rows and
operators. I have not quoted a timing: a laptop measurement of the gap says more
about the laptop than about Spark.
Good to know: codegen has limits. An operator that does not support it, such
as a Python UDF, breaks the chain into separate stages. A query with more than
spark.sql.codegen.maxFields (100) columns in a stage, or a generated method
larger than the JIT will compile (spark.sql.codegen.hugeMethodLimit, 65,535
bytes of bytecode), falls back to slower paths, which is one reason very wide
tables can be unexpectedly slow. And the Exchange in the plan above still says
200, because spark.sql.shuffle.partitions is still the 1.x default; AQE
coalesced after the fact rather than planning a better number in the first
place.
Vectorized Parquet reader
How it works. The row-based reader decoded one record at a time into a row
object. The vectorized reader decodes one column at a time, a batch of
values (4,096 by default, spark.sql.parquet.columnarReaderBatchSize) into a
column vector, which is a tight loop over a primitive array that modern CPUs run
very fast. The batch then feeds the plan either as columns or, through the
ColumnarToRow step you see above file scans in plans, as rows.
Good to know: it is on by default
(spark.sql.parquet.enableVectorizedReader). Until 3.3 it only handled flat
columns, so a table with nested structs, arrays or maps silently used the slow
row reader; that is one of the quieter reasons upgrades to 3.3 speed up
nested-data jobs.
Native CSV source
The spark-csv package moved into Spark itself.
How it works. The reader is built on the univocity parser. Each line is split into fields and converted to the schema’s types, and the mode decides what happens to a line that does not convert. The same messy file, read three ways:
import os, tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("csv20").getOrCreate()
d = tempfile.mkdtemp()
with open(os.path.join(d, "o.csv"), "w") as f:
f.write("order_id,amount\n1,120\n2,oops\n3,45\n")
schema = "order_id INT, amount INT, _corrupt_record STRING"
(spark.read.option("header", "true").schema(schema)
.option("mode", "PERMISSIVE").option("columnNameOfCorruptRecord", "_corrupt_record")
.csv(d).show())
dropped = spark.read.option("header", "true").schema("order_id INT, amount INT").option("mode", "DROPMALFORMED").csv(d)
print("count():", dropped.count(), " len(collect()):", len(dropped.collect()))
try:
spark.read.option("header", "true").schema("order_id INT, amount INT").option("mode", "FAILFAST").csv(d).collect()
except Exception as e:
print([l for l in str(e).splitlines() if "MALFORMED" in l][0].split("SparkException: ")[-1][:100])
# +--------+------+---------------+
# |order_id|amount|_corrupt_record|
# +--------+------+---------------+
# | 1| 120| NULL|
# | 2| NULL| 2,oops|
# | 3| 45| NULL|
# +--------+------+---------------+
# count(): 3 len(collect()): 2
# [MALFORMED_RECORD_IN_PARSING.WITHOUT_SUGGESTION] Malformed records are detected in record parsing: [
spark.stop()
PERMISSIVE, the default, keeps the bad row with nulls and puts the raw line in
the corrupt-record column, so nothing is lost silently. DROPMALFORMED discards
it. FAILFAST stops the job.
Good to know: look at the DROPMALFORMED line. count() said 3 and
collect() returned 2 rows, from the same DataFrame. count() does not need
any column values, so Spark skips converting the fields, never discovers that
oops is not an integer, and counts the line. Any check of “how many rows
survived” must look at the parsed values, or use PERMISSIVE and count the rows
where the corrupt-record column is not null. The same applies to JSON.
Hive-style bucketing
Writing a table with bucketBy pre-hashes rows into a fixed number of files by
key. Two tables bucketed the same way on the join key can be joined with no
shuffle at all:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("bucket20").config("spark.sql.autoBroadcastJoinThreshold", "-1").getOrCreate()
for name in ["orders_b", "payments_b"]:
spark.range(0, 100_000).selectExpr("id % 1000 AS cust_id", "id AS v") \
.write.bucketBy(8, "cust_id").sortBy("cust_id").mode("overwrite").saveAsTable(name)
df = spark.table("orders_b").join(spark.table("payments_b"), "cust_id")
plan = df._jdf.queryExecution().executedPlan().toString()
print("Exchange in plan:", "Exchange" in plan)
# Exchange in plan: False
spark.sql("DROP TABLE orders_b"); spark.sql("DROP TABLE payments_b")
spark.stop()
How it works. A sort-merge join shuffles both sides so that equal keys land in
the same partition. Bucketing does that shuffle once, at write time: each row
goes to bucket hash(cust_id) mod 8, and the table’s metadata records the bucket
column and count. At read time, Spark knows bucket n of one table can only
match bucket n of the other, so it reads them as eight partitions already
aligned, and skips the Exchange. With sortBy, each bucket file is also
pre-sorted, which can let Spark skip the sort as well.
Good to know: bucketing only works with saveAsTable, because the bucket spec
lives in the catalog, and both sides need the same bucket count on the join key.
Each write task produces a file per bucket it touches, so 200 tasks writing into
8 buckets can create 1,600 small files; repartition by the bucket column before
writing. Spark’s bucket hashing is not compatible with Hive’s, and table formats
such as Iceberg offer a more flexible version (bucket partition transforms).
Structured Streaming, experimental
DStreams, from 0.7, modelled a stream as a sequence of small RDDs, so you wrote RDD code and the engine could not optimise across batches.
How it works. Structured Streaming models a stream as a table that keeps
growing. You write an ordinary DataFrame query against it, and the engine
works out how to run that query incrementally: each trigger it plans the query
over only the new data, updates whatever state the query needs (running
aggregates, open windows, join buffers), and writes the change to the sink. The
same groupBy works on a static DataFrame and on a stream:
import json, os, tempfile
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("ss20").getOrCreate()
spark.conf.set("spark.sql.shuffle.partitions", "4")
src = tempfile.mkdtemp()
with open(os.path.join(src, "batch1.json"), "w") as f:
for city, amount in [("pune", 10), ("goa", 5), ("pune", 7)]:
f.write(json.dumps({"city": city, "amount": amount}) + "\n")
events = spark.readStream.schema("city STRING, amount INT").json(src)
totals = events.groupBy("city").agg(F.sum("amount").alias("total"))
q = (totals.writeStream
.outputMode("complete")
.format("memory").queryName("city_totals")
.trigger(availableNow=True)
.start())
q.awaitTermination()
spark.sql("SELECT * FROM city_totals ORDER BY city").show()
# +----+-----+
# |city|total|
# +----+-----+
# | goa| 5|
# |pune| 17|
# +----+-----+
spark.stop()
Two details in that example come from later releases. trigger(availableNow=True)
arrived in 3.3: it processes everything available and then stops, which makes
streaming code testable as a batch job. And the file source needs a schema up
front, because a stream cannot infer one from files that do not exist yet.
The three parts of every streaming query are the source (files, Kafka,
sockets, the rate test source), the trigger (how often to run: as fast as
possible by default, processingTime="1 minute", or availableNow) and the
sink with its output mode:
completerewrites the whole result table every batch, which only makes sense for small aggregates.appendemits a row once and never again, which for an aggregation means waiting until the row can no longer change, and that requires a watermark (2.1).updateemits only the rows that changed in this batch.
Good to know: the checkpoint location is what makes a query restartable, and it is tied to the query: change the query’s stateful shape and the old checkpoint can no longer be used. The internals are in Structured Streaming internals.
Spark 2.1, December 2016
Event-time watermarks (SPARK-18124)
How it works. Streaming data arrives out of order. A record’s event time,
when it happened, can be well behind its processing time, when Spark sees
it, because of network delays, offline devices or retries. An event-time window
aggregation therefore cannot know when a window is complete, and without a rule
for closing windows its state grows without bound. A watermark is that rule: the
maximum event time seen so far, minus a delay you declare. A window is
finalised, emitted in append mode and dropped from state, once the watermark
passes its end. A record that arrives for a window that has already been
finalised is discarded. Three micro-batches show all of it:
import json, os, tempfile
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("wm21").getOrCreate()
spark.conf.set("spark.sql.shuffle.partitions", "4")
src, out, chk = tempfile.mkdtemp(), tempfile.mkdtemp(), tempfile.mkdtemp()
def drop(name, rows):
with open(os.path.join(src, name), "w") as f:
for ts, city in rows:
f.write(json.dumps({"ts": ts, "city": city}) + "\n")
events = (spark.readStream.schema("ts TIMESTAMP, city STRING").json(src)
.withWatermark("ts", "10 minutes")
.groupBy(F.window("ts", "5 minutes"), "city").count())
def run():
q = (events.writeStream.outputMode("append").format("parquet")
.option("path", out).option("checkpointLocation", chk)
.trigger(availableNow=True).start())
q.awaitTermination()
return q
drop("1.json", [("2026-01-01 10:01:00", "pune"), ("2026-01-01 10:03:00", "pune")])
run()
drop("2.json", [("2026-01-01 10:30:00", "goa")]) # advances the watermark to 10:20
run()
drop("3.json", [("2026-01-01 10:02:00", "pune")]) # 28 minutes late: dropped
q = run()
(spark.read.parquet(out).selectExpr("window.start", "window.end", "city", "count")
.orderBy("start").show(truncate=False))
print("watermark:", q.lastProgress["eventTime"].get("watermark"))
spark.stop()
On 4.1.3, and identically on 3.5.9:
+-------------------+-------------------+----+-----+
|start |end |city|count|
+-------------------+-------------------+----+-----+
|2026-01-01 10:00:00|2026-01-01 10:05:00|pune|2 |
+-------------------+-------------------+----+-----+
watermark: 2026-01-01T10:20:00.000Z
The 10:00 window was emitted with a count of 2, not 3. The late 10:02 record arrived after the watermark had moved to 10:20, so its window was already closed and the record was dropped. The goa window for 10:30 is still open, waiting for the watermark to pass 10:35, so it has not been written at all. That is the watermark trade-off in one table: bounded state and final results, paid for with the late record.
Good to know: choosing the delay is a business decision: longer delays accept
more late data but keep more state and emit results later. When a query reads
several streams, each with its own watermark, Spark uses the minimum by
default, so the slowest stream holds everyone back
(spark.sql.streaming.multipleWatermarkPolicy=max changes that, from 2.4). And
my first attempt at this used the memory sink: the second run failed with
This query does not support recovering from checkpoint location, because the
memory sink cannot resume from a checkpoint.
Kafka 0.10 source for Structured Streaming (SPARK-17346)
The Kafka source that is still the standard way to stream from Kafka. It needs
the spark-sql-kafka-0-10 package on the classpath:
events = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("subscribe", "orders")
.option("startingOffsets", "earliest")
.option("maxOffsetsPerTrigger", 100_000)
.load()
.selectExpr("CAST(key AS STRING) AS k", "CAST(value AS STRING) AS v",
"topic", "partition", "offset", "timestamp"))
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.13:4.1.3 app.py
How it works. The source applies the direct-stream design from 1.3 to
Structured Streaming. Each trigger it decides an offset range per topic
partition, writes that decision to the checkpoint’s offset log before reading
anything, and reads each range as one Spark partition. Every record arrives as
binary key and value columns plus its topic, partition, offset and timestamp,
and you decode the payload yourself with CAST, from_json (below) or
from_avro (2.4).
Good to know: offsets are tracked in the query’s checkpoint, not committed
back to Kafka, so consumer-group lag tools do not see a Structured Streaming
query’s progress; monitor lastProgress or a StreamingQueryListener instead.
maxOffsetsPerTrigger caps each batch, which matters on the first run against a
topic with a large backlog. failOnDataLoss (default true) stops the query
when offsets it expected have been deleted by retention, which is usually what
you want to know about.
Stable offset-log format (SPARK-17829)
How it works. A streaming query’s checkpoint directory holds an offsets/ log
(the input range each batch is about to process), a commits/ log (batches that
finished), a metadata file with the query ID, the sources/ directory, and the
state/ directory for stateful operators. This release gave the offset log a
versioned, stable format, so a checkpoint written by one release can be resumed
by later ones.
Good to know: never copy a checkpoint between queries or edit it by hand to
“skip” data; change startingOffsets on a new checkpoint instead. The layout is
walked through file by file in
Structured Streaming internals.
from_json and to_json (SPARK-18351)
Parse a JSON string column into a struct, or serialise a struct back to JSON, without leaving the engine. The workhorse for Kafka payloads:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("json21").getOrCreate()
raw = spark.createDataFrame([('{"city":"pune","amount":120}',), ('{"city":"goa"}',)], "payload STRING")
parsed = raw.select(F.from_json("payload", "city STRING, amount INT").alias("p")).select("p.*")
parsed.show()
parsed.select(F.to_json(F.struct("city", "amount")).alias("json")).show(truncate=False)
# +----+------+
# |city|amount|
# +----+------+
# |pune| 120|
# | goa| NULL|
# +----+------+
# +----------------------------+
# |json |
# +----------------------------+
# |{"city":"pune","amount":120}|
# |{"city":"goa"} |
# +----------------------------+
spark.stop()
How it works. from_json parses each string against the schema you give it,
producing a struct column; fields missing from a document come out NULL, and a
document that does not parse at all becomes a NULL struct by default rather
than an error. to_json walks a struct and writes it back out.
Good to know: to_json leaves out null fields, as the second row shows.
schema_of_json, from 2.4, infers a DDL schema from a sample document, which is
a convenient way to write the schema for from_json. For payloads whose schema
you do not control, 4.0’s VARIANT type is the better fit.
Scalable partition handling (SPARK-17861)
How it works. Before this change, a data source table’s partitions were discovered by listing its directory tree every time a query planned, which took minutes on tables with tens of thousands of partitions. Now partition metadata lives in the metastore, and a query that filters on the partition column asks the metastore for just the matching partitions. The trade-off is that the metastore must be told about partitions written outside Spark:
import os, tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("msck21").getOrCreate()
path = tempfile.mkdtemp()
spark.sql(f"CREATE TABLE events (id INT, day STRING) USING parquet PARTITIONED BY (day) LOCATION '{path}'")
spark.sql("INSERT INTO events VALUES (1, '2026-09-25')")
spark.range(2, 4).selectExpr("CAST(id AS INT) AS id").write.parquet(os.path.join(path, "day=2026-09-26"))
print("before:", [r[0] for r in spark.sql("SHOW PARTITIONS events").collect()], spark.table("events").count())
spark.sql("MSCK REPAIR TABLE events")
print("after: ", [r[0] for r in spark.sql("SHOW PARTITIONS events").collect()], spark.table("events").count())
# before: ['day=2026-09-25'] 1
# after: ['day=2026-09-25', 'day=2026-09-26'] 3
spark.sql("DROP TABLE events")
spark.stop()
The two rows written straight to the day=2026-09-26 directory are invisible
until MSCK REPAIR TABLE registers the partition.
Good to know: this is a classic “my data is there but the query returns
nothing” cause. MSCK REPAIR lists the whole directory, so for a daily load
ALTER TABLE events ADD PARTITION (day='2026-09-26') is much cheaper. Table
formats avoid the problem entirely by committing files and partitions together.
Spark 2.2, July 2017
Structured Streaming declared generally available (SPARK-20844)
The experimental label came off, two releases and a year after 2.0.
Good to know: production advice written between 2.0 and 2.2 is about a moving target; APIs such as the Kafka source options and the watermark semantics changed in that window.
Arbitrary stateful processing: mapGroupsWithState
For logic that windows and aggregations cannot express, such as sessionising user activity with custom rules or tracking a device’s state machine.
How it works. After groupByKey, your function is called once per key per
batch with the key, the new records for it, and a GroupState object holding
that key’s state from earlier batches. The function reads and updates the state,
can set a timeout so that keys that go quiet are called one last time, and
returns output rows. flatMapGroupsWithState is the variant that can emit any
number of rows.
Good to know: it is a Scala and Java API. Python got the equivalent,
applyInPandasWithState, in 3.4, and the State API v2 in 4.0
(transformWithState) replaces both with typed, named state variables, TTLs and
timers.
Kafka sink for Structured Streaming
Streaming queries can write to Kafka: the query’s output needs a value
column, optionally key and topic:
(events.selectExpr("CAST(order_id AS STRING) AS key", "to_json(struct(*)) AS value")
.writeStream.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("topic", "orders_enriched")
.option("checkpointLocation", "/chk/orders_enriched")
.start())
Good to know: the Kafka sink is at-least-once: after a failure, the last batch can be written again. Downstream consumers need to tolerate duplicates, usually by being idempotent on a key.
Cost-based optimizer (SPARK-17075 and related)
Cardinality estimation for filter, join, aggregate, project and limit, using
statistics you compute yourself with ANALYZE TABLE.
How it works. The rule-based optimizer knows the shape of a query but not the
data. The cost-based optimizer estimates how many rows each operator will
produce: a filter’s selectivity from the column’s min, max and distinct count, a
join’s output from the distinct counts of its keys, and so on up the plan. Those
estimates drive join order and join strategy. ANALYZE TABLE with FOR COLUMNS
writes the statistics into the catalog:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("cbo22").getOrCreate()
spark.range(0, 10_000).selectExpr("id", "id % 10 AS bucket").write.mode("overwrite").saveAsTable("t22")
spark.sql("ANALYZE TABLE t22 COMPUTE STATISTICS FOR COLUMNS id, bucket")
for r in spark.sql("DESCRIBE EXTENDED t22 bucket").collect():
if r[0] in ("min", "max", "num_nulls", "distinct_count"):
print(r[0], r[1])
# min 0
# max 9
# num_nulls 0
# distinct_count 10
stats = [r for r in spark.sql("DESCRIBE EXTENDED t22").collect() if r[0] == "Statistics"]
print(stats[0][1])
# 43711 bytes, 10000 rows
spark.sql("DROP TABLE t22")
spark.stop()
With these numbers, WHERE bucket = 3 is estimated at 1/10 of 10,000 rows,
1,000, assuming values are evenly spread. Without the FOR COLUMNS part you get
the table’s size and row count but no column statistics, which is enough for the
broadcast decision and not enough to estimate a filter’s selectivity.
Good to know: the optimizer only uses column statistics when
spark.sql.cbo.enabled is true. On 4.1.3 both it and
spark.sql.cbo.joinReorder.enabled still default to false. Statistics are
also a snapshot: they are not updated when the table changes, and stale
statistics can produce worse plans than none. 2.3’s histograms improve the “evenly
spread” assumption.
Cost-based join reordering (SPARK-17080)
How it works. With statistics and spark.sql.cbo.joinReorder.enabled=true,
Spark searches join orders with dynamic programming, estimating the size of
every intermediate result, and picks the order that keeps intermediate results
smallest, instead of the order you wrote. The search is exponential in the
number of tables, so it only runs for up to spark.sql.cbo.joinReorder.dp.threshold
(12) tables.
Good to know: be realistic about it: the
Catalyst post runs a
three-way join on 4.1.3 where CostBasedJoinReorder ran, with statistics, and
changed nothing: the joins stayed in the order they were written, the filtered
small dimension last. In practice, AQE (3.0 onwards) fixes more bad join
decisions at runtime than CBO prevents at planning time.
Broadcast hints in SQL (SPARK-16475)
The 1.5 DataFrame hint, now in SQL as BROADCAST, BROADCASTJOIN or
MAPJOIN, all synonyms:
SELECT /*+ BROADCAST(c) */ o.*, c.city
FROM orders o JOIN cities c ON o.city_id = c.city_id
Good to know: the hint names the table alias used in the query, c, not
the underlying table name. A misspelt or unresolvable hint is ignored with only a
warning in the log, so check the plan for BroadcastHashJoin rather than
trusting the hint. 3.0 added hints for every other join strategy.
Session-local time zone (SPARK-18350)
spark.sql.session.timeZone decides how TIMESTAMP values are displayed and
how strings without an offset are interpreted, per session instead of per JVM:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("tz22").getOrCreate()
spark.conf.set("spark.sql.session.timeZone", "UTC")
spark.sql("SELECT CAST(TIMESTAMP'2026-01-01 12:00:00' AS STRING) AS ts, current_timezone() AS tz").show()
spark.conf.set("spark.sql.session.timeZone", "Asia/Kolkata")
spark.sql("SELECT from_utc_timestamp(TIMESTAMP'2026-01-01 12:00:00', 'Asia/Kolkata') AS ist, current_timezone() AS tz").show()
# +-------------------+---+
# | ts| tz|
# +-------------------+---+
# |2026-01-01 12:00:00|UTC|
# +-------------------+---+
# +-------------------+------------+
# | ist| tz|
# +-------------------+------------+
# |2026-01-01 17:30:00|Asia/Kolkata|
# +-------------------+------------+
spark.stop()
How it works. Internally a Spark TIMESTAMP is a count of microseconds since
the epoch in UTC, an instant. The session time zone is applied only at the
edges: when a string such as '2026-01-01 12:00:00' is parsed into an instant,
and when an instant is displayed or cast back to a string.
Good to know: set it explicitly, usually to UTC, in every job. Left unset it
follows the JVM’s default time zone, so the same job can produce different
output on a laptop and on a cluster, and daily partitions can shift by a day at
midnight boundaries. 3.4’s TIMESTAMP_NTZ is the type for values that should
never shift.
Task blacklisting (SPARK-8425)
Executors and nodes where tasks keep failing are excluded from scheduling.
How it works. Spark counts task failures per executor and per node. After
spark.excludeOnFailure.task.maxTaskAttemptsPerExecutor failures of the same
task on an executor, that task avoids the executor; after enough failures across
tasks in a stage or application, the whole executor or node is excluded for a
timeout, so a node with a bad disk stops failing task after task.
Good to know: it is off by default. 3.1 renamed the feature to exclude on
failure; the switch is now spark.excludeOnFailure.enabled, and the old
spark.blacklist.* names still work but are deprecated.
Java 7 removed (SPARK-19493)
Java 8 became the minimum, which let Spark’s own code use lambdas and the
java.time API.
pip install pyspark
It belongs in the same list as the engine features. It is the moment PySpark stopped requiring a cluster distribution to try:
pip install pyspark
python -c "from pyspark.sql import SparkSession; print(SparkSession.builder.getOrCreate().range(3).count())"
How it works. The wheel bundles the Spark JARs, so pip install gives you a
complete local Spark that starts a JVM behind the scenes; it still needs a Java
runtime on the machine.
Good to know: the PyPI version must match the cluster’s Spark version when you
submit to a cluster, or you can hit serialisation errors between the two. For
Spark Connect clients, 4.0 added the much smaller pyspark-client package with
no JARs at all.
Spark 2.3, February 2018
A dense release.
Kubernetes scheduler backend, experimental (SPARK-18278)
Spark could run its driver and executors as Kubernetes pods, with the API server as the cluster manager:
spark-submit \
--master k8s://https://<api-server>:6443 \
--deploy-mode cluster \
--conf spark.kubernetes.container.image=apache/spark:4.1.3-python3 \
--conf spark.kubernetes.namespace=etl \
--conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
--conf spark.executor.instances=4 \
local:///opt/spark/examples/src/main/python/pi.py
How it works. spark-submit asks the API server to create a driver pod.
The driver, running with a service account that is allowed to create pods, then
asks the API server for executor pods and talks to them directly. Executor
pods are owned by the driver pod, so deleting the driver cleans up its
executors. Container images are built from the distribution with
bin/docker-image-tool.sh.
Good to know: the local:// scheme means “already inside the image”, not the
machine you run spark-submit on. The service account needs RBAC permission to
create and delete pods in the namespace, which is the most common first-run
failure. Kubernetes support went GA in 3.1.
Native vectorized ORC reader (SPARK-16060)
A native ORC reader, enabled with spark.sql.orc.impl=native, which 2.4 made
the default.
How it works. The same columnar-batch idea as the 2.0 Parquet reader, written against the ORC library directly instead of going through Hive’s row-at-a-time reader and object inspectors.
Data Source API V2, experimental (SPARK-15689)
A redesigned connector API.
How it works. The 1.2 API (V1) gave sources a fixed menu, columns and filters,
and tied them to RDDs and the SQL internals. V2 is a set of Java interfaces a
source implements to declare what it can do: a ScanBuilder that can accept
pushed-down filters, column pruning, limits and aggregates; columnar reads;
streaming reads; and transactional writes with a commit protocol that makes a
write all-or-nothing. It was redesigned again in 3.0, together with the catalog
API, into the form every table format now uses.
Good to know: you rarely implement it yourself; you meet it when a connector advertises “DSv2”. From 4.0 you can write a V2 source in Python; see the Python Data Source API in the 4.0 section.
Continuous processing, experimental
An execution mode for millisecond latency.
How it works. Micro-batch execution plans and schedules a small job per trigger, which puts a floor of roughly 100 milliseconds under end-to-end latency. Continuous processing instead starts long-running tasks that read, process and write records as they arrive, and marks progress with epoch markers that flow through the data, checkpointed at the interval you give. It supports map-like queries only, such as projections and filters, with no aggregations or joins:
import time
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("cont23").getOrCreate()
q = (spark.readStream.format("rate").option("rowsPerSecond", 10).load()
.selectExpr("value", "value % 2 = 0 AS even")
.writeStream.format("memory").queryName("rate_out")
.trigger(continuous="1 second").start())
time.sleep(8)
q.stop()
print("rows received:", spark.sql("SELECT count(*) FROM rate_out").first()[0] > 0)
# rows received: True
spark.stop()
Good to know: the "1 second" is the checkpoint interval, not a batch
interval, and delivery is at-least-once. Continuous processing never left
experimental status; 4.1’s real-time mode is the project’s second attempt at
the same goal.
Stream-to-stream joins
Harder than they look because both sides are unbounded.
How it works. To join an impression with a click, Spark keeps every impression in state in case its click arrives later, and every click in case its impression arrives later. Without a limit, both buffers grow forever. Give each side a watermark and bound the join condition in time, and Spark can work out when a buffered row can no longer find a match and evict it:
import json, os, tempfile
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("ssjoin").getOrCreate()
spark.conf.set("spark.sql.shuffle.partitions", "4")
imp_dir, clk_dir = tempfile.mkdtemp(), tempfile.mkdtemp()
with open(os.path.join(imp_dir, "i.json"), "w") as f:
f.write(json.dumps({"ad_id": 1, "imp_time": "2026-01-01 10:00:00"}) + "\n")
f.write(json.dumps({"ad_id": 2, "imp_time": "2026-01-01 10:00:00"}) + "\n")
with open(os.path.join(clk_dir, "c.json"), "w") as f:
f.write(json.dumps({"click_ad_id": 1, "click_time": "2026-01-01 10:04:00"}) + "\n")
f.write(json.dumps({"click_ad_id": 2, "click_time": "2026-01-01 11:30:00"}) + "\n")
imps = (spark.readStream.schema("ad_id INT, imp_time TIMESTAMP").json(imp_dir)
.withWatermark("imp_time", "1 hour"))
clicks = (spark.readStream.schema("click_ad_id INT, click_time TIMESTAMP").json(clk_dir)
.withWatermark("click_time", "2 hours"))
joined = imps.join(clicks, F.expr("""
click_ad_id = ad_id AND
click_time BETWEEN imp_time AND imp_time + INTERVAL 15 MINUTES"""))
q = joined.writeStream.format("memory").queryName("attributed").trigger(availableNow=True).start()
q.awaitTermination()
spark.sql("SELECT * FROM attributed").show()
# +-----+-------------------+-----------+-------------------+
# |ad_id| imp_time|click_ad_id| click_time|
# +-----+-------------------+-----------+-------------------+
# | 1|2026-01-01 10:00:00| 1|2026-01-01 10:04:00|
# +-----+-------------------+-----------+-------------------+
spark.stop()
Ad 2’s click came 90 minutes after its impression, outside the 15-minute window, so it does not match.
Good to know: the time bound in the condition is what matters operationally.
It tells Spark that an impression older than the clicks’ watermark minus 15
minutes can be evicted from state. Without it, state grows forever. Outer joins
(which 2.3 added too) emit their unmatched rows only when that eviction happens,
so a left outer join’s NULL rows arrive late by design.
Pandas UDFs (SPARK-22216)
They changed the economics of Python on Spark.
How it works. A plain Python UDF pickles rows to a Python worker and back in
small batches, and calls your function once per row. A pandas UDF uses Apache
Arrow, a columnar in-memory format both the JVM and Python understand:
Spark sends a batch of rows (10,000 by default,
spark.sql.execution.arrow.maxRecordsPerBatch) as Arrow columns, pandas wraps
them without copying, your function runs once per batch on whole pd.Series
using vectorized pandas and NumPy code, and the result comes back the same way.
Same language, different order of magnitude of overhead.
The 2.3 API took a PandasUDFType argument. 3.0 replaced it with Python type
hints, and that is the form to write today. A Series -> Series function is a
scalar UDF; applyInPandas hands you a whole group as one pandas DataFrame,
which is how you run arbitrary pandas code per key:
import pandas as pd
from pyspark.sql import SparkSession
from pyspark.sql.functions import pandas_udf
spark = SparkSession.builder.appName("pudf").getOrCreate()
@pandas_udf("double")
def with_gst(amount: pd.Series) -> pd.Series:
return amount * 1.18
df = spark.createDataFrame([(1, 100.0), (2, 250.0)], "id INT, amount DOUBLE")
df.select("id", with_gst("amount").alias("gross")).show()
# +---+-----+
# | id|gross|
# +---+-----+
# | 1|118.0|
# | 2|295.0|
# +---+-----+
def normalise(pdf: pd.DataFrame) -> pd.DataFrame:
pdf["share"] = pdf["amount"] / pdf["amount"].sum()
return pdf
sales = spark.createDataFrame([("pune", 30.0), ("pune", 10.0), ("goa", 5.0)], "city STRING, amount DOUBLE")
sales.groupBy("city").applyInPandas(normalise, "city STRING, amount DOUBLE, share DOUBLE").orderBy("city", "amount").show()
# +----+------+-----+
# |city|amount|share|
# +----+------+-----+
# | goa| 5.0| 1.0|
# |pune| 10.0| 0.25|
# |pune| 30.0| 0.75|
# +----+------+-----+
spark.stop()
Good to know: a scalar pandas UDF must return a Series of the same length as
its input, and cannot assume anything about how rows are split into batches.
With applyInPandas, each group is loaded into a single Python worker’s memory
as one pandas DataFrame, so a skewed key becomes an out-of-memory error in
Python rather than a spill in the JVM. The same Arrow path also speeds up
toPandas() and createDataFrame(pandas_df) when
spark.sql.execution.arrow.pyspark.enabled is on, the default from 4.2.
Histograms for the cost-based optimizer (SPARK-21975)
How it works. The 2.2 CBO estimated a filter’s selectivity assuming values are
spread evenly between min and max. On skewed data that is badly wrong. With
spark.sql.statistics.histogram.enabled=true, ANALYZE TABLE ... FOR COLUMNS
also builds an equi-height histogram: the column’s values are split into
bins (254 by default) that each hold the same number of rows, so a value that
fills many bins is known to be common:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("hist23").getOrCreate()
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")
spark.range(0, 10_000).selectExpr("id", "CASE WHEN id < 9000 THEN 1 ELSE id END AS skewed").write.mode("overwrite").saveAsTable("h23")
spark.sql("ANALYZE TABLE h23 COMPUTE STATISTICS FOR COLUMNS skewed")
for r in spark.sql("DESCRIBE EXTENDED h23 skewed").collect():
if r[0] in ("distinct_count", "histogram") or r[0].startswith("bin_0") or r[0] == "bin_253":
print(r[0], "|", r[1])
# distinct_count | 990
# histogram | height: 39.37007874015748, num_of_bins: 254
# bin_0 | lower_bound: 1.0, upper_bound: 1.0, distinct_count: 1
# bin_253 | lower_bound: 9960.0, upper_bound: 9999.0, distinct_count: 39
spark.sql("DROP TABLE h23")
spark.stop()
With 90% of rows holding the value 1, the histogram’s first bin spans exactly
one value, 1.0 to 1.0, while the last bin covers forty values. A filter on
skewed = 1 is now estimated at around nine-tenths of the table instead of one
row in a thousand.
Good to know: distinct_count is 990 against a true 1,001. Column statistics
use a HyperLogLog sketch, so they are estimates, and they go stale the moment the
table changes. Histograms also make ANALYZE TABLE noticeably slower, since
they need an extra pass to compute percentiles.
Shuffle-plus-repartition correctness fix (SPARK-23207)
A correctness fix, not a feature: certain combinations of shuffle and repartition could return wrong results when tasks were retried.
How it works. repartition(n) without columns distributes rows round-robin,
so which output partition a row goes to depends on the order in which a task
sees its input. If that input came from a shuffle, its order is not
deterministic: a retried task can receive the same rows in a different order,
send them to different partitions, and some rows end up duplicated and others
lost when retried and non-retried output are combined. The fix sorts each
task’s rows locally before the round-robin assignment
(spark.sql.execution.sortBeforeRepartition, on by default), so a retry makes
the same assignment. The DataFrame case was fixed here in 2.3; the RDD case
followed in 2.4 as SPARK-23243.
Good to know: if you are on anything older and you repartition after a shuffle, that is your reason to move, and it outranks every feature in this post. The same class of problem, non-deterministic input to a retried stage, is what 4.1’s checksum-based stage retry addresses more generally.
Scala 2.10 removed (SPARK-19810)
Scala 2.11 became the minimum, and every Scala application and library had to be rebuilt for it.
Spark 2.4, November 2018
The last 2.x feature release, and the long-term home of an enormous amount of production Spark.
Barrier execution mode (SPARK-24374)
Distributed deep learning needs something Spark’s scheduler was built to avoid.
How it works. Normal Spark tasks are independent: they start whenever a slot
is free and are retried one at a time. Distributed training frameworks such as
Horovod or PyTorch’s DDP need the opposite: all workers start together, discover
each other, exchange gradients, and if one fails they all restart. A barrier
stage gangs its tasks: all of them are launched at once, barrier() is a
synchronisation point that every task must reach, getTaskInfos() tells each
task where its peers are, and any failure restarts the whole stage:
from pyspark import BarrierTaskContext
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[4]").appName("barrier24").getOrCreate()
def train(it):
ctx = BarrierTaskContext.get()
ctx.barrier() # every task waits here until all have started
peers = len(ctx.getTaskInfos())
yield (ctx.partitionId(), peers)
print(sorted(spark.sparkContext.parallelize(range(8), 4).barrier().mapPartitions(train).collect()))
# [(0, 4), (1, 4), (2, 4), (3, 4)]
spark.stop()
Every task saw all four peers.
Good to know: a barrier stage needs a free slot for every task at once. Spark
checks that the cluster has enough slots before launching it and, if not,
fails the job after a few retries instead of starting some tasks early. With
four partitions on local[2] the job above would be rejected. 3.5’s
TorchDistributor is built on this mode.
Higher-order functions (SPARK-23899)
The feature most 2.4 users never adopted and should have.
How it works. Before them, transforming an array column meant either
exploding it into rows and regrouping, which is a shuffle, or writing a UDF,
which leaves the engine. A higher-order function takes a lambda expression,
x -> x * 2, which Catalyst compiles like any other expression and applies to
each element inside the row, with no shuffle and no Python:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("hof").getOrCreate()
spark.sql("SELECT transform(array(1, 2, 3), x -> x * 2) AS doubled").show()
# +---------+
# | doubled|
# +---------+
# |[2, 4, 6]|
# +---------+
spark.sql("""
SELECT filter(array(1, 5, 12, 20), x -> x > 4) AS kept,
aggregate(array(1, 2, 3), 0, (acc, x) -> acc + x) AS total,
exists(array(1, 2, 3), x -> x = 2) AS has_two
""").show()
# +-----------+-----+-------+
# | kept|total|has_two|
# +-----------+-----+-------+
# |[5, 12, 20]| 6| true|
# +-----------+-----+-------+
spark.stop()
2.4 shipped transform, filter, aggregate, exists and zip_with; 3.0
added forall and the map versions, transform_keys, transform_values and
map_filter.
Good to know: the same functions exist in the Python DataFrame API since 3.1,
as F.transform, F.filter, F.aggregate and F.exists, taking Python lambdas
that are translated to expressions, not run in Python. The aggregate
accumulator’s type is fixed by its initial value, so aggregate(arr, 0, ...)
over doubles fails; start from 0D or CAST(0 AS DOUBLE).
Built-in Avro data source (SPARK-24768)
Read and write Avro files, and encode or decode Avro binary columns, which is
how Avro-encoded Kafka messages are handled. It ships with Spark but as a
separate module, spark-avro:
import tempfile
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.avro.functions import to_avro, from_avro
spark = SparkSession.builder.appName("avro24").getOrCreate()
path = tempfile.mkdtemp() + "/orders_avro"
df = spark.createDataFrame([(1, "pune"), (2, "goa")], "id INT, city STRING")
df.write.format("avro").save(path)
spark.read.format("avro").load(path).orderBy("id").show()
schema = """{"type": "record", "name": "r", "fields": [
{"name": "id", "type": ["int", "null"]},
{"name": "city", "type": ["string", "null"]}]}"""
(df.select(to_avro(F.struct("id", "city")).alias("bin")) # e.g. a Kafka message value
.select(from_avro("bin", schema).alias("rec"))
.select("rec.*").orderBy("id").show())
# +---+----+
# | id|city|
# +---+----+
# | 1|pune|
# | 2| goa|
# +---+----+
# +---+----+
# | id|city|
# +---+----+
# | 1|pune|
# | 2| goa|
# +---+----+
spark.stop()
spark-submit --packages org.apache.spark:spark-avro_2.13:4.1.3 app.py
How it works. Avro is a row-oriented binary format whose files carry their schema in the header, which suits record-at-a-time systems such as Kafka and data ingestion better than analytics. Values are written in schema order with no field names, so decoding needs the exact schema the writer used.
Good to know: that last point bites. The reader schema must match the
writer’s exactly, and a mismatch does not raise an error. My first version
declared "type": "int", but to_avro writes nullable columns as
["int", "null"] unions, and from_avro quietly decoded both rows as 0 and an
empty string. Also, messages written with a Confluent Schema Registry
serializer carry a 5-byte header in front of the Avro bytes, which has to be
stripped first.
Native ORC reader by default (SPARK-23456)
spark.sql.orc.impl defaults to native, the vectorized reader from 2.3, and
ORC tables created through Hive are read with it too when
spark.sql.hive.convertMetastoreOrc is on (its default from 2.4).
EXCEPT ALL and INTERSECT ALL (SPARK-21274)
Set operations that respect duplicates, as the SQL standard defines them:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("setops24").getOrCreate()
spark.sql("""
SELECT * FROM VALUES (1), (1), (1), (2) AS a(x)
EXCEPT ALL
SELECT * FROM VALUES (1), (2) AS b(x)
""").show()
spark.sql("""
SELECT * FROM VALUES (1), (1), (2) AS a(x)
INTERSECT ALL
SELECT * FROM VALUES (1), (1), (1) AS b(x)
""").show()
# +---+
# | x|
# +---+
# | 1|
# | 1|
# +---+
# +---+
# | x|
# +---+
# | 1|
# | 1|
# +---+
spark.stop()
How it works. EXCEPT ALL removes one matching row for each row on the right:
three 1s minus one 1 leaves two, and the single 2 is cancelled out. INTERSECT
ALL keeps each value as many times as it appears on both sides: two 1s on
the left, three on the right, so two. Spark implements them by counting each
value on both sides and replicating the difference or the minimum.
Good to know: in the DataFrame API they are exceptAll and intersectAll.
Plain EXCEPT and INTERSECT de-duplicate, which is usually not what a
reconciliation query between two copies of a table wants.
PIVOT syntax (SPARK-24035)
Turns row values into columns:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("pivot24").getOrCreate()
spark.createDataFrame(
[("pune", "Q1", 100), ("pune", "Q2", 150), ("goa", "Q1", 40), ("goa", "Q2", 60)],
"city STRING, quarter STRING, amount INT").createOrReplaceTempView("sales")
spark.sql("""
SELECT * FROM sales
PIVOT (sum(amount) FOR quarter IN ('Q1' AS q1, 'Q2' AS q2))
ORDER BY city
""").show()
# +----+---+---+
# |city| q1| q2|
# +----+---+---+
# | goa| 40| 60|
# |pune|100|150|
# +----+---+---+
spark.stop()
How it works. A pivot is a GROUP BY on the columns that are not mentioned
(city), with one conditional aggregate per listed value, roughly
sum(CASE WHEN quarter = 'Q1' THEN amount END) AS q1.
Good to know: list the pivot values explicitly, as above. Without an IN
list, the DataFrame pivot() has to run an extra job to find the distinct
values before it can plan the query, and the column set changes whenever the
data does. There is a limit of 10,000 distinct values
(spark.sql.pivotMaxValues).
Nested schema pruning (SPARK-4502)
Selecting one field of a struct reads only that field’s column from Parquet, not the whole struct:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("nested24").getOrCreate()
path = tempfile.mkdtemp()
spark.sql("SELECT id, named_struct('name', 'u' || id, 'address', named_struct('city', 'pune', 'zip', '411001')) AS user FROM range(10)") \
.write.mode("overwrite").parquet(path)
df = spark.read.parquet(path).select("user.address.city")
df.collect()
scan = [l for l in df._jdf.queryExecution().executedPlan().toString().splitlines() if "FileScan" in l][0]
print(scan[scan.index("ReadSchema"):].strip())
# ReadSchema: struct<user:struct<address:struct<city:string>>>
spark.stop()
The ReadSchema contains only user.address.city: id, user.name and
user.address.zip are never read from disk.
How it works. Parquet stores each leaf field of a nested structure as its own
column, so user.address.city and user.address.zip are separate columns on
disk. Column pruning in 1.x worked at the top level only, so selecting one field
of user read all of user’s leaves. Nested pruning pushes the exact field
paths into the scan.
Good to know: it was opt-in here and has been on by default since 3.0
(spark.sql.optimizer.nestedSchemaPruning.enabled). For wide nested event data
it is often the single biggest I/O saving available.
Bucket pruning (SPARK-23803)
A filter on a bucketed table’s bucket column reads only the matching bucket’s files, the way partition pruning skips directories:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("bprune24").getOrCreate()
spark.range(0, 10_000).selectExpr("id AS cust_id", "id % 100 AS v") \
.write.bucketBy(8, "cust_id").mode("overwrite").saveAsTable("cust_b")
def scan_info(auto):
spark.conf.set("spark.sql.sources.bucketing.autoBucketedScan.enabled", auto)
df = spark.table("cust_b").where("cust_id = 4242")
df.collect()
scan = [l for l in df._jdf.queryExecution().executedPlan().toString().splitlines() if "FileScan" in l][0]
bucketed = scan[scan.index("Bucketed:"):scan.index(", DataFilters")]
selected = scan[scan.index("SelectedBucketsCount"):].strip() if "SelectedBucketsCount" in scan else "-"
print(f"auto={auto:5s} {bucketed} | {selected}")
scan_info("true")
scan_info("false")
# auto=true Bucketed: false (disabled by query planner) | -
# auto=false Bucketed: true | SelectedBucketsCount: 1 out of 8
spark.sql("DROP TABLE cust_b")
spark.stop()
How it works. Rows were assigned to buckets by hash(cust_id) mod 8 at write
time, so a filter cust_id = 4242 can compute the one bucket that can hold the
value and read only its files.
The output shows a behaviour from a later release. By default, the scan says
Bucketed: false (disabled by query planner) and no pruning happens. Since 3.1,
Spark decides per query whether a bucketed scan is worth it
(spark.sql.sources.bucketing.autoBucketedScan.enabled, on by default), and a
bucketed scan only pays off when a join or aggregation can use the bucket
layout; a plain filter does not qualify, so the planner reads the table as
ordinary files, and bucket pruning goes with the bucketed scan. With the rule
off, Spark reads 1 out of 8 buckets.
Good to know: that makes bucket pruning much less useful than it sounds on current versions, because it only happens when the query also does something the bucket layout helps with, or when you turn the rule off. Partitioning, or a table format’s data skipping, is the more reliable way to make point lookups cheap.
Blocks larger than 2 GB (SPARK-24296, SPARK-24307)
It explains a class of failure that simply stops happening after 2.4.
How it works. Spark’s network and storage layers held each block in a single
ByteBuffer, and ByteBuffer is indexed by int, so no block could exceed
2^31 bytes. Any shuffle block, cached partition or replicated block over 2 GB
failed with an overflow such as Size exceeds Integer.MAX_VALUE. 2.4 changed the
transfer and replication paths to stream large blocks in chunks.
Good to know: the fix removes a hard ceiling that used to dictate partition counts. Partitions that large are still slow and put the whole task at the mercy of one executor; aim for tens to low hundreds of megabytes each.
Scala 2.12, experimental (SPARK-14220)
Spark could be built for Scala 2.12. It became the default in 3.0, 3.2 added Scala 2.13, and 4.0 dropped 2.12 so that 2.13 is the only option.
Good to know: Scala minor versions are not binary compatible, which is why
every library JAR carries a _2.12 or _2.13 suffix and why each of these
moves forced the whole ecosystem of connectors to republish.
The 3.x line: the planner stops guessing
The 2.x engine made one plan from estimates and executed it. The estimates were frequently wrong, because nobody knows the selectivity of a filter over data nobody has read yet. The 3.x line’s theme is letting the engine revise itself, and separately, admitting that Python is the primary language. Read the 3.x features in three groups: 3.0 and 3.2 make planning adaptive; 3.1 to 3.4 make SQL stricter and errors structured; 3.4 and 3.5 split the client from the driver with Spark Connect.
Note the gap: 2.4.0 shipped in November 2018 and 3.0.0 in June 2020, nineteen months later. Major versions are where the project spends its compatibility budget, and it does not spend it often.
Spark 3.0, June 2020
Adaptive query execution, rebuilt (SPARK-31412)
How it works. A query is cut into query stages at every shuffle and broadcast. AQE runs the stages that have no dependencies first, and when a stage finishes it has something the planner never had before: the actual size of every shuffle partition it wrote. It then re-optimizes the rest of the plan with those numbers, and repeats after each stage. Four rules use them:
- Coalescing merges runs of small neighbouring shuffle partitions so each
reducer reads about
spark.sql.adaptive.advisoryPartitionSizeInBytes(64 MB). - Join conversion replaces a planned sort-merge join with a broadcast hash join when one side turns out to be small enough.
- Skew splitting cuts an oversized partition of a sort-merge join into several tasks, duplicating the matching partition of the other side.
- Local shuffle reads let a converted broadcast join read shuffle files on the same node instead of fetching them over the network.
A filter whose selectivity the planner cannot know shows the join conversion, and one rule that surprised me. Both tables are large and incompressible on disk, so Spark plans a sort-merge join; at runtime the filter on the dimension side leaves either 40 rows or 2,000:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("aqe30").config("spark.sql.adaptive.enabled", "true").getOrCreate()
p = tempfile.mkdtemp()
spark.range(0, 2_000_000).selectExpr("id AS k", "id % 100 AS v", "md5(CAST(id AS STRING)) AS pad").write.parquet(p + "/facts")
spark.range(0, 2_000_000).selectExpr("id AS k", "sha2(CAST(id AS STRING), 256) AS name").write.parquet(p + "/dims")
facts = spark.read.parquet(p + "/facts")
def run(modulus):
dims = spark.read.parquet(p + "/dims").where(f"k % {modulus} = 7") # selectivity unknowable at plan time
df = facts.join(dims, "k").groupBy("v").count()
df.collect()
plan = df._jdf.queryExecution().executedPlan().toString()
final, initial = plan.split("== Initial Plan ==")
pick = lambda s: [n for n in ["SortMergeJoin", "BroadcastHashJoin"] if n in s]
print(f"{dims.count():>5} dim rows initial: {pick(initial)} final: {pick(final)}")
run(50_000)
run(1_000)
# 40 dim rows initial: ['SortMergeJoin'] final: ['SortMergeJoin']
# 2000 dim rows initial: ['SortMergeJoin'] final: ['BroadcastHashJoin']
spark.stop()
With 2,000 surviving rows, AQE measured the dimension’s shuffle output, found it
far below the 10 MB threshold, and replaced the sort-merge join with a broadcast
hash join. With 40 rows, fewer and smaller, it kept the sort-merge join. The
reason is a rule that is easy to miss:
spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoin (default 0.2). A
side whose shuffle output leaves fewer than 20% of partitions non-empty is not
considered for a runtime broadcast, because a sort-merge join where most
partitions are empty is already cheap. Forty rows hashed into 200 partitions
fill at most 40 of them. My first attempt had a subtler problem: the tables’
Parquet files compressed to under 10 MB, so Spark planned a broadcast join
before AQE had anything to change. Plan-time estimates come from file sizes, so
compressible test data behaves nothing like production data.
Good to know: AQE can only change what comes after a shuffle it has
measured. It cannot change the first stage’s parallelism, cannot fix skew in an
aggregation (only in joins), and only splits a partition once it is larger than
both skewedPartitionThresholdInBytes (256 MB) and five times the median. What
it does and does not revisit is the subject of
Adaptive query execution, which also
has a verified walkthrough of the skew split, including the
threshold that often stops it firing.
Dynamic partition pruning (SPARK-11150)
How it works. A star-schema query filtering a small dimension table cannot prune fact-table partitions at plan time, because the filter is on the other side of a join. DPP adds a runtime filter to the fact scan: before scanning the fact table, Spark computes the set of join keys that survive the dimension filter and prunes every partition whose key is not in the set. When the dimension side is already being broadcast, DPP reuses that broadcast as the key set, so the pruning is almost free. The effect is measurable from the scan’s own metrics. Thirty daily partitions, a dimension marking eight of those days as weekend days, and a query filtering on the dimension:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("dpp30").getOrCreate()
spark.conf.set("spark.sql.adaptive.enabled", "false") # static plan, so the scan's metrics are easy to read
path = tempfile.mkdtemp()
(spark.range(0, 100_000)
.selectExpr("id AS order_id", "CAST(id % 30 AS INT) AS day_id", "id % 1000 AS amount")
.write.partitionBy("day_id").parquet(path + "/orders"))
(spark.createDataFrame([(i, "weekend" if i % 7 in (5, 6) else "weekday") for i in range(30)],
"day_id INT, kind STRING")
.write.parquet(path + "/days"))
spark.read.parquet(path + "/orders").createOrReplaceTempView("orders")
spark.read.parquet(path + "/days").createOrReplaceTempView("days")
query = """
SELECT sum(amount) FROM orders o JOIN days d ON o.day_id = d.day_id
WHERE d.kind = 'weekend'
"""
def partitions_read(dpp):
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", str(dpp).lower())
df = spark.sql(query)
df.collect()
leaves = df._jdf.queryExecution().executedPlan().collectLeaves()
for i in range(leaves.size()):
node = leaves.apply(i)
if node.nodeName().startswith("Scan parquet") and "orders" in node.toString():
return node.metrics().apply("numPartitions").value()
print("partitions read, DPP off:", partitions_read(False))
# partitions read, DPP off: 30
print("partitions read, DPP on: ", partitions_read(True))
# partitions read, DPP on: 8
spark.stop()
Same result on 3.5.9. Eight is exactly the number of weekend days in the range, so the fact scan read nothing it did not need.
Good to know: DPP only prunes partitions, so the fact table must be
partitioned on the join key; for row-level pruning of unpartitioned tables see
the 3.3 bloom filters. My first version of this example registered the dimension
straight from createDataFrame, and the pruning filter came out as
dynamicpruningexpression(true), which means pruning was planned and then
thrown away. An in-memory local relation has no size statistics, so Spark
assumed it was huge and broadcast the fact table instead. DPP can only prune
the side that is being probed, so with the sides reversed there was nothing to
prune. DPP depends on the join strategy, and the join strategy depends on
statistics.
Join hints for every strategy
2.x only had BROADCAST; 3.0 added a hint for every physical join strategy,
which makes a join-strategy experiment a one-word change:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("hints30").getOrCreate()
spark.range(0, 100_000).createOrReplaceTempView("a")
spark.range(0, 100_000).createOrReplaceTempView("b")
for hint in ["BROADCAST(b)", "MERGE(b)", "SHUFFLE_HASH(b)", "SHUFFLE_REPLICATE_NL(b)"]:
df = spark.sql(f"SELECT /*+ {hint} */ * FROM a JOIN b ON a.id = b.id")
plan = df._jdf.queryExecution().executedPlan().toString()
node = next(l.strip(" +-*()0123456789").split(" ")[0] for l in plan.splitlines() if "Join" in l or "Cartesian" in l)
print(f"{hint:26s}{node}")
# BROADCAST(b) BroadcastHashJoin
# MERGE(b) SortMergeJoin
# SHUFFLE_HASH(b) ShuffledHashJoin
# SHUFFLE_REPLICATE_NL(b) CartesianProduct
spark.stop()
How it works. The four strategies differ in what moves and what must fit in memory. A broadcast hash join ships the small side to every executor and never shuffles the large one. A sort-merge join shuffles and sorts both sides and needs no side to fit in memory. A shuffled hash join shuffles both sides but skips the sort, building a hash table from each partition of the smaller side, which must fit. Shuffle-replicate nested loop is a Cartesian product, the only option for joins without an equality condition.
Good to know: a hint is a request, not an order. If the hinted strategy cannot execute the join, such as a broadcast of the preserved side of an outer join, Spark logs a warning and picks another. When both sides carry different hints, the precedence is broadcast, then merge, then shuffle hash, then nested loop.
Accelerator-aware scheduling (SPARK-24615)
Executors and tasks can request GPUs, or any named resource, and the scheduler places tasks only where the resource is free.
How it works. Each executor runs a discovery script at startup that prints
the addresses of the devices it owns, for example GPUs 0 and 1. The
scheduler then treats them like CPU cores: a task needing a quarter of a GPU is
only placed where a quarter is free, and it is told which device it was given:
--conf spark.executor.resource.gpu.amount=1 \
--conf spark.task.resource.gpu.amount=0.25 \
--conf spark.executor.resource.gpu.discoveryScript=/opt/spark/examples/src/main/scripts/getGpusResources.sh
Inside a task, TaskContext.get().resources()["gpu"].addresses tells your code
which device to use.
Good to know: 3.1 added stage-level scheduling, so one stage can ask for GPU
executors while the rest of the job uses cheaper CPU-only ones, through
ResourceProfileBuilder and rdd.withResources(...). The RAPIDS Accelerator
plugin uses this mechanism to run whole SQL plans on GPUs.
Pandas UDFs with Python type hints (SPARK-28264)
How it works. The kind of UDF is now inferred from the function’s type hints,
instead of a PandasUDFType argument that was easy to get wrong:
pd.Series -> pd.Seriesis a scalar UDF, called once per Arrow batch.Iterator[pd.Series] -> Iterator[pd.Series]is called once per task with an iterator of batches, so expensive setup, such as loading a model, happens once instead of once per batch.pd.Series -> float(any scalar) is a grouped aggregate, usable ingroupBy().agg()and over windows.
import pandas as pd
from typing import Iterator
from pyspark.sql import SparkSession
from pyspark.sql.functions import pandas_udf
spark = SparkSession.builder.appName("pudf30").getOrCreate()
@pandas_udf("double")
def median_amount(v: pd.Series) -> float: # Series -> scalar: a grouped aggregate
return float(v.median())
@pandas_udf("double")
def scaled(batches: Iterator[pd.Series]) -> Iterator[pd.Series]:
factor = 1.18 # expensive setup (e.g. loading a model) runs once per task
for b in batches:
yield b * factor
df = spark.createDataFrame([("pune", 10.0), ("pune", 30.0), ("pune", 50.0), ("goa", 8.0)], "city STRING, amount DOUBLE")
df.groupBy("city").agg(median_amount("amount").alias("median")).orderBy("city").show()
df.select("city", scaled("amount").alias("gross")).orderBy("city", "gross").show()
# +----+------+
# |city|median|
# +----+------+
# | goa| 8.0|
# |pune| 30.0|
# +----+------+
# +----+------------------+
# |city| gross|
# +----+------------------+
# | goa| 9.44|
# |pune|11.799999999999999|
# |pune| 35.4|
# |pune| 59.0|
# +----+------------------+
spark.stop()
Good to know: the 2.3 section shows the scalar form and applyInPandas. A
grouped-aggregate pandas UDF loads each group into memory in one piece, like
applyInPandas, and cannot be partially aggregated before the shuffle the way
built-in aggregates are, so on large groups a built-in such as
percentile_approx is much cheaper than a pandas median.
Catalog plugin API (SPARK-31121)
Multiple named catalogs in one session, each backed by a plug-in class. It is
the basis of every modern table format connector, and why a table can be
addressed as catalog.database.table.
How it works. A catalog is a class implementing the CatalogPlugin family of
interfaces (TableCatalog, SupportsNamespaces, and optionally views,
functions and procedures), registered under a name with configuration:
--conf spark.sql.catalog.lake=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.lake.type=rest \
--conf spark.sql.catalog.lake.uri=http://polaris:8181/api/catalog
SELECT * FROM lake.sales.orders;
USE lake.sales; -- change the current catalog and namespace
SHOW TABLES;
When a query names lake.sales.orders, Spark asks the lake plug-in to load
table orders in namespace sales, and the table object it returns declares
what it supports: batch reads, streaming, row-level deletes, and so on. Every
spark.sql.catalog.lake.* setting is passed to the plug-in as its options.
Good to know: the built-in spark_catalog is always there, so unqualified
names keep working. The
Iceberg with Spark Connect
post registers a catalog the same way, with a Hadoop-type warehouse instead of
a REST service.
Proleptic Gregorian calendar (SPARK-26651)
The sleeper.
How it works. The Gregorian calendar replaced the Julian one in October 1582,
and the switch skipped ten days: the day after 4 October 1582 was 15 October.
Spark 2.x used a hybrid calendar with that jump, like java.util.Date. Spark
3.0 switched to the proleptic Gregorian calendar, which projects Gregorian
rules backwards forever, matching Java 8’s java.time, the SQL standard and most
other engines. On 3.5.9 and 4.1.3 alike:
SELECT date_add(DATE'1582-10-04', 1) AS day_after,
datediff(DATE'1582-10-15', DATE'1582-10-04') AS gap;
-- 1582-10-05 11
Every date exists and the arithmetic is ordinary. Under the hybrid calendar of Spark 2.x the answers were the historical ones: the day after 4 October 1582 was the 15th, the two dates were one day apart, and a date inside the gap, such as 10 October 1582, did not exist and was silently moved forward.
Good to know: the risk is data written by 2.x. The same stored day number
means a different date under the two calendars, so dates before 1582, and
timestamps before 1900 in some formats, can shift when read by 3.x. Spark
refuses to guess: reading ambiguous old values from files written by 2.x can
raise a SparkUpgradeException (error class
INCONSISTENT_BEHAVIOR_CROSS_VERSION) until you choose, with
spark.sql.parquet.datetimeRebaseModeInRead set to LEGACY (rebase from the
old calendar) or CORRECTED (read as-is). If you store historical dates, or
placeholder dates such as 0001-01-01, test before upgrading rather than after.
That read setting only matters for files that do not say how they were
written. I tried to reproduce a shift on 3.5.9 and 4.1.3 by writing
DATE'1000-01-01' with datetimeRebaseModeInWrite=LEGACY and reading it back
under both read modes, and got 1000-01-01 both times. The reason is in the
Parquet footer: a LEGACY write adds an org.apache.spark.legacyDateTime key,
and a CORRECTED write does not, so a 3.x or 4.x reader knows exactly how to
interpret the file whatever the read mode says. The read mode is a decision for
files from 2.x and from other writers, which carry no such marker.
Structured Streaming UI (SPARK-29543)
How it works. A dedicated tab lists every streaming query, active and
finished, and charts per batch: input rate, processing rate, input rows, batch
duration broken into its phases, and aggregated state rows and memory. The
numbers come from the same StreamingQueryProgress events that
query.lastProgress returns.
Good to know: the one chart to watch is input rate against processing rate. If input stays above processing, batches take longer and longer and the query falls behind; no amount of waiting fixes it.
Java 11 and Hadoop 3 (SPARK-24417, SPARK-23534)
Spark ran on Java 11 and built against Hadoop 3, which brought the S3A connector improvements, including the S3A committers that make writing to object storage safe and fast.
Good to know: the release notes state that Spark 3.0 is roughly twice as fast as 2.4 on a 30 TB TPC-DS benchmark. That is the project’s own measurement on its own hardware. Your workload is not TPC-DS.
Spark 3.1, March 2021
There is no Spark 3.1.0. The line starts at 3.1.1; the release notes for 3.1.1 describe it as the second release of the 3.x line.
Kubernetes declared generally available (SPARK-33005)
Three years after its experimental introduction in 2.3, which is the clearest example of the project’s pace on deployment features.
How it works. GA brought the pieces production needed: pod templates, so anything Spark’s settings do not cover (tolerations, sidecars, init containers, volumes) can be set in ordinary Kubernetes YAML; executor loss detection and cleanup; and support for dynamic allocation through shuffle tracking:
--conf spark.kubernetes.driver.podTemplateFile=driver-template.yaml \
--conf spark.kubernetes.executor.podTemplateFile=executor-template.yaml \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.dynamicAllocation.shuffleTracking.enabled=true
# executor-template.yaml
apiVersion: v1
kind: Pod
spec:
tolerations:
- key: spot
operator: Exists
nodeSelector:
workload: spark
Good to know: Spark’s own settings win over the template where both set the same field, such as memory and cores. Executors on Kubernetes have no external shuffle service by default, which is why shuffle tracking, or 3.1’s decommissioning below, matters there.
Project Zen
A push to make PySpark feel like a Python library rather than a Scala library with Python bindings.
How it works. Three strands: type hints throughout the API, which IDEs and
mypy understand; a rewritten documentation site with Python-first examples;
and first-class dependency management, shipping a packed conda or virtualenv
environment to executors so that the Python packages there match the driver’s:
python -m venv env && . env/bin/activate && pip install pandas pyarrow venv-pack
venv-pack -o env.tar.gz
spark-submit --archives env.tar.gz#environment \
--conf spark.pyspark.python=./environment/bin/python app.py
Good to know: --archives unpacks env.tar.gz in each executor’s working
directory under the name after #, which is why the Python path is relative.
This is the reliable way to get a specific pandas or NumPy version onto a cluster
you do not administer.
CHAR and VARCHAR types (SPARK-33480)
Real length-checked types rather than aliases for STRING. CHAR(n) pads on
read and VARCHAR(n) rejects values that are too long on write:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("char31").getOrCreate()
spark.sql("CREATE TABLE codes (c CHAR(5), v VARCHAR(3)) USING parquet")
spark.sql("INSERT INTO codes VALUES ('ab', 'xyz')")
spark.sql("SELECT concat('[', c, ']') AS c, length(c) AS len FROM codes").show()
# +-------+---+
# | c|len|
# +-------+---+
# |[ab ]| 5|
# +-------+---+
try:
spark.sql("INSERT INTO codes VALUES ('ab', 'toolong')")
except Exception as e:
print(str(e).split("\n")[0][:110])
# [EXCEED_LIMIT_LENGTH] Exceeds char/varchar type length limitation: 3. SQLSTATE: 54006
spark.sql("DROP TABLE codes")
spark.stop()
How it works. On disk both are stored as ordinary strings. The type is kept in
the table’s metadata, and Spark adds a length check on write, and padding on
read for CHAR. Data frames created in code never have these types; they only
come from table definitions.
Good to know: the padding is the part that surprises people migrating from
engines where CHAR comparisons ignore trailing spaces. length(c) is 5, not
- Unless you need the check,
STRINGis simpler.
ANSI mode raises errors instead of returning null (SPARK-33275)
This opens an argument that runs across the rest of the 3.x line and is settled in 4.0. Silently returning null on overflow or a bad cast is convenient and hides data corruption. ANSI mode raises instead. In 3.1 it is opt-in:
spark.conf.set("spark.sql.ansi.enabled", "true")
How it works. With the flag on, the affected expressions (arithmetic, casts,
array and map indexing, some date functions) generate code that checks for
overflow and invalid input and throws, instead of producing NULL or wrapping
around. The generated-code example in the 2.0 section shows the effect on 4.x:
id * 3 compiles to multiplyExact, which throws on overflow.
Good to know: the 4.0 section runs the same statement on 2.4, 3.5 and 4.1 to
show what changes, and lists the try_ functions that give the old behaviour
back where you want it.
Explicit cast rules under ANSI mode (SPARK-33354)
How it works. ANSI mode also changes which casts are allowed at all. A cast that can never make sense, such as a timestamp to a boolean, is rejected when the query is analysed; a cast that can fail for some values, such as a string to a date, is allowed but throws at runtime on a bad value:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("cast31").getOrCreate()
for q in ["SELECT CAST(TIMESTAMP'2026-01-01 00:00:00' AS BOOLEAN)", "SELECT CAST('2026-13-45' AS DATE)"]:
try:
spark.sql(q).collect(); print("ok", q)
except Exception as e:
print(str(e).split("\n")[0][:120])
spark.conf.set("spark.sql.ansi.enabled", "false")
print(spark.sql("SELECT CAST('2026-13-45' AS DATE) AS d").first().d)
# [DATATYPE_MISMATCH.CAST_WITH_CONF_SUGGESTION] Cannot resolve "CAST(TIMESTAMP '2026-01-01 00:00:00' AS BOOLEAN)" due to d...
# [CAST_INVALID_INPUT] The value '2026-13-45' of the type "STRING" cannot be cast to "DATE" because it is malformed. ...
# None
spark.stop()
The timestamp-to-boolean cast is rejected before the query runs, with an error
class, DATATYPE_MISMATCH.CAST_WITH_CONF_SUGGESTION, whose message suggests
turning ANSI mode off. The invalid date is only discovered when the value is
read, so it raises CAST_INVALID_INPUT at runtime. With ANSI off, the same cast
quietly returns None, which is how a month number of 13 becomes a null date
nobody notices.
Good to know: when a cast can legitimately fail, as with dates typed by
people, use try_cast or try_to_date to get NULL for the bad values
deliberately, and count them.
Shuffled hash join improvements, including full outer joins (SPARK-32461)
Shuffled hash join learned full outer joins and code generation:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("shj31").getOrCreate()
spark.createDataFrame([(1, "a"), (2, "b")], "id INT, l STRING").createOrReplaceTempView("l")
spark.createDataFrame([(2, "x"), (3, "y")], "id INT, r STRING").createOrReplaceTempView("r")
df = spark.sql("SELECT /*+ SHUFFLE_HASH(r) */ l.id AS lid, l, r.id AS rid, r FROM l FULL OUTER JOIN r ON l.id = r.id")
df.orderBy("lid", "rid").show()
print([n for n in ["ShuffledHashJoin", "SortMergeJoin"] if n in df._jdf.queryExecution().executedPlan().toString()])
# +----+----+----+----+
# | lid| l| rid| r|
# +----+----+----+----+
# |NULL|NULL| 3| y|
# | 1| a|NULL|NULL|
# | 2| b| 2| x|
# +----+----+----+----+
# ['ShuffledHashJoin']
spark.stop()
How it works. A full outer join must also emit the rows of the build side that never matched. Shuffled hash join now tracks which hash-table entries were matched while streaming the other side, then emits the unmatched ones at the end of each partition. It avoids the sort that sort-merge join pays, at the cost of holding one side’s partition in memory.
Good to know: AQE can also pick shuffled hash join by itself: when every
partition of one side is below
spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold, it replaces a
sort-merge join with a shuffled hash join at runtime.
Node decommissioning, experimental (SPARK-20624)
How it works. When a node is about to go away, such as a spot instance receiving its termination notice or a Kubernetes node being drained, Spark marks its executors as decommissioning: no new tasks are scheduled on them, running tasks are allowed to finish, and their cached RDD blocks and shuffle files are migrated to other executors, so the work does not have to be recomputed:
--conf spark.decommission.enabled=true \
--conf spark.storage.decommission.enabled=true \
--conf spark.storage.decommission.shuffleBlocks.enabled=true \
--conf spark.storage.decommission.rddBlocks.enabled=true
Good to know: migration needs somewhere to put the blocks; with
spark.storage.decommission.fallbackStorage.path set to an object-store path,
blocks can survive even when no other executor has room. It pairs naturally with
running executors on spot capacity.
Spark 3.2, October 2021
pandas API on Spark (SPARK-34849)
The former Koalas project, merged in as pyspark.pandas: the pandas API,
executed by Spark, so pandas code can scale past one machine’s memory:
import pyspark.pandas as ps
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ps32").getOrCreate()
psdf = ps.DataFrame({"city": ["pune", "pune", "goa"], "amount": [10, 20, 30]})
print(psdf.groupby("city").amount.sum().sort_index().to_dict())
print(type(psdf.to_spark()).__name__)
print(psdf.describe().loc["mean", "amount"])
# {'goa': 30, 'pune': 30}
# DataFrame
# 20.0
spark.stop()
How it works. A pandas-on-Spark DataFrame is a Spark DataFrame plus an
index, the thing pandas has and Spark does not. Each pandas method is
translated into Spark operations over the data and the index columns, and
results are only brought to the driver when you call something like
to_pandas(). psdf.to_spark() and sdf.pandas_api() convert between the two
worlds without copying data.
Good to know: it needs pandas installed on the cluster; on a clean
apache/spark image the import fails with
[PACKAGE_NOT_INSTALLED] Pandas >= 2.2.0 must be installed. The index is also
the main cost: operations that need a sequential index, such as positional
iloc, force a global ordering, which is expensive. Set the
compute.default_index_type option to distributed when you do not need a
sequential index, and avoid row-by-row apply, which is a Python UDF underneath.
Adaptive query execution enabled by default (SPARK-33679)
The most important line in the 3.x changelog.
How it works. spark.sql.adaptive.enabled became true, together with
partition coalescing and skew-join handling. Every query with a shuffle now runs
stage by stage with re-planning in between, as described in the 3.0 entry.
Good to know: if you upgrade from 3.1 to 3.2 with no config changes, your
query plans change. They usually improve. They do change: the number of tasks
after a shuffle is no longer spark.sql.shuffle.partitions but whatever
coalescing decided, and joins can switch strategy between runs as data grows.
Read the final plan in the SQL tab, not the initial one, when debugging.
Push-based shuffle (SPARK-30602)
How it works. In a normal shuffle, each reducer fetches a small block from every map task, so M map tasks and R reducers mean M × R small random reads, which is slow on disks and busy networks. With push-based shuffle, map tasks also push their blocks to remote shuffle services, which merge all blocks for the same reducer into one file as they arrive. Reducers then read a few large sequential files, falling back to the original blocks for anything that was not merged in time. It was built at LinkedIn for very large shuffles, and needs YARN with the external shuffle service:
--conf spark.shuffle.push.enabled=true \
--conf spark.shuffle.service.enabled=true
Good to know: it pays off for shuffles with many small blocks (large M × R); for small jobs the extra pushing is overhead. Kubernetes deployments usually reach for a remote shuffle service such as Apache Celeborn instead.
RocksDB state store (SPARK-34198)
The fix for the other half of the streaming state problem.
How it works. The default state store keeps each partition’s state as a hash map in the executor’s JVM heap and writes changes to the checkpoint as delta files. A large stateful stream therefore becomes a garbage collection problem. The RocksDB provider keeps state in a RocksDB instance on the executor’s local disk, with its own off-heap cache, and uploads its files to the checkpoint on each commit. State can then grow far past what the heap can hold, trading some read latency:
spark.conf.set(
"spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider",
)
Good to know: the provider is recorded in the checkpoint, so the setting has
to be made before the query first runs; switching an existing query needs a new
checkpoint. spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled
(3.5) uploads only changes on each commit, which cuts checkpoint time for large
state. The 4.0 State API v2 requires this provider.
Session windows (SPARK-10816)
A tumbling window has a fixed length; a session window stays open while events keep arriving within a gap and closes after a quiet period, so its length depends on the data. It works on batch DataFrames too:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("sess32").getOrCreate()
clicks = spark.createDataFrame(
[("u1", "2026-01-01 10:00:00"), ("u1", "2026-01-01 10:03:00"),
("u1", "2026-01-01 10:30:00"), ("u2", "2026-01-01 10:01:00")],
"user STRING, ts STRING").withColumn("ts", F.to_timestamp("ts"))
(clicks.groupBy("user", F.session_window("ts", "10 minutes"))
.count()
.select("user", "session_window.start", "session_window.end", "count")
.orderBy("user", "start")
.show(truncate=False))
# +----+-------------------+-------------------+-----+
# |user|start |end |count|
# +----+-------------------+-------------------+-----+
# |u1 |2026-01-01 10:00:00|2026-01-01 10:13:00|2 |
# |u1 |2026-01-01 10:30:00|2026-01-01 10:40:00|1 |
# |u2 |2026-01-01 10:01:00|2026-01-01 10:11:00|1 |
# +----+-------------------+-------------------+-----+
spark.stop()
User u1’s first session ends at 10:13: the last event at 10:03 plus the 10-minute gap. The 10:30 click starts a new session because nothing arrived between 10:03 and 10:13.
How it works. Each event starts as its own window from its timestamp to timestamp plus the gap; Spark sorts events by key and time and merges windows that overlap. The gap can be an expression instead of a constant, so each key can have its own timeout:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("gap32").getOrCreate()
clicks = spark.createDataFrame(
[("mobile", "2026-01-01 10:00:00"), ("mobile", "2026-01-01 10:04:00"),
("web", "2026-01-01 10:00:00"), ("web", "2026-01-01 10:04:00")],
"device STRING, ts STRING").withColumn("ts", F.to_timestamp("ts"))
gap = F.when(F.col("device") == "mobile", "10 minutes").otherwise("2 minutes")
(clicks.groupBy("device", F.session_window("ts", gap)).count()
.select("device", "session_window.start", "session_window.end", "count")
.orderBy("device", "start").show(truncate=False))
# +------+-------------------+-------------------+-----+
# |device|start |end |count|
# +------+-------------------+-------------------+-----+
# |mobile|2026-01-01 10:00:00|2026-01-01 10:14:00|2 |
# |web |2026-01-01 10:00:00|2026-01-01 10:02:00|1 |
# |web |2026-01-01 10:04:00|2026-01-01 10:06:00|1 |
# +------+-------------------+-------------------+-----+
spark.stop()
The same two clicks, four minutes apart, are one session on mobile and two on the web.
Good to know: in a streaming query, session windows need a watermark, and a
session is emitted in append mode only once the watermark passes its end,
which is the last event plus the gap. Long gaps mean late results.
ANSI mode declared generally available (SPARK-35030)
Still off by default, but supported for production use, with the full set of
try_ functions available for the places where you want nulls back.
ANSI INTERVAL types (SPARK-27790)
The old catch-all CalendarInterval was split into two comparable, storable
types: year-month and day-time. Subtracting two dates now yields a typed
interval:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("int32").getOrCreate()
spark.sql("""
SELECT INTERVAL '1-2' YEAR TO MONTH AS ym,
INTERVAL '3 04:05:06' DAY TO SECOND AS dt,
typeof(DATE'2026-03-31' - DATE'2026-01-01') AS diff_type
""").show(truncate=False)
# +----------------------------+-----------------------------------+------------+
# |ym |dt |diff_type |
# +----------------------------+-----------------------------------+------------+
# |INTERVAL '1-2' YEAR TO MONTH|INTERVAL '3 04:05:06' DAY TO SECOND|interval day|
# +----------------------------+-----------------------------------+------------+
spark.stop()
How it works. A year-month interval is stored as a number of months, and a day-time interval as a number of microseconds. Each is a single number, so they can be compared, sorted, summed and written to Parquet like any other numeric type.
Good to know: the split exists because a month has no fixed length. “One month
and three days” cannot be compared with “thirty-three days”, so a type that
mixes them cannot be ordered, and a column of it cannot be sorted or stored.
DATE'2026-01-31' + INTERVAL '1' MONTH is 2026-02-28: month arithmetic clamps
to the end of the month.
Scala 2.13 support (SPARK-34218)
Spark published Scala 2.13 builds alongside 2.12. 4.0 made 2.13 the only option.
Spark 3.3, June 2022
Row-level runtime filtering with bloom filters (SPARK-32268)
This extends the dynamic pruning idea from partitions to rows.
How it works. A bloom filter is a compact bit array that answers “is this
key possibly in the set?” with no false negatives and a small, tunable rate of
false positives. Spark builds one from the join keys of the filtered small side,
then adds a might_contain(bloom_filter, key) filter to the large side’s scan,
so rows that cannot possibly match are dropped before they are shuffled. It buys
a smaller shuffle at the cost of building the filter, which is why it is gated on
size thresholds. Lowering the thresholds shows it firing:
import tempfile
from pyspark.sql import SparkSession
spark = (SparkSession.builder.appName("bloom33")
.config("spark.sql.autoBroadcastJoinThreshold", "-1")
.config("spark.sql.adaptive.enabled", "false")
.config("spark.sql.optimizer.runtime.bloomFilter.applicationSideScanSizeThreshold", "1KB")
.config("spark.sql.optimizer.runtime.bloomFilter.creationSideThreshold", "10MB")
.getOrCreate())
p = tempfile.mkdtemp()
spark.range(0, 2_000_000).selectExpr("id AS k", "id % 100 AS v").write.parquet(p + "/big")
spark.range(0, 1000).selectExpr("id * 7 AS k", "CAST(id AS STRING) AS tag").write.parquet(p + "/small")
big, small = spark.read.parquet(p + "/big"), spark.read.parquet(p + "/small").where("tag LIKE '1%'")
df = big.join(small, "k")
df.collect()
print("bloom filter applied:", "might_contain" in df._jdf.queryExecution().executedPlan().toString())
# bloom filter applied: True
spark.stop()
Good to know: the plan marker to look for is might_contain. The default
thresholds only fire it when the large side’s scan is over 10 GB and the
filtered small side is under 10 MB, which is why an ordinary laptop test shows
nothing until the thresholds are lowered. It helps most when the small side’s
filter is selective and the join would otherwise shuffle a huge table; when the
small side could simply be broadcast, a broadcast join is better still. It is on
by default from 3.4.
Error classes (SPARK-38781)
The unglamorous change that improved daily life most.
How it works. Every error Spark raises is defined once, in a JSON file inside
the distribution, with a stable name such as DIVIDE_BY_ZERO or
UNRESOLVED_COLUMN.WITH_SUGGESTION, a message template, and (from 3.4) a
SQLSTATE. Exceptions carry the name and the parameters separately from the
formatted message, so code can match on them and messages can be reworded
without breaking anyone:
from pyspark.sql import SparkSession
from pyspark.errors import PySparkException
spark = SparkSession.builder.appName("err33").getOrCreate()
try:
spark.sql("SELECT no_such_col FROM range(1)").collect()
except PySparkException as e:
print(e.getCondition(), e.getSqlState())
# UNRESOLVED_COLUMN.WITH_SUGGESTION 42703
spark.stop()
Good to know: getCondition() is the 4.0 name. On 3.x the same method is
getErrorClass(), which still works on 4.x but emits a FutureWarning saying
it is deprecated. e.getMessageParameters() returns the values that filled the
template, such as the missing column’s name, which is what to log instead of
parsing the message text.
Complex types in the vectorized Parquet reader (SPARK-34863)
How it works. Until 3.3, a query that read any struct, array or map column from Parquet fell back to the row-based reader for the whole scan. The vectorized reader learned Parquet’s repetition and definition levels, which encode nesting, so nested columns are decoded in batches too.
Good to know: it is controlled by
spark.sql.parquet.enableNestedColumnVectorizedReader. Nested-heavy tables,
such as event data with arrays of structs, are where upgrades to 3.3 often
speed up without any code change.
The hidden _metadata column (SPARK-37273)
Every file-based source exposes a hidden struct describing the file each row came from. It only appears when you select it:
import tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("meta33").getOrCreate()
path = tempfile.mkdtemp()
spark.range(0, 4).repartition(2).write.mode("overwrite").parquet(path)
df = spark.read.parquet(path)
df.select("_metadata").printSchema()
# root
# |-- _metadata: struct (nullable = false)
# | |-- file_path: string (nullable = false)
# | |-- file_name: string (nullable = false)
# | |-- file_size: long (nullable = false)
# | |-- file_block_start: long (nullable = false)
# | |-- file_block_length: long (nullable = false)
# | |-- file_modification_time: timestamp (nullable = false)
# | |-- row_index: long (nullable = false)
(df.groupBy("_metadata.file_name").count()
.selectExpr("left(file_name, 10) AS file_prefix", "count AS rows")
.show())
# +-----------+----+
# |file_prefix|rows|
# +-----------+----+
# | part-00000| 2|
# | part-00001| 2|
# +-----------+----+
spark.stop()
How it works. The fields are filled by the scan itself from the file it is
reading, so they cost nothing unless selected, and filters on them, such as
_metadata.file_modification_time > ..., can be used to skip files.
Good to know: it is a richer replacement for input_file_name(), which
returns only the path. Size, modification time and the row’s position inside
the file come with it, so finding the one file with a corrupt value, and the row
inside it, is now a WHERE clause. A common ingestion pattern is to write
_metadata.file_path into the target table as a lineage column.
Profiler for Python and pandas UDFs (SPARK-37443)
How it works. Spark runs cProfile inside the Python workers on the executors,
collects the statistics per UDF, and brings them back to the driver. The
original 3.3 switch was spark.python.profile=true, with results printed by
sc.show_profiles(). 4.0 replaced it with a SQL-level setting that also works
over Spark Connect:
import io, contextlib
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
spark = SparkSession.builder.appName("prof40").getOrCreate()
spark.conf.set("spark.sql.pyspark.udf.profiler", "perf")
@udf("long")
def slow_square(x):
return sum(x for _ in range(x % 50))
spark.range(2000).select(slow_square("id")).collect()
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
spark.profile.show(type="perf")
out = buf.getvalue().splitlines()
print(out[0] if out else "no output")
print(next((l.strip() for l in out if "function calls" in l), "-"))
print(next((l.strip() for l in out if "slow_square" in l or "<genexpr>" in l), "-"))
# ============================================================
# 59000 function calls in ... seconds
# 51000 ... h07_profiler.py:9(<genexpr>)
spark.stop()
The times are elided above; they are laptop numbers. The call counts are what the profile is for: 51,000 of the 59,000 calls are the generator expression on line 9, which is exactly where the UDF spends its work.
Good to know: profiling adds overhead to every UDF call, so switch it on for a
diagnostic run, not in production. "memory" instead of "perf" gives a
line-by-line memory profile, which needs the memory-profiler package on the
executors.
Spark 3.4, April 2023
Spark Connect, Python client (SPARK-39375)
The architectural change of the 3.x line, even though it arrives near the end of it.
How it works. Until 3.4, a Spark application meant a driver JVM in your process: PySpark launched one and talked to it over a local socket, and your application’s lifetime was the driver’s lifetime. Connect replaces that with a gRPC protocol. The client builds DataFrames exactly as before, but each operation only extends an unresolved logical plan encoded as protocol buffers. When an action runs, the whole plan is sent to the Connect server, which analyses, optimizes and executes it, and streams the results back as Arrow batches. The client becomes thin, and the driver becomes a server that many clients can share. Start the server that ships in every distribution, then connect with nothing but a URL:
/opt/spark/sbin/start-connect-server.sh --master "local[4]"
pip install pyspark-client # or the full pyspark, plus grpcio and friends
from pyspark.sql import SparkSession
spark = SparkSession.builder.remote("sc://localhost:15002").getOrCreate()
print(type(spark).__module__)
# pyspark.sql.connect.session
df = spark.range(5).selectExpr("id", "id * id AS sq")
df.show()
# +---+---+
# | id| sq|
# +---+---+
# | 0| 0|
# | 1| 1|
# | 2| 4|
# | 3| 9|
# | 4| 16|
# +---+---+
spark.stop()
The session class is pyspark.sql.connect.session, not the classic one, and
there is no JVM in the client process.
Good to know: sparkContext and RDDs do not exist on this side of the wire,
which is the first thing to check when porting code to Connect. Analysis also
happens later: a misspelt column raises its error at the first action or
df.schema call, not when the DataFrame is defined. For several releases some
things also worked in a classic session and not in a Connect one; two posts here
work through that in practice, for
Iceberg and for
Hudi.
DEFAULT values for columns (SPARK-38334)
A column can declare a default, so an INSERT can omit it or say DEFAULT
explicitly:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("default34").getOrCreate()
spark.sql("CREATE TABLE users (id INT, country STRING DEFAULT 'IN', active BOOLEAN DEFAULT true) USING parquet")
spark.sql("INSERT INTO users (id) VALUES (1)")
spark.sql("INSERT INTO users VALUES (2, 'US', DEFAULT)")
spark.sql("SELECT * FROM users ORDER BY id").show()
# +---+-------+------+
# | id|country|active|
# +---+-------+------+
# | 1| IN| true|
# | 2| US| true|
# +---+-------+------+
spark.sql("DROP TABLE users")
spark.stop()
How it works. The default is stored in the table’s metadata as an expression
and filled in by the analyzer whenever an INSERT, UPDATE or MERGE leaves
the column out. Adding a column with a default to an existing table also makes
old rows read that value, without rewriting their files.
Good to know: defaults must be constant expressions, such as literals; they cannot refer to other columns. Support depends on the table format, and the built-in Parquet, ORC, JSON and CSV tables all have it.
TIMESTAMP_NTZ, timestamp without time zone (SPARK-35662)
Spark’s TIMESTAMP is an instant: it is stored in UTC and displayed in the
session time zone, so the same value prints differently in different zones.
TIMESTAMP_NTZ is a wall-clock reading that never shifts:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ntz34").getOrCreate()
for tz in ["UTC", "Asia/Kolkata"]:
spark.conf.set("spark.sql.session.timeZone", tz)
r = spark.sql("""SELECT CAST(TIMESTAMP_NTZ'2026-01-01 09:00:00' AS STRING) AS ntz,
CAST(CAST(TIMESTAMP_LTZ'2026-01-01 09:00:00 UTC' AS TIMESTAMP) AS STRING) AS ltz""").first()
print(f"{tz:13s} ntz={r.ntz} ltz={r.ltz}")
# UTC ntz=2026-01-01 09:00:00 ltz=2026-01-01 09:00:00
# Asia/Kolkata ntz=2026-01-01 09:00:00 ltz=2026-01-01 14:30:00
spark.stop()
How it works. Both are stored as a count of microseconds. For TIMESTAMP
(also called TIMESTAMP_LTZ) the count is from the epoch in UTC, and the
session time zone is applied on display. For TIMESTAMP_NTZ the count describes
the local date and time directly, and no time zone is ever applied.
Good to know: use TIMESTAMP_NTZ for values that are meant as local times,
such as a store’s opening hours or a date of birth with a time, and plain
TIMESTAMP for events that happened at an instant.
spark.sql.timestampType=TIMESTAMP_NTZ makes TIMESTAMP in DDL and literals
mean the no-time-zone type for a whole session.
Lateral column aliases (SPARK-27561)
The small feature people notice immediately, because every other SQL engine already had it:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("lca").getOrCreate()
spark.sql("SELECT 1 AS a, a + 1 AS b").show()
# +---+---+
# | a| b|
# +---+---+
# | 1| 2|
# +---+---+
spark.stop()
Before 3.4, referring to a in the same SELECT list that defined it was an
error, and you repeated the expression or nested a subquery. The old behaviour
is still one setting away, which shows what the feature changed. On 3.5.9:
spark.conf.set("spark.sql.lateralColumnAlias.enableImplicitResolution", "false")
spark.sql("SELECT 1 AS a, a + 1 AS b").show()
# [UNRESOLVED_COLUMN.WITHOUT_SUGGESTION] A column or function parameter with name `a` cannot be resolved. ; line 1 pos 15;
How it works. The analyzer resolves a name first against the input columns,
and only if that fails against aliases defined earlier in the same SELECT
list. It works in aggregate queries too, such as
SELECT sum(amount) AS total, total / count(*) AS avg_amount.
Good to know: because real columns win, an alias with the same name as an input column is not used by later expressions; the input column is. Pick alias names that do not collide.
Parameterized SQL (SPARK-41271)
The security feature of the release.
How it works. Values are bound as literals by the analyzer, never spliced into the SQL text, so an input containing a quote cannot change the query:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("params34").getOrCreate()
spark.createDataFrame([("pune", 1200), ("goa", 300), ("pune", 50)], "city STRING, amount INT") \
.createOrReplaceTempView("orders")
spark.sql("SELECT city, amount FROM orders WHERE city = :city AND amount > :min ORDER BY amount",
args={"city": "pune", "min": 100}).show()
# +----+------+
# |city|amount|
# +----+------+
# |pune| 1200|
# +----+------+
spark.sql("SELECT * FROM orders WHERE amount > ? ORDER BY amount", args=[200]).show()
# +----+------+
# |city|amount|
# +----+------+
# | goa| 300|
# |pune| 1200|
# +----+------+
spark.stop()
Named markers came in 3.4 and positional ? markers in 3.5. Parameters can only
stand where a value goes. For a table or column name, wrap the parameter in the
IDENTIFIER clause, also from 3.4:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ident34").getOrCreate()
spark.range(3).createOrReplaceTempView("events_2026")
spark.sql("SELECT count(*) AS n FROM IDENTIFIER(:tbl)", args={"tbl": "events_2026"}).show()
# +---+
# | n|
# +---+
# | 3|
# +---+
spark.stop()
Good to know: stop building SQL with f-strings. Besides injection, bound parameters also avoid quoting bugs with dates, decimals and strings containing quotes.
UNPIVOT and melt (SPARK-38864)
The reverse of a pivot: columns become rows:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("unpivot34").getOrCreate()
wide = spark.createDataFrame([("pune", 100, 150), ("goa", 40, 60)], "city STRING, q1 INT, q2 INT")
wide.unpivot("city", ["q1", "q2"], "quarter", "amount").orderBy("city", "quarter").show()
# +----+-------+------+
# |city|quarter|amount|
# +----+-------+------+
# | goa| q1| 40|
# | goa| q2| 60|
# |pune| q1| 100|
# |pune| q2| 150|
# +----+-------+------+
spark.stop()
The SQL form is
SELECT * FROM wide UNPIVOT (amount FOR quarter IN (q1, q2)).
How it works. Each input row is expanded into one output row per listed column, carrying the identifier columns, the column’s name and its value. It is a single projection with no shuffle.
Good to know: the unpivoted columns must have compatible types, since they
end up in one column. UNPIVOT drops rows whose value is NULL unless you
write UNPIVOT INCLUDE NULLS.
Bloom filter joins enabled by default (SPARK-38841)
The 3.3 runtime filter, now on:
spark.sql.optimizer.runtime.bloomFilter.enabled defaults to true.
SQLSTATE codes on error classes (SPARK-41994)
Every error class carries a five-character SQLSTATE, such as 42703 for an
unresolved column in the 3.3 example above. The first two characters are the
class (42 is a syntax or access rule violation, 22 a data exception such as
division by zero), which is the standard that JDBC and ODBC clients and SQL
tools understand.
Memory profiler for UDFs (SPARK-40281)
A line-by-line memory profile of Python UDFs, using the memory-profiler
package: which lines of the UDF allocated how much. In 4.x it is the "memory"
mode of the profiler shown in the 3.3 entry.
Spark 3.5, September 2023
The long-term support release of the 3.x line in practice, and where most organisations that have not moved to 4.x are sitting. Its maintenance line is still active: 3.5.9 shipped in July 2026, after 4.2.0.
Scala and Go clients for Spark Connect (SPARK-42554, SPARK-43351)
Connect stopped being Python-only. The Scala client is a separate artifact,
spark-connect-client-jvm, and the Go client a separate repository; both speak
the same protocol:
val spark = org.apache.spark.sql.SparkSession.builder().remote("sc://localhost:15002").getOrCreate()
spark.range(5).show()
How it works. Because the protocol, not a language binding, is the interface, a client in any language only needs to build plan protobufs and read Arrow results. The Scala client can run inside an application with its own dependencies and its own Scala version, without Spark’s JARs on its classpath.
Good to know: Scala UDFs over Connect must be serialisable and their classes available on the server, which is the main porting cost for Scala code.
Structured Streaming on Connect (SPARK-42938)
readStream and writeStream work from Connect clients in Python and Scala.
How it works. The streaming query runs on the server; the client holds a
handle to it, for status, lastProgress, stop() and awaitTermination().
foreachBatch and streaming listeners written in Python are run in a Python
process on the server side, since the client may not stay connected.
pandas API on the Connect client (SPARK-42497)
pyspark.pandas works over a remote session.
Arrow-optimized Python UDFs (SPARK-40307)
The middle step between the plain UDF and the pandas UDF. You keep writing a function over plain Python values, one row at a time, but the rows travel between the JVM and Python as Arrow batches instead of pickled rows. It is one keyword:
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
spark = SparkSession.builder.appName("arrowpy35").getOrCreate()
@udf(returnType="string", useArrow=True)
def mask(email: str) -> str:
user, domain = email.split("@")
return user[0] + "***@" + domain
spark.createDataFrame([("ranga@example.com",)], "email STRING").select(mask("email").alias("m")).show()
# +----------------+
# | m|
# +----------------+
# |r***@example.com|
# +----------------+
spark.stop()
How it works. The operator in the plan changes from BatchEvalPython to
ArrowEvalPython. Transfer is columnar and cheaper, but your function is still
called once per row with Python objects, so it is faster than a pickled UDF and
slower than a pandas UDF that works on whole Series.
Good to know: in 3.5 through 4.1 this is opt-in, per function or through
spark.sql.execution.pythonUDF.arrow.enabled, which reads false on both 3.5.9
and 4.1.3. In 4.2 that setting defaults to true, so every plain @udf takes
this path unless you turn it off. Arrow coerces values that do not match the
declared return type differently from pickling, so run your UDF tests once when
you switch.
Python user-defined table functions (SPARK-43798)
They fill a gap every Python team had hit: a UDF returns one value per row, but
some logic produces many rows from one input, such as splitting, exploding a
nested payload or calling an API that returns a list. A UDTF is a class whose
eval yields rows, usable from both the DataFrame API and SQL:
from pyspark.sql import SparkSession
from pyspark.sql.functions import udtf, lit
spark = SparkSession.builder.appName("udtf35").getOrCreate()
@udtf(returnType="word STRING, length INT")
class SplitWords:
def eval(self, text: str):
for w in text.split():
yield w, len(w)
SplitWords(lit("spark tables from python")).show()
# +------+------+
# | word|length|
# +------+------+
# | spark| 5|
# |tables| 6|
# | from| 4|
# |python| 6|
# +------+------+
spark.udtf.register("split_words", SplitWords)
spark.sql("SELECT * FROM split_words('hello udtf')").show()
# +-----+------+
# | word|length|
# +-----+------+
# |hello| 5|
# | udtf| 4|
# +-----+------+
spark.stop()
How it works. The class is instantiated once per task. eval is called once
per input row and can yield zero or more output rows; terminate, if defined,
is called once after the last row, so a UDTF can also aggregate. A UDTF can take
a whole query as input with TABLE(...), receiving each row as a Row:
from pyspark.sql import SparkSession
from pyspark.sql.functions import udtf
spark = SparkSession.builder.appName("udtf35b").getOrCreate()
@udtf(returnType="city STRING, orders INT, total INT")
class CityTotals:
def __init__(self):
self.acc = {}
def eval(self, row):
c = self.acc.setdefault(row["city"], [0, 0])
c[0] += 1; c[1] += row["amount"]
def terminate(self):
for city, (n, total) in sorted(self.acc.items()):
yield city, n, total
spark.udtf.register("city_totals", CityTotals)
spark.createDataFrame([("pune", 10), ("goa", 5), ("pune", 7)], "city STRING, amount INT").createOrReplaceTempView("orders")
spark.sql("SELECT * FROM city_totals(TABLE(orders)) ORDER BY city").show()
spark.sql("SELECT * FROM city_totals(TABLE(orders) PARTITION BY city) ORDER BY city").show()
# +----+------+-----+
# |city|orders|total|
# +----+------+-----+
# | goa| 1| 5|
# |pune| 1| 7|
# |pune| 1| 10|
# +----+------+-----+
# +----+------+-----+
# |city|orders|total|
# +----+------+-----+
# | goa| 1| 5|
# |pune| 2| 17|
# +----+------+-----+
spark.stop()
The first query is wrong in an instructive way: pune appears twice with one
order each. Without PARTITION BY, the input rows are spread across tasks
however Spark likes, each task gets its own instance of the class, and
terminate runs once per instance, so each task reported its own partial
totals. PARTITION BY city guarantees that all rows for a city reach the same
instance, and the totals come out right.
Good to know: treat a UDTF with terminate like a groupBy: if its output
depends on seeing a group’s rows together, the call must say how to group them
with PARTITION BY (and ORDER BY if the order matters). A UDTF can also
declare its output schema dynamically with an analyze static method instead
of a fixed returnType.
Distributed PyTorch training (SPARK-42471)
TorchDistributor in pyspark.ml.torch launches a PyTorch distributed
training function across Spark executors:
from pyspark.ml.torch.distributor import TorchDistributor
def train(lr):
import torch.distributed as dist
dist.init_process_group("gloo")
... # ordinary DistributedDataParallel code
return "done"
TorchDistributor(num_processes=4, local_mode=False, use_gpu=False).run(train, 1e-3)
How it works. It starts a barrier stage (2.4) with one task per process, sets
the environment variables PyTorch’s torchrun expects (the master address,
world size and rank of each worker), and runs your function in each. The data
loading is up to your function, typically reading a prepared dataset from
storage.
Good to know: the example is not runnable without PyTorch installed on every
executor. Use local_mode=True to run all processes on the driver while
developing.
PySpark errors on error classes (SPARK-42986)
Errors raised by PySpark itself, not just by the JVM, got error classes such as
PACKAGE_NOT_INSTALLED or CANNOT_CONVERT_TYPE, so Python-side failures can be
matched the same way as engine errors, through PySparkException.
The 4.x line: SQL becomes the surface again
The 4.x theme is that after a decade of enriching the DataFrame API, the project turned back to SQL and started adding things that SQL engines have and Spark did not: variant data, collations, user-defined functions written in SQL, procedural scripting, and a pipeline syntax. Alongside that, Spark Connect becomes the default way new components are built, and Python moves onto Arrow everywhere.
Spark 4.0, May 2025
The break. Read the removals first, because they are what stop an upgrade from compiling or starting.
Removed: Scala 2.12 (SPARK-45314)
Scala 2.13 is the only supported version.
Good to know: this is the removal that breaks builds hardest, because it is
not a Spark-only change: every JVM dependency you bring, including every
connector, needs a _2.13 artifact. Check your whole dependency tree for
_2.12 suffixes before you plan the upgrade, not after.
Removed: JDK 8 and JDK 11 (SPARK-45315)
JDK 17 is the minimum, and 4.0 also runs on JDK 21.
Good to know: the runtime on every node and in every container image has to move, and JVM flags that worked on 8 or 11, such as old GC options, can stop the JVM from starting on 17.
Removed: Mesos support (SPARK-44442)
Move to Standalone, YARN or Kubernetes.
Removed: Python 3.8 (SPARK-47993)
Python 3.9 is the minimum for PySpark.
Deprecated: SparkR (SPARK-49347)
Deprecated, not removed. 4.2 still ships it, but new R work belongs in
sparklyr or another language.
ANSI SQL mode on by default (SPARK-44444)
The change that will break jobs, and it is the right change. On 4.1.3:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ansi").getOrCreate()
print(spark.conf.get("spark.sql.ansi.enabled"))
# true
try:
spark.sql("SELECT 1/0").collect()
except Exception as e:
print(type(e).__name__, str(e).split("\n")[0])
# ArithmeticException [DIVIDE_BY_ZERO] Division by zero. Use `try_divide` to tolerate ...
spark.sql("SELECT try_divide(1, 0) AS v").show()
# +----+
# | v|
# +----+
# |NULL|
# +----+
spark.stop()
How it works. ANSI mode changes three things at once. Arithmetic checks
for overflow and division by zero and throws instead of wrapping or returning
NULL. Casts follow the SQL standard’s table of allowed conversions,
rejecting impossible ones at analysis time and throwing on invalid values at
runtime (the 3.1 section shows both). Implicit conversions fail loudly too:
'1' = 1 is still true, because the string is cast to a number, but on 4.1.3
'abc' = 1 raises CAST_INVALID_INPUT where 3.x quietly compared against
NULL and filtered the row out. The clearest way to see what changed is to run
one statement on both demo lines:
SELECT 1/0 AS a, CAST('abc' AS INT) AS b, 2147483647 + 1 AS c
On 3.5.9, where spark.sql.ansi.enabled is still false:
+----+----+-----------+
| a| b| c|
+----+----+-----------+
|NULL|NULL|-2147483648|
+----+----+-----------+
On 4.1.3 the same statement raises, on the division, during constant folding before any data is touched:
[DIVIDE_BY_ZERO] Division by zero. Use `try_divide` to tolerate divisor being 0 and return NULL instead.
Run the overflow on its own and 4.1.3 names that too:
[ARITHMETIC_OVERFLOW] integer overflow. Use 'try_add' to tolerate overflow and return NULL instead.
Look at column c. The division and the cast at least produced NULL, which a
downstream null check might catch. The integer overflow produced
-2147483648, a perfectly valid-looking negative number, from adding one to a
positive number. No null check catches that. ANSI mode is less about the
exceptions than about closing that hole.
For each operation that now raises, there is a try_ function that restores
the old null-returning behaviour where you genuinely want it:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ansi40").getOrCreate()
spark.sql("SELECT try_cast('abc' AS INT) AS a, try_add(2147483647, 1) AS b, try_element_at(array(1,2), 5) AS c").show()
# +----+----+----+
# | a| b| c|
# +----+----+----+
# |NULL|NULL|NULL|
# +----+----+----+
spark.stop()
Good to know: try_add returns NULL on overflow rather than wrapping, so
even the explicit opt-out is safer than the old default. ANSI mode only governs
arithmetic overflow, division by zero, invalid casts and out-of-range indexing.
A string that parses as a number but means something else still passes
silently. For the upgrade itself, run your test suite with ANSI on while still
on 3.5 (spark.sql.ansi.enabled=true); every error it raises is a place where
your pipeline was producing nulls or wrong numbers. Setting it back to false
on 4.x is possible, but treat that as a temporary switch.
The VARIANT data type (SPARK-45827)
The answer to a question every lakehouse team has asked: how do you store JSON whose schema you do not control, without either flattening it or storing a string and paying to parse it on every read?
How it works. The theory is the same one that makes Parquet fast. A JSON
string has to be parsed from the start every time you read one field from it. A
VARIANT value is parsed once on write into two binary parts: a metadata
dictionary holding each distinct object key once, and a value in which object
keys are replaced by indexes into that dictionary, objects carry a sorted table
of field offsets, and numbers are stored in binary rather than as text.
Extracting user.id becomes a binary search and a few offset lookups instead of
a scan through the text:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("variant").getOrCreate()
spark.sql("""
SELECT variant_get(parse_json('{"user":{"id":7,"city":"pune"}}'), '$.user.id', 'int') AS user_id
""").show()
# +-------+
# |user_id|
# +-------+
# | 7|
# +-------+
spark.stop()
The : path syntax is shorter than variant_get and reads like the JSON path it
is. schema_of_variant tells you what a value actually contains, which is the
first thing to run on data whose schema you do not control:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("variant40").getOrCreate()
spark.sql("""
SELECT j:user.city::string AS city,
j:tags[1]::string AS second_tag,
schema_of_variant(j) AS inferred
FROM (SELECT parse_json('{"user":{"id":7,"city":"pune"},"tags":["a","b"]}') AS j)
""").show(truncate=False)
# +----+----------+-------------------------------------------------------------------+
# |city|second_tag|inferred |
# +----+----------+-------------------------------------------------------------------+
# |pune|b |OBJECT<tags: ARRAY<STRING>, user: OBJECT<city: STRING, id: BIGINT>>|
# +----+----------+-------------------------------------------------------------------+
spark.stop()
The real use is a VARIANT column in a table holding records of different
shapes. Each row can have different fields, and queries pick out what they need:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("variant40b").getOrCreate()
spark.sql("CREATE TABLE events (id INT, payload VARIANT) USING parquet")
spark.sql("""INSERT INTO events VALUES
(1, parse_json('{"type":"click","page":"/home","ms":120}')),
(2, parse_json('{"type":"buy","sku":"A-7","price":49.5}')),
(3, parse_json('{"type":"click","page":"/cart"}'))""")
spark.sql("""
SELECT id, payload:type::string AS type,
try_variant_get(payload, '$.price', 'double') AS price,
payload:ms::int AS ms
FROM events WHERE payload:type::string = 'click' OR payload:price IS NOT NULL
ORDER BY id
""").show()
spark.sql("SELECT key, value FROM events, LATERAL variant_explode(payload) WHERE id = 2 ORDER BY key").show()
# +---+-----+-----+----+
# | id| type|price| ms|
# +---+-----+-----+----+
# | 1|click| NULL| 120|
# | 2| buy| 49.5|NULL|
# | 3|click| NULL|NULL|
# +---+-----+-----+----+
# +-----+-----+
# | key|value|
# +-----+-----+
# |price| 49.5|
# | sku|"A-7"|
# | type|"buy"|
# +-----+-----+
spark.sql("DROP TABLE events")
spark.stop()
A path that does not exist in a row returns NULL (ms for rows 2 and 3), and
variant_explode turns one variant object into key-value rows, which is how you
discover the fields a column actually holds.
Good to know: ::int and variant_get raise an error if a value exists but
cannot be cast to the requested type; try_variant_get returns NULL instead,
which is the safer choice on data you do not control. Values printed by
variant_explode are still variants, which is why strings keep their JSON
quotes. 4.1’s shredding goes further, storing frequently seen paths as ordinary
typed Parquet columns so that they get column statistics and pruning.
String collations (SPARK-46830)
A string column can declare its own comparison rules, so case-insensitive
matching stops requiring lower() on both sides:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("collation").getOrCreate()
spark.sql("SELECT 'Spark' COLLATE UTF8_LCASE = 'spark' AS eq").show()
# +----+
# | eq|
# +----+
# |true|
# +----+
spark.stop()
How it works. A collation is part of the string type, so every operation
on the column, equality, GROUP BY, DISTINCT, joins, sorting and string
functions such as contains, follows it. Declared on a column, it applies
without anyone having to remember it:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("collate40").getOrCreate()
spark.sql("CREATE TABLE users (name STRING COLLATE UTF8_LCASE, city STRING) USING parquet")
spark.sql("INSERT INTO users VALUES ('Ravi', 'pune'), ('RAVI', 'goa'), ('ravi', 'pune'), ('Sita', 'goa')")
spark.sql("SELECT name, count(*) AS n FROM users GROUP BY name ORDER BY n DESC").show()
spark.sql("SELECT count(*) AS n FROM users WHERE name = 'rAvI'").show()
spark.sql("SELECT 'Äpfel' COLLATE UNICODE_CI_AI = 'apfel' AS ci_ai, collation('a' COLLATE UNICODE_CI) AS coll").show(truncate=False)
# +----+---+
# |name| n|
# +----+---+
# |Ravi| 3|
# |Sita| 1|
# +----+---+
# +---+
# | n|
# +---+
# | 3|
# +---+
# +-----+-------------------------+
# |ci_ai|coll |
# +-----+-------------------------+
# |true |SYSTEM.BUILTIN.UNICODE_CI|
# +-----+-------------------------+
spark.sql("DROP TABLE users")
spark.stop()
The three spellings of Ravi group together, and the filter matches all three.
UTF8_LCASE is a fast, locale-free lowercase comparison; the UNICODE family
comes from the ICU library, with _CI for case-insensitive and _AI for
accent-insensitive, so Äpfel equals apfel.
Good to know: which of the grouped values is shown (Ravi here) is not
defined; any of the equal spellings can be returned. Comparing columns with two
different collations is an error, which forces you to choose one explicitly
with COLLATE.
SQL user-defined functions (SPARK-46057) and session variables (SPARK-42849)
Together they close a long-standing gap for teams whose logic lives in SQL rather than in Scala or Python:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("sqludf").getOrCreate()
spark.sql("DECLARE OR REPLACE cutoff INT DEFAULT 100")
spark.sql("SET VAR cutoff = 250")
spark.sql("""
CREATE OR REPLACE TEMPORARY FUNCTION over_cutoff(amount INT)
RETURNS BOOLEAN RETURN amount > cutoff
""")
spark.sql("SELECT over_cutoff(300) AS big, over_cutoff(100) AS small").show()
# +----+-----+
# | big|small|
# +----+-----+
# |true|false|
# +----+-----+
spark.stop()
How it works. A SQL UDF’s body is a SQL expression or query, stored in the catalog, and inlined into the calling query during analysis. Unlike a Python UDF there is no separate process and nothing opaque: the optimizer sees through it, pushes filters and folds constants as if you had written the expression yourself. A function can also return a table, which makes it a parameterized view:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("sqltvf40").getOrCreate()
spark.createDataFrame([(1, "pune", 1200), (2, "goa", 300), (3, "pune", 450)], "id INT, city STRING, amount INT").createOrReplaceTempView("orders")
spark.sql("""
CREATE OR REPLACE TEMPORARY FUNCTION orders_over(min_amount INT)
RETURNS TABLE (id INT, city STRING, amount INT)
RETURN SELECT id, city, amount FROM orders WHERE amount > min_amount
""")
spark.sql("SELECT * FROM orders_over(400) ORDER BY id").show()
# +---+----+------+
# | id|city|amount|
# +---+----+------+
# | 1|pune| 1200|
# | 3|pune| 450|
# +---+----+------+
spark.stop()
Session variables are typed values that live for the session, declared with
DECLARE, changed with SET VAR, and usable anywhere a literal can go.
Good to know: a permanent SQL function (without TEMPORARY) is stored in the
catalog and available to every user of it, including from BI tools, which is
the main reason to prefer it over a Python UDF for shared business rules.
SQL pipe syntax (SPARK-49555)
A query written as a sequence of steps in execution order, instead of a SELECT
whose clauses run in an order that does not match how they are written.
How it works. Each |> operator takes the table produced so far and applies
one transformation: WHERE, SELECT, EXTEND (add columns), SET (replace
columns), AGGREGATE ... GROUP BY, JOIN, ORDER BY, LIMIT and more. The
query reads top to bottom in the order it runs:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("pipe40b").getOrCreate()
spark.createDataFrame([(1, "pune", 1200), (2, "goa", 300), (3, "pune", 450), (4, "goa", 900)], "id INT, city STRING, amount INT").createOrReplaceTempView("orders")
spark.sql("""
FROM orders
|> WHERE amount > 350
|> EXTEND amount * 1.18 AS gross
|> AGGREGATE count(*) AS n, round(sum(gross), 1) AS gross_total GROUP BY city
|> ORDER BY gross_total DESC
|> LIMIT 1
""").show()
# +----+---+-----------+
# |city| n|gross_total|
# +----+---+-----------+
# |pune| 2| 1947.0|
# +----+---+-----------+
spark.stop()
This is where my first guess was wrong. I wrote this:
FROM VALUES (1), (2), (3) AS t(x) |> WHERE x > 1 |> SELECT sum(x) AS s
and got [PIPE_OPERATOR_CONTAINS_AGGREGATE_FUNCTION]. Aggregation has its own
pipe operator, and |> SELECT refuses to take an aggregate function. The
correct form:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("pipe").getOrCreate()
spark.sql("FROM VALUES (1), (2), (3) AS t(x) |> WHERE x > 1 |> AGGREGATE sum(x) AS s").show()
# +---+
# | s|
# +---+
# | 5|
# +---+
spark.stop()
Good to know: pipe syntax and ordinary SQL mix freely; any subquery can be written either way, and the optimizer produces the same plan for both, so there is no performance difference. It is a readability feature, most useful for long queries that are otherwise a stack of nested subqueries.
Built-in XML data source (SPARK-44265)
It retires the external spark-xml package:
import os, tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("xml40").getOrCreate()
d = tempfile.mkdtemp()
with open(os.path.join(d, "books.xml"), "w") as f:
f.write("""<catalog>
<book id="b1"><title>Learning Spark</title><price>40.5</price></book>
<book id="b2"><title>Spark Internals</title><price>55.0</price></book>
</catalog>""")
df = spark.read.option("rowTag", "book").xml(d)
df.printSchema()
# root
# |-- _id: string (nullable = true)
# |-- price: double (nullable = true)
# |-- title: string (nullable = true)
df.orderBy("_id").show()
# +---+-----+---------------+
# |_id|price| title|
# +---+-----+---------------+
# | b1| 40.5| Learning Spark|
# | b2| 55.0|Spark Internals|
# +---+-----+---------------+
spark.stop()
How it works. rowTag names the element that becomes a row. Child elements
become columns, repeated children become arrays, nested elements become
structs, and attributes become columns prefixed with _ (the
attributePrefix option). Types are inferred like JSON’s, so price came out
as a double. from_xml and schema_of_xml do the same for XML strings in a
column.
Good to know: files are split between tasks by searching for the row tag, so very large single-document files still parallelise. As with JSON, pass an explicit schema in production; inference reads all the data.
Python Data Source API
The more consequential addition for anyone who has had to write a connector.
How it works. Before 4.0, a new source meant a Scala or Java implementation of
the Data Source V2 interfaces, packaged as a JAR. Now it is a Python class. For
reading, schema declares the columns, partitions says how to split the work,
and read runs on the executors and yields rows for one partition:
from pyspark.sql import SparkSession
from pyspark.sql.datasource import DataSource, DataSourceReader, InputPartition
class CountdownDataSource(DataSource):
@classmethod
def name(cls):
return "countdown"
def schema(self):
return "n INT, label STRING"
def reader(self, schema):
return CountdownReader(int(self.options.get("start", 3)))
class CountdownReader(DataSourceReader):
def __init__(self, start):
self.start = start
def partitions(self):
return [InputPartition(i) for i in range(2)]
def read(self, partition):
for n in range(self.start, 0, -1):
if n % 2 == partition.value:
yield n, f"p{partition.value}"
spark = SparkSession.builder.appName("pyds40").getOrCreate()
spark.dataSource.register(CountdownDataSource)
spark.read.format("countdown").option("start", 5).load().orderBy("n").show()
# +---+-----+
# | n|label|
# +---+-----+
# | 1| p1|
# | 2| p0|
# | 3| p1|
# | 4| p0|
# | 5| p1|
# +---+-----+
spark.stop()
The label column shows the two partitions each producing their own rows in
parallel. Writing follows the same two-level shape Spark’s own writers use:
write runs once per partition on the executors and returns a commit message,
and commit is called once with all the messages, only after every partition
succeeded, which is where a sink makes the write visible:
import glob, json, os, tempfile
from pyspark.sql import SparkSession
from pyspark.sql.datasource import DataSource, DataSourceWriter, WriterCommitMessage
class JsonLinesSink(DataSource):
@classmethod
def name(cls):
return "jsonl_sink"
def writer(self, schema, overwrite):
return JsonLinesWriter(self.options["path"])
class Done(WriterCommitMessage):
def __init__(self, n): self.n = n
class JsonLinesWriter(DataSourceWriter):
def __init__(self, path): self.path = path
def write(self, rows): # runs on executors, once per partition
from pyspark import TaskContext
pid = TaskContext.get().partitionId()
n = 0
with open(os.path.join(self.path, f"part-{pid}.jsonl"), "w") as f:
for r in rows:
f.write(json.dumps(r.asDict()) + "\n"); n += 1
return Done(n)
def commit(self, messages): # called once, after every partition succeeded
with open(os.path.join(self.path, "_SUCCESS"), "w") as f:
f.write(str(sum(m.n for m in messages)))
spark = SparkSession.builder.appName("pyds40w").getOrCreate()
spark.dataSource.register(JsonLinesSink)
out = tempfile.mkdtemp()
spark.range(10).selectExpr("id", "id * id AS sq").repartition(3).write.format("jsonl_sink").option("path", out).mode("append").save()
print(sorted(os.path.basename(f) for f in glob.glob(out + "/*")))
print("rows recorded by commit():", open(os.path.join(out, "_SUCCESS")).read())
# ['_SUCCESS', 'part-0.jsonl', 'part-1.jsonl', 'part-2.jsonl']
# rows recorded by commit(): 10
spark.stop()
Good to know: options arrive as strings, hence the int(). The reader and
writer objects are pickled and sent to executors, so anything they hold must be
picklable, and clients such as database connections should be created inside
read or write, not in __init__. commit runs in a Python worker process,
not in your script, so it should record its result somewhere durable, as the
marker file here does, rather than print it. The same
interface has streamReader and streamWriter methods, so a REST API or an
internal service can become a streaming source without any JVM code.
Native plotting API
DataFrame.plot draws charts directly from a Spark DataFrame, with Plotly as the
backend:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("plot40").getOrCreate()
df = spark.createDataFrame([("pune", 1650), ("goa", 300)], "city STRING, revenue INT")
fig = df.plot.bar(x="city", y="revenue")
print(type(fig).__module__, type(fig).__name__)
# plotly.graph_objs._figure Figure
spark.stop()
How it works. Each chart type decides how much data it needs. Bar and line
charts bring back only a limited number of rows (spark.sql.pyspark.plotting.max_rows),
while histograms and box plots compute their bins and quartiles on the cluster
and bring back only the summary, so plotting a billion-row column does not pull
a billion rows to the driver.
Good to know: it needs plotly installed on the driver, and returns an
ordinary Plotly figure, so fig.show() in a notebook or fig.write_html(...)
work as usual.
Arbitrary State API v2 (transformWithState)
The streaming entry that replaces flatMapGroupsWithState.
How it works. The old API gave each key a single state object that you
serialised and replaced as a whole. The new one gives a stateful processor
any number of typed, named state variables, each stored separately in RocksDB:
ValueState for one value, ListState for an appendable list, and MapState
for a keyed map, each with an optional time-to-live after which entries expire
by themselves. Processors can also register timers, in processing time or
event time, which call the processor back for a key even when no new data
arrives. Here a processor keeps a running maximum per city across two separate
runs of the query:
import json, os, tempfile
import pandas as pd
from typing import Iterator
from pyspark.sql import SparkSession
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
class RunningMax(StatefulProcessor):
def init(self, handle: StatefulProcessorHandle) -> None:
self.best = handle.getValueState("best", "amount INT")
def handleInputRows(self, key, rows: Iterator[pd.DataFrame], timerValues) -> Iterator[pd.DataFrame]:
current = self.best.get()[0] if self.best.exists() else 0
for pdf in rows:
current = max(current, int(pdf["amount"].max()))
self.best.update((current,))
yield pd.DataFrame({"city": [key[0]], "best": [current]})
def close(self) -> None:
pass
spark = SparkSession.builder.appName("tws40").getOrCreate()
spark.conf.set("spark.sql.shuffle.partitions", "2")
spark.conf.set("spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
src, chk = tempfile.mkdtemp(), tempfile.mkdtemp()
emitted = []
def run(name, rows):
with open(os.path.join(src, name), "w") as f:
for c, a in rows:
f.write(json.dumps({"city": c, "amount": a}) + "\n")
q = (spark.readStream.schema("city STRING, amount INT").json(src)
.groupBy("city")
.transformWithStateInPandas(RunningMax(), "city STRING, best INT", "Update", "None")
.writeStream.outputMode("update")
.foreachBatch(lambda batch, batch_id: emitted.extend((batch_id, *r) for r in batch.collect()))
.option("checkpointLocation", chk).trigger(availableNow=True).start())
q.awaitTermination()
run("1.json", [("pune", 10), ("goa", 4)])
run("2.json", [("pune", 3), ("goa", 9)])
for row in sorted(emitted):
print(row)
# (0, 'goa', 4)
# (0, 'pune', 10)
# (1, 'goa', 9)
# (1, 'pune', 10)
spark.stop()
In the second batch pune received only a 3, and the processor still emitted 10, read back from state written by the first run.
Good to know: it requires the RocksDB state store provider, and on a minimal
image it also needs protobuf installed. Without it the Python worker crashes
with cannot import name 'descriptor' from 'google.protobuf', surfaced as
TransformWithStateInPySpark driver worker exited unexpectedly. The fourth
argument, "None", is the time mode; "ProcessingTime" or "EventTime" enables
timers.
State data source
Makes streaming state inspectable.
How it works. Point spark.read.format("statestore") at a checkpoint and the
state of a stateful operator becomes an ordinary DataFrame of key and value
structs, read at the latest committed batch or at any batch you name with
batchId. That turns “why did my streaming aggregate emit that?” from guesswork
into a query:
import json, os, tempfile
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("state40").getOrCreate()
spark.conf.set("spark.sql.shuffle.partitions", "2")
src, chk = tempfile.mkdtemp(), tempfile.mkdtemp()
with open(os.path.join(src, "a.json"), "w") as f:
for c in ["pune", "goa", "pune"]:
f.write(json.dumps({"city": c}) + "\n")
q = (spark.readStream.schema("city STRING").json(src).groupBy("city").count()
.writeStream.outputMode("update").format("noop")
.option("checkpointLocation", chk).trigger(availableNow=True).start())
q.awaitTermination()
spark.read.format("statestore").load(chk).selectExpr("key.city", "value.count").orderBy("city").show()
# +----+-----+
# |city|count|
# +----+-----+
# | goa| 1|
# |pune| 2|
# +----+-----+
spark.stop()
A companion state-metadata source lists what a checkpoint holds, which is
where to start when you do not know what a query keeps:
spark.read.format("state-metadata").load(chk) \
.select("operatorName", "stateStoreName", "numPartitions", "minBatchId", "maxBatchId").show()
# +--------------+--------------+-------------+----------+----------+
# | operatorName|stateStoreName|numPartitions|minBatchId|maxBatchId|
# +--------------+--------------+-------------+----------+----------+
# |stateStoreSave| default| 2| 0| 0|
# +--------------+--------------+-------------+----------+----------+
Good to know: reading state does not interfere with the running query, so you
can inspect a production checkpoint while it runs. With several stateful
operators, pick one with the operatorId option; state-metadata tells you the
IDs.
Spark Kubernetes Operator (SPARK-45923)
An official operator, released from its own repository, that manages Spark
applications and clusters as Kubernetes custom resources instead of
spark-submit commands.
How it works. You kubectl apply a SparkApplication resource describing the
job: image, main file, driver and executor sizes, Spark settings. The operator
watches for these resources, submits the application, tracks its state in the
resource’s status, and applies the restart policy on failure. Recurring jobs
and long-running Spark clusters are resources too, so a whole platform can be
managed with GitOps tools.
Good to know: it replaces the community operators people used before; if you run one of those, the resource format is different, so plan a migration rather than a drop-in swap.
Spark 4.1, December 2025
Spark Declarative Pipelines (SPARK-51727)
The largest addition.
How it works. You declare datasets, materialized views and streaming
tables, together with the queries that produce them. Spark builds a dataflow
graph from the references between them, then works out the execution order,
runs independent datasets in parallel, manages checkpoints for the streaming
ones, and retries failures. Materialized views are recomputed from their inputs
on each run; streaming tables are fed incrementally by flows, so each run
processes only new data. It ships with a spark-pipelines CLI. Here is a
complete pipeline project, mixing Python and SQL definitions. The spec file
names the project and tells the CLI where to find the definitions:
# spark-pipeline.yml
name: sales_pipeline
storage: file:///tmp/sales_pipeline/checkpoints
libraries:
- glob:
include: transformations/**
# transformations/sales.py
from pyspark import pipelines as dp
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
spark = SparkSession.active()
@dp.materialized_view
def raw_orders() -> DataFrame:
return spark.sql("""
SELECT * FROM VALUES (1, 'pune', 1200), (2, 'goa', 300), (3, 'pune', 450)
AS t(order_id, city, amount)""")
@dp.materialized_view
def city_revenue() -> DataFrame:
return spark.read.table("raw_orders").groupBy("city").agg(F.sum("amount").alias("revenue"))
-- transformations/big_orders.sql
CREATE MATERIALIZED VIEW big_orders AS
SELECT * FROM raw_orders WHERE amount > 400;
spark-pipelines run from the project directory produces this, trimmed of
timestamps:
Found 2 files matching glob 'transformations/**/*'
Registering SQL file /w/pipe/transformations/big_orders.sql...
Importing /w/pipe/transformations/sales.py...
Starting run...
Flow spark_catalog.default.raw_orders is QUEUED.
Flow spark_catalog.default.big_orders is QUEUED.
Flow spark_catalog.default.city_revenue is QUEUED.
Flow spark_catalog.default.raw_orders is RUNNING.
Flow spark_catalog.default.big_orders is RUNNING.
Flow spark_catalog.default.city_revenue is RUNNING.
Flow spark_catalog.default.big_orders has COMPLETED.
Flow spark_catalog.default.city_revenue has COMPLETED.
Run is COMPLETED.
Nobody declared an order. The CLI read spark.read.table("raw_orders") in one
file and FROM raw_orders in another, built the dependency graph from them,
ran raw_orders first, then ran the two downstream views in parallel. Reading
the results afterwards gives goa 300, pune 1650 for city_revenue and orders
1 and 3 for big_orders.
Good to know: two things cost me time.
- The CLI is a Spark Connect client. On a stock
apache/spark:4.1.3-python3image it fails twice before it runs, first withModuleNotFoundError: No module named 'yaml'and then with[PACKAGE_NOT_INSTALLED] grpcio >= 1.48.1 must be installed.pip install pyyaml grpciofixes both. Adopting Declarative Pipelines means adopting Spark Connect’s client dependencies. - Query functions may only describe a DataFrame. My first
raw_ordersusedspark.createDataFrame([...], "order_id INT, city STRING, amount INT"), and registration failed with[ATTEMPT_ANALYSIS_IN_PIPELINE_QUERY_FUNCTION] Operations that trigger DataFrame analysis or execution are not allowed in pipeline query functions.A query function is called while the graph is being registered, before anything runs. Creating a DataFrame from local data with a DDL schema string asks the server to analyse that schema, so it is refused.spark.sql,spark.read.tableand DataFrame transformations only build a plan, so they are allowed. The same rule forbidscount(),collect()andshow()inside a definition.
Besides materialized_view, the pyspark.pipelines module has
create_streaming_table, append_flow, temporary_view and create_sink, for
streaming tables fed by one or more flows and for writing to external systems.
spark-pipelines dry-run validates the graph without running it, which is the
check to put in CI.
Real-time mode for Structured Streaming (SPARK-53736)
A new trigger for sub-second end-to-end latency.
How it works. Micro-batch execution plans and schedules a small job per trigger, which puts a floor of around a hundred milliseconds or more under latency. Real-time mode keeps the micro-batch model for fault tolerance, with offsets and state committed per batch, but runs each batch as long-lived tasks that process records the moment they arrive and stream results downstream without waiting for the batch to end. The trigger’s interval is how long each such batch lasts, which sets how often progress is checkpointed, not the latency.
Good to know: in PySpark the trigger is trigger(realTime="5 seconds"), and
that keyword is in the 4.2.0 DataStreamWriter.trigger signature but not in
4.1.3’s, which only has processingTime, once, continuous and
availableNow. From Python, plan on 4.2 for real-time mode. It is the successor
to 2.3’s continuous processing, which never left experimental status.
SQL scripting enabled by default, and GA (SPARK-54499)
Spark SQL becomes a procedural language:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("scripting").getOrCreate()
print(spark.conf.get("spark.sql.scripting.enabled"))
# true
spark.sql("""
BEGIN
DECLARE total INT DEFAULT 0;
SET total = (SELECT sum(x) FROM VALUES (1), (2), (3) AS t(x));
SELECT total AS grand_total;
END
""").show()
# +-----------+
# |grand_total|
# +-----------+
# | 6|
# +-----------+
spark.stop()
How it works. A BEGIN ... END block is parsed as a script and executed by the
driver statement by statement. Each SQL statement inside runs as a normal Spark
query, and control flow (IF, CASE, WHILE, REPEAT, LOOP, FOR) and
local variables are handled between them, so the logic that used to live in a
Python wrapper around spark.sql calls now lives in the engine. Loops work the
way they do in other SQL dialects:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("loop41").getOrCreate()
spark.sql("""
BEGIN
DECLARE i INT DEFAULT 1;
DECLARE acc STRING DEFAULT '';
WHILE i <= 3 DO
SET acc = acc || CASE WHEN i % 2 = 0 THEN 'even ' ELSE 'odd ' END;
SET i = i + 1;
END WHILE;
SELECT trim(acc) AS parity;
END
""").show()
# +------------+
# | parity|
# +------------+
# |odd even odd|
# +------------+
spark.stop()
A FOR loop iterates over the rows of a query, with each column available on
the loop variable:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("for41").getOrCreate()
spark.createDataFrame([("pune", 1650), ("goa", 300)], "city STRING, revenue INT").createOrReplaceTempView("rev")
spark.sql("""
BEGIN
DECLARE report STRING DEFAULT '';
FOR r AS SELECT city, revenue FROM rev ORDER BY city DO
SET report = report || r.city || '=' || r.revenue || ' ';
END FOR;
IF length(report) > 0 THEN
SELECT trim(report) AS report;
ELSE
SELECT 'empty' AS report;
END IF;
END
""").show(truncate=False)
# +-----------------+
# |report |
# +-----------------+
# |goa=300 pune=1650|
# +-----------------+
spark.stop()
Good to know: a script returns the result of its last SELECT. A loop body
runs one Spark query per iteration, so a FOR over a million rows is a million
queries; use loops for orchestration over a handful of items (tables,
partitions, dates), and set-based SQL for the data itself. LEAVE and ITERATE
exit or continue a labelled loop early.
VARIANT enabled by default, GA, with shredding (SPARK-54454)
How it works. The 4.0 type, declared stable. Shredding writes the
commonly occurring fields of a variant column as separate, typed Parquet
columns next to the binary variant, with the rest left in the binary form. A
query reading v:user.id can then read just that column, with its min/max
statistics for data skipping, like any other column.
Good to know: shredding is a storage-layout decision taken on write; readers
see the same VARIANT values either way, so it can be introduced without
changing queries.
Recursive common table expressions
WITH RECURSIVE lets a query refer to itself, which matters for anyone who has
been faking a hierarchy traversal with a loop in Python:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("cte").getOrCreate()
spark.sql("""
WITH RECURSIVE nums(n) AS (
SELECT 1
UNION ALL
SELECT n + 1 FROM nums WHERE n < 5
)
SELECT collect_list(n) AS ns FROM nums
""").show()
# +---------------+
# | ns|
# +---------------+
# |[1, 2, 3, 4, 5]|
# +---------------+
spark.stop()
How it works. The CTE has an anchor query, run once, and a recursive query that refers to the CTE’s own name. Spark runs the anchor, then runs the recursive part repeatedly, each time over only the rows produced by the previous round, until a round produces nothing, and unions all the rounds. The typical use is walking a hierarchy, such as an org chart, carrying the path and depth along:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("cte41b").getOrCreate()
spark.createDataFrame([(1, "Asha", None), (2, "Bala", 1), (3, "Chen", 1), (4, "Dev", 2), (5, "Esha", 4)],
"id INT, name STRING, manager_id INT").createOrReplaceTempView("emp")
spark.sql("""
WITH RECURSIVE chain(id, name, level, path) AS (
SELECT id, name, 0, name FROM emp WHERE manager_id IS NULL
UNION ALL
SELECT e.id, e.name, c.level + 1, c.path || ' > ' || e.name
FROM emp e JOIN chain c ON e.manager_id = c.id
)
SELECT level, name, path FROM chain ORDER BY level, name
""").show(truncate=False)
# +-----+----+------------------------+
# |level|name|path |
# +-----+----+------------------------+
# |0 |Asha|Asha |
# |1 |Bala|Asha > Bala |
# |1 |Chen|Asha > Chen |
# |2 |Dev |Asha > Bala > Dev |
# |3 |Esha|Asha > Bala > Dev > Esha|
# +-----+----+------------------------+
spark.stop()
Good to know: each round is a separate Spark job, so a hierarchy a hundred levels deep costs a hundred jobs. Spark stops a runaway recursion at a configurable depth limit instead of looping forever, which also protects you from cycles in the data.
Stored procedures API for catalogs (SPARK-44167)
How it works. Catalogs can expose procedures, named routines with typed
arguments that run inside the catalog’s implementation and are invoked with
CALL. Iceberg’s maintenance procedures had this syntax through its own SQL
extensions; 4.1 makes it a Spark API that any catalog can implement:
CALL lake.system.rewrite_data_files(table => 'sales.orders');
Good to know: procedures are how table formats expose maintenance, such as compaction, snapshot expiry and orphan-file cleanup, so a native API means fewer format-specific SQL extensions to configure.
KLL and Theta sketches
Approximate aggregates that answer two common questions in bounded memory: how many distinct values (Theta), and what is the p99 (KLL quantiles):
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("sketch41").getOrCreate()
spark.range(0, 1_000_000).selectExpr("id % 50000 AS user_id", "id % 1000 AS latency_ms").createOrReplaceTempView("hits")
spark.sql("""
SELECT theta_sketch_estimate(theta_sketch_agg(user_id)) AS approx_users,
kll_sketch_get_quantile_bigint(kll_sketch_agg_bigint(latency_ms), 0.99) AS p99_latency
FROM hits
""").show()
# +------------+-----------+
# |approx_users|p99_latency|
# +------------+-----------+
# | 50213| 993|
# +------------+-----------+
spark.stop()
The true answers are 50,000 distinct users and a p99 of 990. The Theta estimate is within half a per cent, and the KLL quantile is within the sketch’s rank error.
How it works. An exact distinct count or percentile needs every value in one place. A sketch is a small, fixed-size summary, kilobytes, that can be built in parallel on each partition and merged. Theta sketches keep a sample of the smallest hash values seen, from which the number of distinct values can be estimated, and they also support set operations: union, intersection and difference. KLL sketches keep a compacted sample of values at several levels of detail, from which any quantile can be read with a bounded rank error.
Good to know: store the sketch itself, the binary output of
theta_sketch_agg, in a table, and combine stored sketches later with
theta_union_agg. That is what makes “distinct users this month” cheap when
you only keep daily aggregates; distinct counts cannot be summed, but sketches
can be merged. 4.2 adds tuple sketches, which carry a summary value with each
key.
JDBC driver for Spark Connect (SPARK-53484)
A JDBC driver that speaks the Connect protocol, with URLs of the form
jdbc:sc://host:15002, so JDBC-based tools can reach a Connect server without
the Thrift server.
Good to know: it gives BI and SQL tools the same thin-client path Python and Scala clients have, and the same isolation from the server’s JVM.
Arrow-native UDF and UDTF decorators (SPARK-52214, SPARK-52979)
The endpoint of the line that started with pandas UDFs in 2.3:
import pyarrow as pa
from pyspark.sql import SparkSession
from pyspark.sql.functions import arrow_udf
spark = SparkSession.builder.appName("arrowudf").getOrCreate()
@arrow_udf("long")
def plus_one(v: pa.Array) -> pa.Array:
return pa.compute.add(v, 1)
spark.range(3).select(plus_one("id").alias("p")).show()
# +---+
# | p|
# +---+
# | 1|
# | 2|
# | 3|
# +---+
spark.stop()
How it works. A pandas UDF receives Arrow data and converts it to pandas
Series, which costs a copy for some types and brings pandas’ own type quirks,
such as integers with nulls becoming floats. arrow_udf hands your function the
PyArrow arrays directly and takes arrays back, so there is no conversion at all
and nulls stay nulls.
Good to know: use pyarrow.compute functions inside, not Python loops over the
array, or you are back to row-at-a-time speed with extra steps.
Checksum-based shuffle stage retry (SPARK-51756)
How it works. When a shuffle map stage is recomputed, for example after an executor is lost, a non-deterministic stage (one using random numbers, round-robin repartitioning or unordered input) can produce different output from the first attempt. Downstream tasks that already read part of the old output would then combine two different versions. Spark now compares checksums of the recomputed output with the original, and when they differ it retries the whole downstream stage instead of mixing old and new output, avoiding silently incorrect results.
Good to know: it is the general version of the 2.3 fix for repartition after shuffle, and it relies on the per-partition checksums the 1.1 shuffle-files example shows on disk.
Spark ML on Connect, GA for Python (SPARK-51236)
pyspark.ml pipelines train and predict over a remote Connect session: model
fitting runs on the server, and the fitted model lives there, referenced from
the client.
Spark 4.2, July 2026
The current release.
Geospatial GEOMETRY and GEOGRAPHY types (SPARK-51658)
How it works. GEOMETRY holds shapes on a flat plane, where distances are
Euclidean in the coordinate system’s units; GEOGRAPHY holds shapes on the
earth’s surface, where distances follow the curve of the globe. Each carries a
spatial reference identifier (SRID) naming its coordinate system, and in
Spark 4.2 the SRID is part of the type itself. This is where I have to be
precise, because the release note summary and the shipped function registry do
not quite agree. The note says “ST_* functions, WKB/WKT and Parquet
read/write”. What is actually registered as a built-in in 4.2.0 is five
functions:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("geo").getOrCreate()
fns = [r[0] for r in spark.sql("SHOW FUNCTIONS").collect()]
print(sorted(f for f in fns if f.startswith("st_")))
# ['st_asbinary', 'st_geogfromwkb', 'st_geomfromwkb', 'st_setsrid', 'st_srid']
spark.stop()
There is no ST_Point, no ST_AsText, and no ST_Distance in the built-in
registry, and there is no cast from a WKT string either:
CAST('POINT(1 2)' AS GEOMETRY(4326)) fails with
DATATYPE_MISMATCH.CAST_WITHOUT_SUGGESTION. The entry point in 4.2.0 is
well-known binary (WKB), the standard binary encoding of shapes:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("geo2").getOrCreate()
# WKB for POINT(73.8567 18.5204), little-endian
wkb = "0101000000ed9e3c2cd4765240a1d634ef38853240"
spark.sql(f"""
SELECT typeof(st_geomfromwkb(unhex('{wkb}'))) AS geom_type,
typeof(st_geogfromwkb(unhex('{wkb}'))) AS geog_type,
st_srid(st_setsrid(st_geomfromwkb(unhex('{wkb}')), 4326)) AS srid
""").show(truncate=False)
# +-----------+---------------+----+
# |geom_type |geog_type |srid|
# +-----------+---------------+----+
# |geometry(0)|geography(4326)|4326|
# +-----------+---------------+----+
spark.stop()
Good to know: geometry(0) means SRID 0, an unspecified reference system.
GEOGRAPHY defaults to 4326, which is WGS 84, the latitude and longitude used
by GPS. The bare type name does not parse: CAST(NULL AS GEOMETRY) is a syntax
error while CAST(NULL AS GEOMETRY(4326)) is fine, because the SRID is part of
the type’s grammar. What 4.2 delivers is a standard type that table formats and
Parquet can store; for spatial functions today, libraries such as Apache Sedona
remain the tool.
Change data capture: the CHANGES clause (SPARK-55668)
How it works. A change feed returns the rows inserted, updated and deleted in a table between two versions or timestamps, each tagged with the kind of change, so downstream jobs can process only what changed. The syntax is not the one I guessed; the grammar takes a version or timestamp range:
SELECT * FROM sales CHANGES FROM VERSION 0 TO VERSION 1
On a built-in Parquet table that parses and analyses, then stops:
[UNSUPPORTED_FEATURE.CHANGE_DATA_CAPTURE] The feature is not supported:
Catalog spark_catalog does not support Change Data Capture (CDC). SQLSTATE: 0A000
Good to know: that is the correct behaviour and it tells you exactly what the feature is. Spark 4.2 ships the SQL surface, the parser and the analyzer rules for CDC; the capability itself is delegated to the Data Source V2 catalog. You get change feeds when your table format implements them, which is the same division of labour Spark used for time travel. Plain Parquet tables have no versions, so they can never support it.
Auto CDC in Declarative Pipelines (SPARK-56249)
A declarative way to apply a change feed to a target table as SCD Type 1.
How it works. A change feed is a stream of rows that each describe an insert,
update or delete of a key, possibly out of order. Applying it correctly means,
for each key, keeping only the latest change by a sequencing column and turning
it into an upsert or a delete. Auto CDC does that for you: you name the source,
the keys and the ordering column, and the pipeline works out the merges. In
4.2.0 it is create_auto_cdc_flow in pyspark.pipelines, which does not exist
in 4.1.3:
from pyspark import pipelines as dp
dp.create_streaming_table("customers")
dp.create_auto_cdc_flow(
target="customers",
source="customers_cdc_feed",
keys=["customer_id"],
sequence_by="updated_at",
)
The 4.2.0 signature also takes apply_as_deletes, a condition that marks a
change row as a delete, and column_list or except_column_list to choose which
columns are written. stored_as_scd_type accepts only 1 in this release, so
history-keeping SCD Type 2 is not yet available.
Good to know: SCD Type 1 keeps only the current row per key; if you need the history of every change, keep the raw change feed as its own table too.
Arrow-optimized Python UDFs and Arrow IPC on by default (SPARK-54555)
The change that affects everyone without being asked for. On 4.1.3,
spark.sql.execution.arrow.pyspark.enabled is false; on 4.2.0 it is true,
and spark.sql.execution.pythonUDF.arrow.enabled is true as well. Reading the
defaults directly on 4.2.0:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("defaults42").getOrCreate()
print(spark.version)
for k in ["spark.sql.execution.arrow.pyspark.enabled",
"spark.sql.execution.pythonUDF.arrow.enabled",
"spark.sql.ansi.enabled", "spark.sql.adaptive.enabled",
"spark.sql.scripting.enabled"]:
print(k, spark.conf.get(k))
# 4.2.0
# spark.sql.execution.arrow.pyspark.enabled true
# spark.sql.execution.pythonUDF.arrow.enabled true
# spark.sql.ansi.enabled true
# spark.sql.adaptive.enabled true
# spark.sql.scripting.enabled true
spark.stop()
How it works. The first setting moves toPandas() and
createDataFrame(pandas_df) onto Arrow; the second moves every plain @udf from
the pickled BatchEvalPython path to the columnar ArrowEvalPython path
described in the 3.5 entry. Both need PyArrow, and pandas, installed. On the
official apache/spark:4.2.0-python3 image, which has neither, a plain UDF still
works, with a warning:
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
spark = SparkSession.builder.appName("arrow42").getOrCreate()
print(spark.version, "pythonUDF.arrow =", spark.conf.get("spark.sql.execution.pythonUDF.arrow.enabled"))
try:
import pyarrow; print("pyarrow", pyarrow.__version__)
except ImportError:
print("pyarrow not installed")
@udf("string")
def shout(s):
return s.upper()
df = spark.createDataFrame([("pune",)], "city STRING").select(shout("city").alias("c"))
plan = df._jdf.queryExecution().executedPlan().toString()
print([n for n in ["ArrowEvalPython", "BatchEvalPython"] if n in plan])
print(df.collect())
# 4.2.0 pythonUDF.arrow = true
# pyarrow not installed
# RuntimeWarning: Arrow optimization failed to enable because PyArrow or Pandas is not installed.
# Falling back to a non-Arrow-optimized UDF.
# ['BatchEvalPython']
# [Row(c='PUNE')]
spark.stop()
Good to know: five switches, each flipped in a different release: adaptive execution in 3.2, ANSI in 4.0, scripting in 4.1, and the two Arrow settings in 4.2. The silent fallback above means a cluster without PyArrow keeps the old, slow path while the setting claims otherwise; install PyArrow on the executors to get the speed-up. And because the Arrow path coerces mismatched return types differently, run your UDF tests once after upgrading.
NEAREST BY top-k ranking join (SPARK-56395)
Nearest-neighbour search as a planner-visible operator, instead of a cross join with a sort, which is how people have been doing vector search on Spark.
How it works. The shape is
(APPROX | EXACT) NEAREST <k> BY (DISTANCE | SIMILARITY) <expr>, attached where
a join condition would go. For each row on the left, it keeps the k rows on
the right with the smallest distance, or largest similarity. EXACT computes
every pair; APPROX allows the engine to use faster approximate methods. For
each query row, keep the two documents with the closest score:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("nearest42").getOrCreate()
spark.createDataFrame([(1, 0.10), (2, 0.40), (3, 0.90)], "q_id INT, score DOUBLE").createOrReplaceTempView("queries")
spark.createDataFrame([("a", 0.05), ("b", 0.35), ("c", 0.50), ("d", 0.95)], "doc STRING, score DOUBLE").createOrReplaceTempView("docs")
spark.sql("""
SELECT q.q_id, d.doc
FROM queries q JOIN docs d
EXACT NEAREST 2 BY DISTANCE abs(q.score - d.score)
ORDER BY q.q_id, d.doc
""").show()
# +----+---+
# |q_id|doc|
# +----+---+
# | 1| a|
# | 1| b|
# | 2| b|
# | 2| c|
# | 3| c|
# | 3| d|
# +----+---+
spark.stop()
Query 1 at 0.10 gets a (distance 0.05) and b (0.25); query 3 at 0.90 gets
d (0.05) and c (0.40).
Good to know: it is top-k per left row, not a global top-k. In a real vector
search the distance expression would be a function over two embedding arrays,
and APPROX would let the engine trade exactness for speed.
Data Source V2 transaction management (SPARK-55855)
How it works. An API for connectors to group several reads and writes into one
transaction that commits or aborts together. It is what lets a table format make
multi-statement operations, and operations such as MERGE that read and write
the same table, atomic through Spark rather than through format-specific
extensions.
Path-based name resolution: SET PATH and CURRENT_PATH() (SPARK-54806)
A search path for unqualified names, like search_path in Postgres.
How it works. When a query names orders without a schema, Spark tries each
entry on the path in order until the name resolves. The default path is the
built-in and session namespaces followed by the current schema, and SET PATH
puts other schemas in front. It is off by default in 4.2.0; without
spark.sql.path.enabled=true, SET PATH fails with
[UNSUPPORTED_FEATURE.SET_PATH_WHEN_DISABLED]:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("path42").config("spark.sql.path.enabled", "true").getOrCreate()
print(spark.sql("SELECT current_path() AS p").first().p)
# system.builtin,system.session,spark_catalog.default
spark.sql("CREATE DATABASE IF NOT EXISTS sales")
spark.sql("CREATE TABLE IF NOT EXISTS sales.orders (id INT) USING parquet")
spark.sql("INSERT INTO sales.orders VALUES (42)")
spark.sql("SET PATH = spark_catalog.sales, DEFAULT_PATH")
print(spark.sql("SELECT current_path() AS p").first().p)
# spark_catalog.sales,system.builtin,system.session,spark_catalog.default
spark.sql("SELECT * FROM orders").show()
# +---+
# | id|
# +---+
# | 42|
# +---+
spark.sql("DROP TABLE sales.orders"); spark.sql("DROP DATABASE sales")
spark.stop()
Good to know: the current schema is still default, yet orders resolved to
sales.orders through the path. That is convenient and also a source of
surprise when two schemas on the path hold tables with the same name; the first
one wins.
Metric views (SPARK-54119)
How it works. A metric view defines measures (aggregations such as revenue or distinct customers) and dimensions (the columns they can be grouped by) once, in the catalog. Queries then ask for a measure by dimension, and the engine generates the aggregation, so every dashboard and report computes “revenue” the same way. It is a semantic layer inside Spark rather than in a separate BI tool.
Rebuilt web UI (SPARK-55760)
A new UI with a dark mode and side-by-side plan comparison, the kind of comparison the codegen and AQE examples in this post had to print by hand.
Java 25 (SPARK-51167)
Spark builds and runs on Java 25, the newest long-term-support JDK.
Which defaults changed underneath you?
Features are opt-in. Defaults are not. These are the settings to check before an upgrade, and every current value below was read from a live 4.1.3 or 4.2.0 session rather than from documentation.
spark.shuffle.manager: changed in 1.2 from hash to sort
The most consequential default flip in Spark’s history: one file plus an index per map task instead of one file per reducer. The setting is now internal and absent from the 4.1.3 configuration reference; sort shuffle is simply how Spark shuffles.
spark.shuffle.blockTransferService: changed in 1.2 from nio to netty
A faster network layer for fetching shuffle blocks. Netty is now the only implementation and the key is gone from the 4.1.3 configuration reference.
spark.io.compression.codec: changed in 1.1 from lzf to snappy
The codec for shuffle files, broadcast variables and spilled data. The current
default is lz4, changed again in a release whose notes do not record it.
spark.sql.autoBroadcastJoinThreshold: changed in 1.2 from 10000 to 10485760
Tables estimated below this size are broadcast to every executor instead of
shuffled. The current value is still 10485760b, 10 MB, eleven years later.
spark.sql.shuffle.partitions: unchanged, 200
The number of partitions every SQL shuffle writes. No release note records a
change, and it is still 200. AQE coalesces partitions after the shuffle, but
this value still decides how many are written in the first place.
spark.sql.orc.impl: changed in 2.4 from hive to native
ORC files are read with Spark’s own vectorized reader instead of Hive’s
libraries. The current value is native.
spark.sql.adaptive.enabled: changed in 3.2 from false to true
Adaptive query execution re-plans after each shuffle. Upgrading from 3.1 to 3.2 changes your query plans even if you change nothing else.
spark.sql.optimizer.dynamicPartitionPruning.enabled: true from 3.0
Dynamic partition pruning was introduced switched on, so fact-table scans are pruned from dimension filters without any configuration.
spark.sql.optimizer.runtime.bloomFilter.enabled: changed in 3.4 from false to true
Bloom filters built from a join’s small side are pushed into the large side’s scan when the size thresholds are met.
spark.sql.ansi.enabled: changed in 4.0 from false to true
Overflow, division by zero and invalid casts raise errors instead of returning
NULL or wrapping around. This is the default change most likely to break a
job, and the one most likely to be exposing a real data problem when it does.
spark.sql.scripting.enabled: changed in 4.1 from false to true
Multi-statement SQL scripts with variables and control flow run without any configuration.
spark.sql.execution.arrow.pyspark.enabled: changed in 4.2 from false to true
toPandas() and createDataFrame(pandas_df) transfer data through Arrow. It
reads false on 4.1.3 and true on 4.2.0.
spark.sql.execution.pythonUDF.arrow.enabled: changed in 4.2 from false to true
Plain Python UDFs exchange rows with the JVM as Arrow batches instead of pickled
rows. It reads true on 4.2.0.
Three of these change results rather than performance. ANSI mode turns silent nulls into errors. The proleptic Gregorian calendar in 3.0 shifts pre-1582 dates. And the 2.3 fix for shuffle plus repartition (SPARK-23207) changed wrong answers into right ones, which is also a behaviour change if something downstream was tuned to the wrong ones.
What was removed, and when?
Upgrades fail on removals more often than they fail on features.
Java 7, removed in 2.2
Java 8 became the minimum, so builds and clusters still on Java 7 had to move first.
Hadoop 2.5 and earlier, removed in 2.2
Hadoop 2.6 or later became the minimum for the Hadoop client libraries Spark builds against.
Scala 2.10, removed in 2.3
Scala 2.11 became the minimum, and every Scala application and library had to be rebuilt for it.
Python 3.7, removed in 3.5
PySpark required Python 3.8 or later.
Python 3.8, removed in 4.0
PySpark requires Python 3.9 or later.
Scala 2.12, removed in 4.0
Scala 2.13 is the only option. Every JVM dependency, including every connector,
needs a _2.13 build.
JDK 8 and JDK 11, removed in 4.0
JDK 17 is the minimum, so the runtime on every node and in every container image has to move.
Mesos, removed in 4.0
The cluster manager Spark was born on is gone. Move to Standalone, YARN or Kubernetes.
R 3.x, removed in 4.2
SparkR requires R 4.x.
SparkR was deprecated in 4.0 (SPARK-49347) and has not been removed: 4.2 still lists SparkR changes, including Java 25 support. Deprecated is not gone, but it is a clear signal about where to put new R work.
The Scala 2.12 removal in 4.0 is the one that breaks builds hardest, because it is not a Spark-only change. Every JVM dependency you bring, including every connector, needs a 2.13 artifact.
Common misconceptions
“Adaptive query execution is a Spark 3 feature.” The idea shipped in 1.6 (SPARK-9858) and only chose reducer counts. The implementation you know was written for 3.0 (SPARK-31412) and switched on by default in 3.2 (SPARK-33679). The distinction matters when reading old tuning advice, which may be about the 1.6 version.
“Structured Streaming arrived in 2.0.” It arrived as an experimental API in 2.0 and was declared generally available in 2.2 (SPARK-20844), two releases and a year later. Production advice written between those points is about a moving target.
“Kubernetes has been supported since 2.3.” Experimentally. The release note for 2.3 says configurations, container images and entrypoints were expected to change. It went GA in 3.1 (SPARK-33005).
“Spark 4 enables ANSI mode, so my job will fail loudly on bad data.” Only for the operations ANSI mode governs, which are arithmetic overflow, division by zero and invalid casts. A string that parses as a number but means something else still passes silently, as it should.
“spark.sql.shuffle.partitions does not matter any more, because AQE
coalesces.” AQE coalesces partitions down after the shuffle has been written.
The value still decides how many partitions get written in the first place, and
it is still 200. The plan output earlier in this post shows exactly that:
Exchange hashpartitioning(bucket#1L, 200) followed by AQEShuffleRead
coalesced.
“Spark 4.2 ships geospatial functions, so I can port my PostGIS queries.”
Five st_* functions are registered as built-ins in 4.2.0, and WKT is not among
the entry points. Check SHOW FUNCTIONS against the queries you intend to port
before planning the work.
Which version should you be on?
A short version of the advice, for the three positions most teams are actually in.
On 2.4. The correctness fix in 2.3 (SPARK-23207) is behind you, which is the good news. Everything else argues for moving: no adaptive execution, no dynamic partition pruning, no Kubernetes GA, and a Scala version no current connector ships for. The jump to 3.5 is large but well trodden.
On 3.1 to 3.3. The single biggest gain available to you is AQE by default in
3.2, and you can have most of it today by setting
spark.sql.adaptive.enabled=true and measuring. Do that first; it tells you
what the upgrade is worth before you spend it.
On 3.5. You are on the practical long-term-support release, and it is still being maintained: 3.5.9 shipped on 16 July 2026, two days after 4.2.0. Moving to 4.x costs you a Scala 2.13 rebuild, a JDK 17 or later runtime, and an ANSI-mode audit. Budget the ANSI audit seriously and treat it as a data-quality exercise rather than a compatibility chore, because every error it raises is a place your pipeline was silently producing nulls.
On 4.0 or 4.1. Moving within the 4.x line is ordinary. Watch for the Arrow defaults in 4.2 changing Python UDF behaviour, which is usually an improvement and is still a change worth a test run.
Frequently asked questions
Which Spark version should I move to from 2.4? 3.5 if you need a staged migration, because it keeps Scala 2.12 as an option and gets you adaptive execution, dynamic partition pruning and Kubernetes GA. 4.2 if you are willing to do the Scala 2.13 and JDK 17 work once rather than twice.
Is Spark 3.5 still maintained? Yes. 3.5.9 was released on 16 July 2026, two days after 4.2.0. The 3.5 line is the practical long-term-support release for organisations that have not yet moved to 4.x.
What actually breaks when I move to Spark 4? Four things, in descending order of how often they bite: Scala 2.12 artifacts that have no 2.13 build, a JDK older than 17, ANSI mode turning silent nulls into raised errors, and Mesos deployments having nowhere to go. The first two are build failures you find immediately. The third is a runtime failure you find in production unless you audit for it.
Do I have to use Spark Connect on Spark 4?
No. The classic in-process driver is still the default, and spark-submit works
as it always has. Connect is opt-in, with one caveat: some newer components are
Connect clients themselves, including the spark-pipelines CLI for Declarative
Pipelines, which needs grpcio installed even when your queries do not.
When did adaptive query execution actually become the default? 3.2, via SPARK-33679. It existed but was off by default in 3.0 and 3.1, and a much smaller version of the idea shipped in 1.6 as SPARK-9858. Tuning advice that predates 3.2 usually assumes it is off.
Has SparkR been removed? Not yet. It was deprecated in 4.0 (SPARK-49347) and is still present in 4.2, which added Java 25 support to it. Support for R 3.x was dropped in 4.2 (SPARK-57767), so an old R runtime will stop you before the deprecation does.
Why does the SQL pipe operator reject my aggregate function?
Because |> SELECT does not accept aggregates. Use |> AGGREGATE instead. The
error class is PIPE_OPERATOR_CONTAINS_AGGREGATE_FUNCTION and it names the fix.
Why does my plan show dynamicpruningexpression(true)?
Dynamic partition pruning was planned and then dropped. The usual cause is that
the table you expected to be pruned ended up as the broadcast (build) side of
the join, often because the other side has no statistics and was assumed to be
huge. Give the dimension real statistics, by writing it to a table or running
ANALYZE TABLE, and check which side the BroadcastExchange is on.
How do I unit-test a streaming query?
Use trigger(availableNow=True), available from 3.3, with a file source in a
temporary directory. The query processes what is there and stops, so
awaitTermination() returns and you can assert on the output. Use a Parquet or
foreachBatch sink rather than memory if the test restarts the query, because
the memory sink cannot recover from a checkpoint.
Why does spark-pipelines refuse my createDataFrame call?
Pipeline query functions run while the dataflow graph is registered, and may
only describe a DataFrame, not analyse or execute one. createDataFrame with a
schema string triggers analysis on the server, so it fails with
ATTEMPT_ANALYSIS_IN_PIPELINE_QUERY_FUNCTION. Build literal data with
spark.sql("SELECT * FROM VALUES ...") instead, and keep count(), collect()
and show() out of definitions.
The mental model
If you remember one thing, make it the shape rather than the list.
Spark spent its first line learning that the RDD hid too much from the engine, its second line replacing it with a schema the engine could compile, its third line admitting that plans made from estimates are wrong and building a mechanism to revise them, and its fourth line detaching the client from the driver and pouring the accumulated capability back into SQL.
Every feature in this post is an instance of one of those four moves. When the next release lands, that is the question to ask of each headline: which move is this, and does it change what my plans do or only what I am able to write?
References
- Spark release archive, the index every release note in this post was read from
- Spark 3.0.0 release notes for adaptive query execution, dynamic partition pruning and the TPC-DS claim quoted above
- Spark 4.0.0 release notes for ANSI by default,
VARIANT, collations and the Scala, JDK and Mesos removals - Spark 4.1.0 release notes for Declarative Pipelines, real-time mode and SQL scripting
- Spark 4.2.0 release notes for geospatial types, change data capture and the Arrow defaults
- Spark SQL migration guide for the behaviour changes each version introduced
- Structured Streaming programming guide for output modes, watermarks and stream-stream join state bounds
- Python Data Source API for the reader, writer and streaming interfaces used in the 4.0 example
- Spark Declarative Pipelines programming guide for the spec file, decorators and the restrictions on query functions
- Structured Streaming internals for what the checkpoint and state store look like on disk
- Spark joins in depth for the join strategies the 3.0 hints select
- Apache Spark architecture for the runtime the shuffle and scheduling changes act on
- Adaptive query execution for what the 3.x re-planning does and does not revisit
Trademarks
Apache Spark, Apache Hudi, Apache Iceberg, Apache Parquet, Apache Avro, Apache ORC, Apache Arrow, Apache Kafka, Apache Hadoop, Apache Mesos, Apache YuniKorn and Apache are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.
Found this useful?
These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.