Serialization: Java, Kryo, Tungsten and Arrow
Where Spark serializes data and code, and which serializer matters for RDDs, DataFrames and Python.
On this page
spark.serializer mostly affects RDD code, where Kryo is much faster than the default Java serialization. For Python, Apache Arrow replaces row-by-row pickling with columnar batches.Where serialization happens
| What | Serialized with |
|---|---|
| Tasktask: The work for one partition in one stage, run on one CPU core. Learn more → closures (your functions and the variables they capture) | Java serialization on the JVM; cloudpickle for Python functions |
| DataFrame shufflesshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more → and broadcasts | Tungsten binary rows (UnsafeRow), copied as bytes |
| Cached DataFrames | Compressed in-memory columnar batches |
| RDD shuffles, broadcasts, serialized RDD caching | spark.serializer: Java by default, or Kryo |
| Rows to and from Python workers | Arrow for pandas UDFsUDF: User-defined function: your own Python function applied to each row. Flexible but much slower than built-in functions. Learn more →, and for regular Python UDFs when Arrow-optimized (the default from Spark 4.2); otherwise pickle in batches |
Why DataFrames mostly do not care
A DataFrame row in Tungsten format is a compact byte layout: a null bitmap, fixed-width slots for numbers, and offsets to variable-length values like strings. Operators read fields straight from those bytes, and shuffling a row means copying its bytes. There are no JVM objects per row to create, serialize or garbage-collect. This is a big part of why DataFrames beat RDDs of objects, and why changing spark.serializer rarely changes a pure DataFrame job.
Java vs Kryo for RDDs
Java serialization (default)
Serializable. Slow, and output is large, because it writes full class metadata.Kryo
spark = (SparkSession.builder .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .config("spark.kryo.registrationRequired", "true") # fail fast on unregistered classes .config("spark.kryo.classesToRegister", "com.shop.Order,com.shop.Customer") .getOrCreate())
A "Kryo serialization failed: Buffer overflow" error means one object is larger than spark.kryoserializer.buffer.max (64 MB); raise it, or avoid huge single records such as giant collected arrays.
Task not serializable
The classic error: a function shipped to executorsexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more → captures an object that cannot be serialized, such as a database connection, a client, or (in Scala) the enclosing class. Create such objects on the executor, once per partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more →:
def save_partition(rows): conn = make_connection() # created on the executor for row in rows: conn.write(row.asDict()) conn.close() orders.foreachPartition(save_partition)
Arrow for Python
A row-at-a-time Python UDF pickles rows in the JVM, unpickles them in Python, calls your function once per row, and pickles the results back. Apache Arrow moves whole column batches in a format both sides read without conversion.
| Where | Setting |
|---|---|
toPandas() and createDataFrame(pandas_df) | spark.sql.execution.arrow.pyspark.enabled: off by default up to Spark 4.1, on by default from Spark 4.2 |
pandas UDFs (@pandas_udf), mapInPandas, applyInPandas | Always Arrow |
| Regular Python UDFs (Spark 3.5+) | @udf(useArrow=True) per UDF, or spark.sql.execution.pythonUDF.arrow.enabled for the session: off by default in 3.5, 4.0 and 4.1, on by default from Spark 4.2. useArrow overrides the session setting. |
| Batch size | spark.sql.execution.arrow.maxRecordsPerBatch (10,000); lower it if Python workers run out of memory |
Common mistakes
- Switching to Kryo to speed up a DataFrame job — DataFrames use Tungsten binary rows; the serializer setting barely matters there.
- Capturing a connection or client in a closure — Task not serializable, or one connection per row. Create it per partition on the executor.
- toPandas without Arrow on large results — Slow row-by-row conversion on top of collecting everything to the driverdriver: The process that runs your program, plans the work and sends tasks to executors. Learn more →.
What you learned
- Where Spark serializes data and code
- Java vs Kryo serialization, and why it matters mostly for RDDs
- How DataFrames avoid object serialization with Tungsten binary rows
- How Arrow speeds up data movement between the JVM and Python
Key takeaways
- Spark serializes closures, shuffles, broadcasts, cached RDDs and rows sent to Python.
- DataFrames use Tungsten binary rows, so spark.serializer mostly matters for RDDs.
- Kryo is faster and smaller than Java serialization; register classes.
- Create non-serializable resources per partition on executors.
- Arrow moves column batches between the JVM and Python.
Check yourself
3 questionsWhy does spark.serializer rarely affect a pure DataFrame job?
Show the answer
Rows are already Tungsten binary and shuffle as bytes. There are no per-row JVM objects to serialize.
A function in foreach uses a database connection created on the driver. What happens?
Show the answer
Task not serializable, or the connection is unusable on executors. Connections cannot be serialized; create them per partition on the executor.
What does Arrow replace for pandas UDFs?
Show the answer
Row-by-row pickling with columnar batches. Arrow batches are read directly by both the JVM and pandas.
Keep going
Up next · lesson 24 of 30 · 3 min readDynamic partition pruning
Skip whole partitions of a fact table at runtime using a filter on the dimension it joins to.
Related lessons
RDDs, DataFrames and DatasetsPySpark functions · 4 min read
UDFs and pandas UDFsSpark internals · 4 min read
The memory model
Previous: How aggregations execute
Primary sources: Tuning: data serialization · Apache Arrow in PySpark