skip to content

How do you configure Avro Serdes in a Kafka Streams application, and how does it differ from a plain @KafkaListener?

level: seniorimportance: should knowfreq 45%

answer

  1. Serde = serializer + deserializer in one
  2. SpecificAvroSerde / GenericAvroSerde
  3. schema.registry.url at StreamsConfig top level
  4. manual serde.configure(map, isKey)
  5. default.deserialization.exception.handler not ErrorHandlingDeserializer

basics

~10 s

Streams uses Serde objects (a serializer+deserializer pair), not separate serializer/deserializer classes. Set default.value.serde to SpecificAvroSerde (or GenericAvroSerde), give it schema.registry.url, and pass it explicitly in Consumed/Produced for repartition/state-store topics.

solid answer

~40 s

Kafka Streams serializes at many internal points (repartitions, state-store changelogs, joins), so it works with Serde objects — a bundled serializer+deserializer. You set DEFAULT_VALUE_SERDE_CLASS_CONFIG to io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde (or GenericAvroSerde) and DEFAULT_KEY_SERDE_CLASS_CONFIG similarly. Because the registry URL must reach the serde, you also put schema.registry.url in the StreamsConfig so it's passed to the serde's configure(). For per-operator control you supply the serde explicitly via Consumed.with(...)/Produced.with(...)/Materialized.with(...) — configuring a programmatically-created SpecificAvroSerde requires calling serde.configure(Map.of("schema.registry.url", url), isKey). Deserialization errors don't crash a listener; instead Streams uses default.deserialization.exception.handler (LogAndContinue vs LogAndFail), and production errors use default.production.exception.handler. So the abstraction (Serde not serde class), the config keys, and the error-handling hook all differ from a @KafkaListener.

go deeper

for a junior

Know Streams uses Serde objects and needs schema.registry.url.

for a middle

Name SpecificAvroSerde/GenericAvroSerde and the default.*.serde config keys.

for a senior

Configure per-operator serdes and the deserialization exception handler; explain the manual configure() pitfall.

for a principal

Map the full difference from listener-based serdes including internal-topic serialization and DLT strategy in Streams.

## Serde vs serializer/deserializer A plain consumer needs a `value.deserializer`; a plain producer needs a `value.serializer`. Kafka Streams both reads and writes — and writes a lot internally (repartition topics, changelog topics for state stores, join intermediates). So Streams uses a **Serde** (`org.apache.kafka.common.serialization.Serde`): a single object exposing both a `serializer()` and a `deserializer()`. Confluent ships Avro serdes for Streams: `SpecificAvroSerde` and `GenericAvroSerde` (package `io.confluent.kafka.streams.serdes.avro`). ## Default configuration In your `StreamsConfig`/properties: ``` default.key.serde = io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde default.value.serde = io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde schema.registry.url = http://localhost:8081 ``` When Streams instantiates a default serde via its no-arg constructor, it calls `configure(configs, isKey)` passing the full StreamsConfig — which is why `schema.registry.url` must be present at the top level so the serde can find the registry. ## Per-operator serdes The default serde is overridden anywhere you pass one explicitly: - `Consumed.with(keySerde, valueSerde)` on `stream()`/`table()`. - `Produced.with(...)` on `to()`/`through()`. - `Materialized.with(...)` / `Grouped.with(...)` / `Joined.with(...)` for state stores and joins. If you construct a `SpecificAvroSerde` yourself with `new SpecificAvroSerde<>()`, it is **not** configured — you must call `serde.configure(Map.of("schema.registry.url", url), isKey)` before use, where `isKey` is true for key serdes. Forgetting this is a classic NPE/registry-URL-missing bug. ## Error handling differs fundamentally There's no listener container, so the poison-pill/ErrorHandlingDeserializer machinery doesn't apply. Instead: - `default.deserialization.exception.handler` — `LogAndContinueExceptionHandler` (skip the bad record, keep processing) or `LogAndFailExceptionHandler` (default; shut the task down). You can implement a custom handler to DLT. - `default.production.exception.handler` — handles serialization/produce failures on output. - (Newer Streams also exposes a processing exception handler for user-code errors.) ## Summary of differences from @KafkaListener | Aspect | @KafkaListener | Kafka Streams | |---|---|---| | Abstraction | serializer + deserializer classes | Serde objects | | Registry URL plumbing | consumer/producer props | StreamsConfig top level + serde.configure() | | Typed output | specific.avro.reader=true | SpecificAvroSerde vs GenericAvroSerde | | Bad-record handling | ErrorHandlingDeserializer + DefaultErrorHandler + DLT | default.deserialization.exception.handler | ## Edge cases - Internal repartition/changelog topics are serialized with the resolved serdes; a missing registry URL surfaces deep inside the topology, not at startup. - SpecificAvroSerde needs the generated classes on the classpath, same as specific.avro.reader=true.

  • You did new SpecificAvroSerde<>() and pass it to Consumed.with — why does it fail with a missing registry URL?
    A programmatically constructed serde isn't auto-configured. You must call serde.configure(Map.of("schema.registry.url", url), isKey) before passing it, or it has no registry client.
  • How do you skip undecodable records in Streams instead of crashing the task?
    Set default.deserialization.exception.handler to LogAndContinueExceptionHandler (the default LogAndFail stops the task), or write a custom handler that routes the bad record to a DLT.

saying these in an interview costs you the question

  • Trying to use ErrorHandlingDeserializer in Streams — it's a listener-container construct, not a Streams hook.
  • Passing a manually-built SpecificAvroSerde without calling configure().
  • Assuming Streams only deserializes at the source — it serializes internal repartition/changelog topics too.

context