In Kafka Streams, how do you handle deserialization failures, and when would you choose LogAndContinue vs LogAndFail?
answer
- default.deserialization.exception.handler
- CONTINUE = skip+advance, FAIL = stop thread
- LogAndFail is the DEFAULT
- custom handler → DLQ then CONTINUE
- separate ProductionExceptionHandler for output
basics
~10 sSet default.deserialization.exception.handler. LogAndContinueExceptionHandler logs the bad record and skips it (offset advances); LogAndFailExceptionHandler logs and stops the stream thread. Choose Continue for best-effort tolerance, Fail when no record may be silently dropped.
solid answer
~40 sKafka Streams handles record-level deserialization failures via the `default.deserialization.exception.handler` config, which takes a `DeserializationExceptionHandler` implementation. Two are built in: `LogAndContinueExceptionHandler` returns `CONTINUE` — it logs the corrupt record and skips it, letting the offset advance so the topology keeps processing; and `LogAndFailExceptionHandler` (the default) returns `FAIL`, which throws and shuts down the affected `StreamThread`. You choose `LogAndContinue` for best-effort pipelines where occasional malformed records are acceptable and availability matters more than completeness (e.g. clickstream analytics). You choose `LogAndFail` when silently dropping a record is unacceptable — financial or audit data — so a human investigates before any data loss. For nuanced policy you implement a custom handler that quarantines to a DLQ and then returns `CONTINUE`, getting both progress and durability. Note this handler covers **deserialization**; production-side errors use a separate `ProductionExceptionHandler`.
go deeper
Know the two handlers exist: one skips bad records, one stops the stream.
Name the config key, that LogAndFail is default, and the CONTINUE vs FAIL trade-off.
Discuss custom DLQ handlers, silent-loss alerting, and separation from production/processing handlers.
Set org-wide policy mapping data criticality to handler choice and design the DLQ + replay tooling.
## Where this fits Kafka Streams reads records, deserializes them with the source-node Serdes, and runs your topology. A record whose bytes can't be deserialized is the **poison pill** at the source. Unhandled, it would crash the stream thread on every attempt. Streams exposes a pluggable hook for exactly this. ## The config ```properties default.deserialization.exception.handler=org.apache.kafka.streams.errors.LogAndContinueExceptionHandler ``` The value is a class implementing `org.apache.kafka.streams.errors.DeserializationExceptionHandler`, whose `handle(...)` method returns a `DeserializationHandlerResponse`: **`CONTINUE`** (skip the record, advance the offset) or **`FAIL`** (rethrow, stop the thread). ## The two built-ins - **`LogAndFailExceptionHandler`** — the **default**. Returns `FAIL`. The `StreamThread` dies; if all threads die the app stops. Nothing is silently lost: an operator must act. Use when correctness/completeness is non-negotiable. - **`LogAndContinueExceptionHandler`** — returns `CONTINUE`. Logs (at WARN) and **skips** the record. The consumer position advances past it so the topology keeps flowing. Use for best-effort/high-availability pipelines where dropping the occasional malformed record is acceptable. ## Choosing | Concern | LogAndContinue | LogAndFail | |---|---|---| | Availability | High — never stops for bad data | Lower — stops on first bad record | | Data completeness | May silently drop records | No silent loss | | Good for | Analytics, metrics, clickstream | Finance, ledgers, audit, compliance | | Visibility | Only a log line (easy to miss) | Loud failure forces attention | The risk of `LogAndContinue` is **silent data loss** — a log line is easy to overlook, so pair it with metrics/alerting on the skip count. ## Going beyond the built-ins — custom handler / DLQ A custom `DeserializationExceptionHandler` can publish the raw bytes to a **dead-letter topic** (using a small KafkaProducer inside the handler) and then return `CONTINUE`. This gives progress *and* durability — nothing is lost, the thread survives. Newer Streams versions also provide richer error-handling context and a `ProcessingExceptionHandler` for errors *inside* processors (distinct from deserialization). ## Scope and related knobs - This handler is for **inbound deserialization**. Errors when *producing* to output topics are handled by `default.production.exception.handler` (`ProductionExceptionHandler`) — don't conflate them. - It does not help with logic exceptions thrown by your processors; those need a `ProcessingExceptionHandler` (or try/catch inside transformers). ## Replay safety Because the handler's decision is deterministic per record, replaying the same input (e.g. resetting with the application-reset tool) reproduces the same skips/quarantines — reprocessing is safe and idempotent with respect to poison pills.
- What is the main operational risk of LogAndContinue, and how do you mitigate it?Silent data loss — skipped records leave only a log line that's easy to miss. Mitigate by emitting metrics/alerts on skip counts, or by using a custom handler that writes the bad bytes to a DLQ before returning CONTINUE.
- Does default.deserialization.exception.handler catch a NullPointerException thrown inside a map() processor?No. That handler only covers deserialization of inbound records. Errors inside processors need a ProcessingExceptionHandler (newer versions) or in-processor try/catch; production/output errors use ProductionExceptionHandler.
saying these in an interview costs you the question
- Saying LogAndContinue is the default (LogAndFail is the default).
- Claiming the deserialization handler also catches processor logic exceptions (it doesn't — that's ProcessingExceptionHandler).
- Using LogAndContinue for financial/audit data without alerting on silent drops.
- Confusing the deserialization handler with the ProductionExceptionHandler for output-side errors.