Skip to content
Great engineers know 4 min read · spark.serializer

Serialization: Java, Kryo, Tungsten and Arrow

Where Spark serializes data and code, and which serializer matters for RDDs, DataFrames and Python.

TL;DR Spark serializes whenever data or code leaves a JVM: task closures, shuffles, broadcasts, cached RDDs, and rows going to Python. DataFrames already store rows in Tungsten's compact binary format, so 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.
Read first: RDDs, DataFrames and Datasets · The memory model (skip if you know them)

Where serialization happens

WhatSerialized 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 broadcastsTungsten binary rows (UnsafeRow), copied as bytes
Cached DataFramesCompressed in-memory columnar batches
RDD shuffles, broadcasts, serialized RDD cachingspark.serializer: Java by default, or Kryo
Rows to and from Python workersArrow 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)

Works with any class implementing Serializable. Slow, and output is large, because it writes full class metadata.

Kryo

Often up to 10 times faster and much more compact. Register your classes; unregistered ones still work but store their full class name with every object. Spark already uses Kryo internally when shuffling RDDs of simple types and strings.
PySpark · Turning on Kryo (JVM RDD code)
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 →:

PySpark · Open connections on executors, not the driver
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.

WhereSetting
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, applyInPandasAlways 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 sizespark.sql.execution.arrow.maxRecordsPerBatch (10,000); lower it if Python workers run out of memory
Why this matters: serialization is CPU time and garbage. Every hop between object representations (JVM objects, bytes, Python objects) costs on every row. The fastest Spark code avoids hops altogether: built-in functions on DataFrames never leave Tungsten's binary format.

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 questions

Why 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 read
Dynamic partition pruning
Skip whole partitions of a fact table at runtime using a filter on the dimension it joins to.

Related lessons

Previous: How aggregations execute

Primary sources: Tuning: data serialization · Apache Arrow in PySpark