Every time Spark moves a Java or Scala object between machines, writes it to disk or packs it into memory as bytes, a serializer turns the object graph into a byte stream. The default is Java serialization, which works with any Serializable class but is slow and verbose. The Spark tuning guide recommends Kryo instead, describing it as often as much as 10x faster and more compact, and notes that the only reason it is not the default is the need to register classes.

Switching is one line of configuration, but the benefit ranges from large to zero depending on which API your job uses, and a careless switch introduces new failure modes. This article explains where spark.serializer applies, how Kryo encodes objects, how to register classes and write custom serializers, how to measure the effect, and how to size buffers and diagnose failures. It assumes basic familiarity with Spark shuffle and Spark memory management.

Advertisement

Where spark.serializer applies, and where it does not

Spark uses more than one serializer. The one you configure with spark.serializer handles data: records of RDDs shuffled between stages, RDD partitions cached with serialized storage levels such as MEMORY_ONLY_SER and MEMORY_AND_DISK_SER, broadcast variable values, and values returned from tasks to the driver, for example by collect(). Task closures, the functions you pass to map and friends, are serialized separately with Java serialization, which is why a non-serializable object captured in a closure fails with Task not serializable no matter which data serializer you choose.

Which serializer handles which bytes in a Spark applicationDriverbuilds tasksTask closuresclosure serializer: JavaExecutorsrun tasksRDD shufflespark.serializerSerialized cacheMEMORY_ONLY_SER etc.Broadcast valuesspark.serializerTask resultse.g. collect()DataFrame / Dataset rowsTungsten UnsafeRow, encodersEncoders.kryo[T]whole object as one binary columnGreen boxes switch to Kryo when spark.serializer = KryoSerializerPurple: Kryo not involved. Amber: Kryo used, but the column becomes opaque to the optimizer
Kryo replaces Java serialization for RDD data paths only. DataFrame and Dataset rows use Tungsten's binary row format and never pass through spark.serializer; Encoders.kryo is the one place the SQL side uses Kryo, at a cost.

The DataFrame and Dataset APIs mostly bypass it. Spark SQL stores rows in Tungsten's binary UnsafeRow format and generates encoder code for case classes and primitives, so shuffles and caches of DataFrames do not call Kryo at all (see Spark Tungsten). A pure DataFrame pipeline therefore gains little or nothing from switching. The jobs that gain are RDD-heavy: graph algorithms, custom partitioning, legacy ML code, and anything that caches or shuffles RDDs of your own classes. Since Spark 2.0.0, Spark already uses Kryo internally when shuffling RDDs of simple types, arrays of simple types or strings, so the gain is concentrated on RDDs of custom classes.

PySpark is a special case. RDD records produced in Python are pickled in the Python worker and reach the JVM as opaque byte arrays, so the JVM serializer has little to work on.

How Kryo encodes an object

Java serialization writes a class descriptor into the stream, including the class name, a serial version id and field names and types, and then the field values. Within one stream later objects of the same class refer back to the descriptor, but the format is still verbose, and deserialization relies on reflection.

Kryo keeps a registry that maps classes to small integer ids and to a serializer for each class. For a registered class it writes the id as a variable-length integer, usually one or two bytes, followed by the fields as written by that class's serializer. For an unregistered class it writes the fully qualified class name instead, which Spark allows by default (spark.kryo.registrationRequired=false). Kryo remembers a name it has already written only until it resets, and by default it resets after every top-level object; Spark writes each shuffled or cached record as a top-level object, so the name is repeated per record. A name like com.example.events.Click costs about 25 bytes, which for small records can exceed the data. Kryo can also write integers in a variable-length form, so small positive values take fewer bytes, though whether it does depends on the IO implementation in use (see spark.kryo.unsafe below).

Spark pre-registers many common Scala classes through the AllScalaRegistrar from Twitter's chill library, plus a set of Spark's own internal classes. Your classes, and arrays of your classes, are not registered unless you register them. Note that Array[Click] is a distinct class from Click and needs its own registration.

Reference tracking (spark.kryo.referenceTracking, default true) makes Kryo remember every object written in a stream so that a second reference to the same object is written as a back-reference. It is required for object graphs with cycles and saves space when the same object appears many times. It costs a hash lookup per object; if your records are plain trees with no shared objects, disabling it is a measurable speed-up, and if they are not, disabling it causes a stack overflow on cycles or silent duplication of shared objects.

Advertisement

Turning it on and registering classes

The minimum is to set the serializer and register your classes with registerKryoClasses. Also set spark.kryo.registrationRequired=true in development and CI: Kryo then throws on the first unregistered class, which turns silent class-name overhead into an actionable error that names the class.

import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession

case class Click(userId: Long, itemId: Int, ts: Long, page: String)
case class Session(userId: Long, clicks: Array[Click])

val conf = new SparkConf()
  .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
  .set("spark.kryo.registrationRequired", "true")   // fail fast on any class you forgot
  .registerKryoClasses(Array(
    classOf[Click], classOf[Array[Click]],           // arrays are separate classes
    classOf[Session], classOf[Array[Session]]
  ))

val spark = SparkSession.builder.config(conf).getOrCreate()

The equivalent without code changes is spark.kryo.classesToRegister, a comma-separated list of class names. For more control, implement a KryoRegistrator and name it in spark.kryo.registrator. A registrator can attach a custom serializer to a hot class and can delegate classes Kryo cannot handle, such as some classes with custom writeObject logic, to Kryo's own Java-serialization bridge.

import com.esotericsoftware.kryo.{Kryo, Serializer}
import com.esotericsoftware.kryo.io.{Input, Output}
import org.apache.spark.serializer.KryoRegistrator

// A hand-written serializer for a hot class: fixed field order, no field metadata.
class ClickSerializer extends Serializer[Click] {
  override def write(k: Kryo, out: Output, c: Click): Unit = {
    out.writeLong(c.userId, true)     // true = optimise for positive values if the Output uses varints
    out.writeInt(c.itemId, true)
    out.writeLong(c.ts, true)
    out.writeString(c.page)
  }
  override def read(k: Kryo, in: Input, t: Class[Click]): Click =
    Click(in.readLong(true), in.readInt(true), in.readLong(true), in.readString())
}

class AppRegistrator extends KryoRegistrator {
  override def registerClasses(k: Kryo): Unit = {
    k.register(classOf[Click], new ClickSerializer)
    k.register(classOf[Array[Click]])
    k.register(classOf[LegacyThing],                 // only Java-serializable: delegate
      new com.esotericsoftware.kryo.serializers.JavaSerializer())
  }
}
// spark-submit --conf spark.kryo.registrator=com.example.AppRegistrator ...

Registration order matters. Ids are assigned in registration order, and the writer and the reader must agree on them. Within one application that is automatic, because driver and executors run the same registration code. It becomes a problem only if serialized bytes outlive the application, which is covered under failure modes.

Worked example: measuring before and after

Do not trust a rule of thumb for your data; measure it. The snippet below writes 1,000 small records one by one to a serialization stream, the way a shuffle writes them, with Java serialization, unregistered Kryo and registered Kryo, and prints the byte counts. Serializing one array instead would hide the per-record class-name cost. Run it in a Spark shell with your real record classes.

import java.io.ByteArrayOutputStream
import org.apache.spark.serializer.{JavaSerializer, KryoSerializer, Serializer}

def streamSize(ser: Serializer, recs: Seq[Click]): Int = {
  val bytes = new ByteArrayOutputStream()
  val out   = ser.newInstance().serializeStream(bytes)
  recs.foreach(r => out.writeObject(r))        // one top-level object per record, like a shuffle
  out.close()
  bytes.size()
}

val sample = (1 to 1000).map(i => Click(i, i % 500, 1727650000000L + i, "/p/" + i))
val base   = new SparkConf(false)
println("java              " + streamSize(new JavaSerializer(base), sample))
println("kryo unregistered " + streamSize(new KryoSerializer(base.clone()), sample))
println("kryo registered   " + streamSize(new KryoSerializer(
  base.clone().registerKryoClasses(Array(classOf[Click]))), sample))

What to expect, reasoning from the encodings rather than quoting a benchmark: Java serialization pays a class descriptor once per stream plus per-object overhead and full-width numbers. Unregistered Kryo repeats the class name per record, so for tiny records it can be surprisingly close to Java. Registered Kryo replaces the name with a one- or two-byte id. The gap is largest for many small objects and smallest for a few large arrays of primitives, where both formats are dominated by the raw data.

Then measure the job, not just the bytes. In the Spark UI compare, for the same input, Shuffle Write Size and Shuffle Read Size on the stages that shuffle RDDs, the Size in Memory of serialized cached RDDs on the Storage tab, and task time including GC. If the shuffle is DataFrame-based, expect no change; that is a sign the job did not need Kryo, not that Kryo failed.

Buffers and the size limits that bite

Each executor core gets a Kryo output buffer that starts at spark.kryoserializer.buffer (64k by default) and grows as needed up to spark.kryoserializer.buffer.max (64m by default, and it must be below 2048m). When Spark serializes into a stream, as in shuffle writes, the buffer flushes to the stream and large totals are fine. The ceiling matters wherever Spark serializes one value into a single buffer rather than a stream. Reports usually involve a very large individual record or a large value sent back to the driver, and the Kryo buffer overflow error names buffer.max as the setting to raise.

Raising the ceiling fixes the symptom. The better fix is often to avoid shipping a huge single object: write results out instead of collecting them, collect aggregates rather than raw rows, or split the value. The initial buffer size rarely needs tuning; the tuning guide's advice to raise it for large objects is about avoiding repeated growth, but the limit that fails jobs is the maximum.

Unsafe IO, relocation and the shuffle path

spark.kryo.unsafe switches Kryo to its unsafe-memory based input and output classes, which the configuration page describes as substantially faster. The Spark 4.2 configuration page lists its default as true; earlier releases documented a different default, so check the page for the version you run before assuming either. The encoded bytes from unsafe and safe modes are not interchangeable, which only matters if bytes are stored and read by a different configuration.

Kryo also matters for which shuffle writer Spark picks. The serialized sort-based writer can sort serialized records without deserializing them, but only if the serializer supports relocating serialized objects in the stream. Kryo does, provided its auto-reset behaviour is left on, which is Spark's default; Java serialization does not. That is one more reason RDD shuffles speed up under Kryo beyond the smaller bytes; the shuffle article explains the writer choice.

Kryo and Datasets: Encoders.kryo

Datasets need an encoder. For case classes and primitives Spark derives one that maps fields to real columns. For a class it cannot derive, you can fall back to Encoders.kryo[T], which serializes the whole object into a single binary column.

import org.apache.spark.sql.Encoders
import spark.implicits._

val opaque  = spark.createDataset(clicks)(Encoders.kryo[Click])   // one binary column "value"
val typed   = clicks.toDS()                                        // product encoder: 4 real columns

opaque.printSchema()   // value: binary
typed.printSchema()    // userId: long, itemId: int, ts: long, page: string
typed.filter($"itemId" === 7).select($"userId")   // prunable, pushdown-able; opaque is not

That fallback works but costs a lot: the optimizer sees one opaque column, so there is no column pruning, no predicate pushdown into Parquet, no readable schema, and every access deserializes the whole object. Treat Encoders.kryo as a bridge for a type you do not control, and convert to a case class at the boundary when you can.

Failure modes seen in practice

SymptomCauseFix
IllegalArgumentException: class is not registeredregistrationRequired=true and a class (often an array or a nested collection type) was missedRegister the named class; repeat until clean
ClassNotFoundException for the registratorThe registrator's jar is not on executor classpathsShip it with --jars or in the application jar
Kryo buffer overflow naming buffer.maxOne value serialized into a single buffer exceeded buffer.max, often a huge record or a large resultAvoid shipping huge single values; if necessary raise buffer.max below 2048m
StackOverflowError while serializingDeep or cyclic graph with referenceTracking disabled, or very deep recursionRe-enable reference tracking; flatten deep structures
Wrong or corrupt values after a deployKryo bytes stored across application versions (checkpointed or externally stored RDD data) read with different class definitions or registration orderNever persist Kryo bytes long term; use Parquet, Avro or Protobuf for durable data
No improvement at allThe job is DataFrame-based, so rows never used spark.serializerExpected; look elsewhere for the bottleneck

The storage rule deserves emphasis. Kryo's default field serializer has no schema evolution story you should rely on: add, remove or reorder fields and old bytes may fail to read, or read wrongly. Kryo is a wire and cache format for one running application, not a storage format.

Trade-offs and when not to bother

Kryo costs a little setup and a registration discipline, and it can break on unusual classes that Java serialization handles. Against that, RDD-heavy jobs typically shuffle and cache fewer bytes, spend less CPU in serialization and create less garbage. If your pipeline is DataFrames end to end, leave the default and spend the effort on partitioning and joins, for example broadcast joins. If it shuffles or caches RDDs of your own classes, turn Kryo on, require registration, and measure.

What to do next

  1. List which stages of your job shuffle or cache RDDs of your own classes; if none do, stop here.
  2. Set spark.serializer to KryoSerializer and spark.kryo.registrationRequired=true in a test environment, and register every class the errors name, including array classes.
  3. Run the size measurement snippet with your real records and record Java, unregistered and registered sizes.
  4. Compare shuffle write size, serialized cache size and task time for one representative run before and after.
  5. Leave reference tracking on unless you have verified your records contain no shared or cyclic objects.
  6. Replace any collect() of large values with a write, and raise spark.kryoserializer.buffer.max only if a single value truly needs it.
  7. Make sure no Kryo-serialized bytes are stored beyond the life of the application.
Key takeaway: spark.serializer controls how Spark encodes RDD data in shuffles, serialized caches, broadcasts and task results, not task closures and not DataFrame rows. Kryo shrinks and speeds up those paths by replacing Java's verbose descriptors with registered integer ids and compact field encodings, and it enables the serialized sort shuffle path. Turn it on for RDD-heavy jobs, require registration so no class silently pays for its name, keep reference tracking unless you are sure, handle buffer overflows by not shipping huge single values, and never keep Kryo bytes as durable data.