In a Spring-managed Kafka Streams app, how do you handle uncaught stream exceptions and customize the client before it starts?
answer
- StreamsUncaughtExceptionHandler: REPLACE_THREAD / SHUTDOWN_CLIENT / SHUTDOWN_APPLICATION
- set on StreamsBuilderFactoryBean before start
- record-level = deserialization/production exception handler props
- KafkaStreamsCustomizer = last-mile before start()
- StateListener for RUNNING->ERROR alerting
basics
~20 sSet a StreamsUncaughtExceptionHandler on the StreamsBuilderFactoryBean to decide whether to replace the thread, shut down the client, or shut down the whole app. For anything you can't set via properties, register a KafkaStreamsCustomizer that receives the KafkaStreams before start().
solid answer
~30 sKafka Streams processing runs on internal StreamThreads. If a thread hits an unrecoverable error, the StreamsUncaughtExceptionHandler decides the response: REPLACE_THREAD (restart the failed thread), SHUTDOWN_CLIENT (stop this instance), or SHUTDOWN_APPLICATION (signal all instances of the application.id to stop). In Spring you register it via StreamsBuilderFactoryBean.setStreamsUncaughtExceptionHandler, before startup. Record-level deserialization/production errors are handled separately through StreamsConfig properties: DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG (e.g. LogAndContinue) and DEFAULT_PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG. For arbitrary last-mile tweaks to the constructed client, use setKafkaStreamsCustomizer(KafkaStreamsCustomizer), which hands you the KafkaStreams instance just before start(). In Boot you post-process the factory bean with a StreamsBuilderFactoryBeanConfigurer bean. A StateListener lets you observe RUNNING/ERROR transitions for alerting.
code
java · 20 linesimport org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler;
import org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse;
import org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class StreamsErrorConfig {
// Spring Boot applies this to the default factory bean before start.
@Bean
public StreamsBuilderFactoryBeanConfigurer configurer() {
return fb -> {
fb.setStreamsUncaughtExceptionHandler(ex ->
StreamThreadExceptionResponse.REPLACE_THREAD);
fb.setKafkaStreamsCustomizer(kafkaStreams ->
kafkaStreams.setGlobalStateRestoreListener(null)); // last-mile hook
};
}
}go deeper
Know there's a way to handle stream errors and that Spring lets you set a handler on the factory bean.
Distinguish thread-fatal (uncaught handler) from record-level (deserialization/production) errors.
Explain the three uncaught responses, KafkaStreamsCustomizer, and StreamsBuilderFactoryBeanConfigurer.
Design resilience: when to REPLACE_THREAD vs SHUTDOWN_*, DLQ strategy, and multi-instance coordination.
**Two distinct error layers — don't conflate them.** 1. **Thread-fatal exceptions (uncaught).** A `KafkaStreams` client runs processing on internal `StreamThread`s. An exception that escapes all processing (e.g. a serialization bug in your topology, a broker auth failure) bubbles up as *uncaught*. Kafka Streams (2.8+) uses a `StreamsUncaughtExceptionHandler` returning one of: - `REPLACE_THREAD` — kill the failed thread and spin up a replacement; the app keeps running with the remaining threads. Good for transient thread-local faults. - `SHUTDOWN_CLIENT` — transition *this* `KafkaStreams` instance to `ERROR` and stop it; other instances of the same `application.id` keep going. - `SHUTDOWN_APPLICATION` — instruct *all* instances sharing the `application.id` to shut down (coordinated). Use for truly fatal, non-recoverable conditions. In Spring you set it on the factory bean: `factoryBean.setStreamsUncaughtExceptionHandler(...)`. Default behavior (if unset) shuts the client down. 2. **Record-level exceptions (handled, per-record).** These are configured as `StreamsConfig` properties, not on the factory bean: - `DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG` — e.g. `LogAndContinueExceptionHandler` (skip poison-pill records) vs. `LogAndFailExceptionHandler` (default, fail fast). Set it via your `KafkaStreamsConfiguration` map. - `DEFAULT_PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG` — decide `CONTINUE` vs `FAIL` when a record can't be produced downstream. - Processing exceptions handler (newer versions) for errors inside processor logic. **Customizing the client before start.** Some things aren't expressible as properties — e.g. registering a global `StateRestoreListener`, or wrapping the client. `StreamsBuilderFactoryBean.setKafkaStreamsCustomizer(KafkaStreamsCustomizer)` gives you a callback with the fully-constructed `KafkaStreams` *immediately before* `start()`. There's also `setInfrastructureCustomizer(KafkaStreamsInfrastructureCustomizer)` to hook the `StreamsBuilder` / `Topology` during build (e.g. add a global store). **Observability.** `setStateListener(KafkaStreams.StateListener)` fires on state transitions (`CREATED` → `REBALANCING` → `RUNNING`, or → `ERROR`). Wire it to metrics/alerts so an instance dropping to `ERROR` pages someone. **Where to register in Spring.** Because the default factory bean is auto-created, you post-process it. Options: (a) inject it by the reserved name `defaultKafkaStreamsBuilder` and set handlers in a `@Bean`/`@PostConstruct`; (b) in Spring Boot, declare a `StreamsBuilderFactoryBeanConfigurer` bean — Boot applies it to the factory bean before start. All configuration must happen *before* start, so do it in a bean that runs during context init, not lazily afterward. **Gotchas.** (1) `SHUTDOWN_APPLICATION` only stops *other instances* if they're the same `application.id` and reachable via the group — it's cooperative, not instantaneous. (2) `LogAndContinue` silently drops bad records — great for resilience, dangerous for correctness if you actually needed those records; pair it with metrics/DLQ. (3) Setting the uncaught handler *after* `start()` has no effect on the already-running client. (4) The uncaught handler and the deserialization handler solve different problems; a poison-pill record won't hit the uncaught handler if a deserialization handler swallows it first. **When to use which.** Transient/thread-local → `REPLACE_THREAD`. Instance-level unrecoverable → `SHUTDOWN_CLIENT`. Bad-config/poison-topology affecting every instance → `SHUTDOWN_APPLICATION`. Bad individual records → deserialization/production exception handlers (skip or DLQ), not the uncaught handler.
- A single malformed record keeps crashing your stream. Which mechanism fixes it — the uncaught handler or something else?Configure DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG (e.g. LogAndContinue) as a StreamsConfig property, ideally routing to a DLQ. The uncaught handler is for thread-fatal errors, not per-record poison pills.
- What does SHUTDOWN_APPLICATION do differently from SHUTDOWN_CLIENT?SHUTDOWN_CLIENT stops only this instance; other instances of the same application.id keep running. SHUTDOWN_APPLICATION cooperatively signals all instances sharing the application.id to stop.
saying these in an interview costs you the question
- Using the uncaught exception handler to skip poison-pill records
- Setting handlers after the KafkaStreams client has already started
- Believing SHUTDOWN_CLIENT stops every instance of the application