skip to content

In Kafka Streams, what is a Serde and how does the Serdes factory class fit in?

level: seniorimportance: should knowfreq 55%

answer

  1. Serde = Serializer + Deserializer pair
  2. Streams reads AND writes (repartition/changelog)
  3. Serdes (trailing s) = static factory class
  4. default.key.serde / default.value.serde
  5. override via Consumed/Produced/Grouped/Materialized.with

basics

~20 s

A Serde<T> bundles a Serializer<T> and a Deserializer<T> into one object, because Streams both reads and writes data. The Serdes factory class provides ready-made ones, e.g. Serdes.String(), Serdes.Long(), and you set defaults with default.key.serde / default.value.serde.

solid answer

~40 s

Kafka Streams continuously reads from and writes to topics (including internal repartition and changelog topics), so it needs both directions at once. A Serde<T> (org.apache.kafka.common.serialization.Serde) is a factory pairing a serializer() and a deserializer() for the same type T. The Serdes utility class (note the trailing 's') exposes static factories for built-ins: Serdes.String(), Serdes.Long(), Serdes.Integer(), Serdes.ByteArray(), Serdes.UUID(), Serdes.Bytes(), plus Serdes.serdeFrom(serializer, deserializer) to wrap a custom pair, and Serdes.WrapperSerde internally. You set application-wide defaults via the StreamsConfig properties default.key.serde and default.value.serde, and you can override per operator by passing Consumed.with(...), Produced.with(...), Grouped.with(...), Materialized.with(...). A mismatch between the configured Serde and the actual bytes surfaces as a SerializationException, optionally routed through a DeserializationExceptionHandler.

go deeper

for a junior

Know a Serde bundles serializer + deserializer and that Serdes.String() exists.

for a middle

Set default.key.serde/value.serde and use Consumed/Produced.with to override.

for a senior

Explain why Streams needs combined serdes (repartition/changelog) and handle per-operator type changes and exception handlers.

for a principal

Design serde strategy across a topology, standardize custom serdes, and govern deserialization-exception handling for resilience.

## Why Streams needs a combined abstraction A plain producer only serializes; a plain consumer only deserializes. **Kafka Streams does both, repeatedly**: it consumes input topics, writes to internal **repartition** topics, persists state to **changelog** topics, and produces output topics. For any given data type it therefore needs a serializer *and* a deserializer, kept together so they can't drift apart. That bundle is a **Serde** (SERializer + DEserializer). ## The Serde interface `org.apache.kafka.common.serialization.Serde<T>` has two accessors: ``` Serializer<T> serializer(); Deserializer<T> deserializer(); ``` Whenever Streams needs to write a `T`, it calls `serializer()`; when it reads, `deserializer()`. The pair guarantees symmetry. ## The Serdes factory (note the 's') `org.apache.kafka.common.serialization.Serdes` is a **utility class of static factory methods** — easy to confuse with `Serde` (no 's'). It returns pre-built Serdes for every built-in type: - `Serdes.String()`, `Serdes.Long()`, `Serdes.Integer()`, `Serdes.Double()`, `Serdes.Short()`, `Serdes.Float()` - `Serdes.ByteArray()`, `Serdes.ByteBuffer()`, `Serdes.Bytes()`, `Serdes.UUID()`, `Serdes.Void()` - `Serdes.serdeFrom(Serializer<T>, Deserializer<T>)` — wrap your own pair into a Serde. - `Serdes.serdeFrom(Class<T>)` — look up a built-in by type. ## Where you wire them in Two levels: 1. **Defaults** via `StreamsConfig`: `default.key.serde` and `default.value.serde` (set to Serde class names). These apply when an operator doesn't specify its own. 2. **Per-operator overrides** using the `…with(...)` parameter objects: `Consumed.with(keySerde, valueSerde)`, `Produced.with(...)`, `Grouped.with(...)`, `Materialized.with(...)`, `Joined.with(...)`, `StreamJoined.with(...)`. This is essential because the key/value types change across a topology (e.g. after a `groupBy` or `map`). ## Failure modes - If the default Serde doesn't match the actual byte content, deserialization throws a `SerializationException`. You can install a `default.deserialization.exception.handler` (e.g. `LogAndContinueExceptionHandler` vs `LogAndFailExceptionHandler`). - On the produce side, `default.production.exception.handler` governs serialization/send failures. - Forgetting a per-operator Serde after a type change silently falls back to the (wrong) default Serde and fails at runtime. ## Custom Serdes For JSON/Avro/Protobuf you either provide a Serde implementation (Confluent's `SpecificAvroSerde`, `KafkaJsonSchemaSerde`, etc.) or build one with `Serdes.serdeFrom(mySerializer, myDeserializer)`. Many custom Serdes read extra config in `configure(Map, isKey)` — note the `isKey` boolean distinguishes key vs value configuration.

  • Why might you override the Serde on a specific operator instead of relying on defaults?
    Because key/value types change across a topology (e.g. after map or groupBy), the global default no longer matches, so you pass Consumed.with/Grouped.with/Materialized.with for that step.
  • What is the difference between Serde and Serdes?
    Serde<T> is the interface bundling a serializer and deserializer; Serdes (with an s) is a utility class of static factory methods that returns prebuilt Serde instances like Serdes.String().

saying these in an interview costs you the question

  • Confusing Serde (interface) with Serdes (factory class).
  • Saying Streams only needs a deserializer or only a serializer — it needs both for repartition/changelog topics.
  • Claiming default.key.serde/default.value.serde cannot be overridden per operator.
  • Forgetting that type changes mid-topology require new Serdes.

context