How do you define a Kafka Streams topology (a KStream/KTable) as a Spring bean, and who starts it?
answer
- @Bean method takes StreamsBuilder param
- builder.stream() -> KStream, builder.table() -> KTable
- never call new KafkaStreams() / start() yourself
- StreamsBuilderFactoryBean is SmartLifecycle, autoStartup=true
- getKafkaStreams() for the live client
basics
~20 sWrite a @Bean method that takes a StreamsBuilder parameter. Spring injects the default StreamsBuilder, you build your KStream/KTable pipeline on it, and the StreamsBuilderFactoryBean builds the topology and auto-starts the KafkaStreams instance after the context refreshes.
solid answer
~40 sOnce @EnableKafkaStreams is active, Spring exposes the default StreamsBuilder from the StreamsBuilderFactoryBean. You declare a @Bean method with a StreamsBuilder parameter (Spring injects it), then define your pipeline: builder.stream("input") gives a KStream; you map/filter/aggregate and .to("output"). A KStream models each record as an independent event; a KTable models the latest value per key (a changelog/upsert view). You do NOT call new KafkaStreams(...) or start() yourself — the StreamsBuilderFactoryBean is a SmartLifecycle that collects everything added to the shared StreamsBuilder, builds the Topology, constructs the KafkaStreams client, and starts it automatically (autoStartup defaults to true) when the ApplicationContext finishes refreshing. The method's return type is often KStream, but returning void is fine since the builder is mutated in place.
code
java · 24 linesimport org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.Produced;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafkaStreams;
@Configuration
@EnableKafkaStreams
public class UppercaseTopology {
// Spring injects the shared default StreamsBuilder; the factory bean
// builds the Topology and starts KafkaStreams automatically.
@Bean
public KStream<String, String> uppercasePipeline(StreamsBuilder builder) {
KStream<String, String> input =
builder.stream("input-topic", Consumed.with(Serdes.String(), Serdes.String()));
input.mapValues(v -> v == null ? null : v.toUpperCase())
.to("output-topic", Produced.with(Serdes.String(), Serdes.String()));
return input;
}
}go deeper
Know that a @Bean with a StreamsBuilder param defines the pipeline and Spring starts it.
Explain the shared builder, autoStartup lifecycle, and KStream vs KTable at wiring level.
Cover accessing getKafkaStreams(), Serde defaults, and why manual start is wrong.
Discuss build-time-fixed topology, multi-builder separation, and failure surfacing at refresh.
**The core idea.** With `@EnableKafkaStreams`, Spring gives you a single shared `StreamsBuilder` (from the default `StreamsBuilderFactoryBean`). Any `@Bean` method that declares a `StreamsBuilder` parameter receives that same instance; whatever nodes you add to it become part of one `Topology`. **KStream vs KTable (minimal definitions, since theory is a sibling topic — here just enough to wire correctly).** A `KStream<K,V>` is an unbounded stream where every record is an independent fact/event. A `KTable<K,V>` is a materialized changelog: for each key it holds the latest value (later records upsert earlier ones). You obtain them from the builder: `builder.stream("topic")` → `KStream`; `builder.table("topic")` → `KTable`. Spring's role stops at supplying the builder and running it; the DSL calls are plain Kafka Streams. **Defining the bean.** Typical shape: ```java @Bean public KStream<String, String> pipeline(StreamsBuilder builder) { KStream<String, String> s = builder.stream("input"); s.mapValues(String::toUpperCase).to("output"); return s; } ``` Returning the `KStream` is a convention; the important side effect is that you mutated the injected `builder`. A `void` method that just adds nodes works identically. **Who builds and starts it — the lifecycle.** `StreamsBuilderFactoryBean` implements `SmartLifecycle`. After the context refreshes and all topology beans have contributed to the builder, the factory bean calls `builder.build()` to produce the `Topology`, constructs a `KafkaStreams` object with it plus the merged `StreamsConfig`, and calls `start()`. You never manage this by hand. `autoStartup` (settable on the factory bean) defaults to `true`; set it `false` to defer start and control it programmatically via `factoryBean.start()` / `.stop()`. **Accessing the running client.** Inject the `StreamsBuilderFactoryBean` (by type or the `defaultKafkaStreamsBuilder` name) and call `getKafkaStreams()` to reach the live `KafkaStreams` — e.g. to query state or read `state()`. It returns `null` before start. **Serdes.** Records are bytes; `builder.stream("topic")` uses the default key/value Serde from config (`StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG` / `DEFAULT_VALUE_SERDE_CLASS_CONFIG`) unless you pass `Consumed.with(keySerde, valueSerde)` explicitly. Serde mismatches are a classic runtime failure. **Gotchas.** (1) Don't instantiate `KafkaStreams` yourself — you'll end up with two clients on the same `application.id` and rebalancing chaos. (2) All topology beans share one builder/`application.id`; splitting into truly independent apps requires separate named `StreamsBuilderFactoryBean`s, each with its own config. (3) If a topology bean throws while building, the factory bean fails to start — surfaced at context refresh. (4) The topology is fixed at build time; you can't add nodes after `start()`. **When to use.** Use topology beans for declarative stateful processing (enrichment, joins, aggregations, windowing) that would be awkward and error-prone to hand-roll with plain listeners.
- Why should you never call new KafkaStreams(...).start() inside a topology bean?The StreamsBuilderFactoryBean already builds and starts a KafkaStreams from the shared builder. A second instance on the same application.id creates a duplicate group member and triggers rebalances and state-store lock contention.
- How do you access the live KafkaStreams instance, e.g. to check its state?Inject the StreamsBuilderFactoryBean and call getKafkaStreams(); it returns the running client (or null before startup). From it you can read state() or query state stores.
saying these in an interview costs you the question
- Manually constructing and starting a KafkaStreams instance inside the bean
- Thinking each @Bean gets its own StreamsBuilder rather than the shared one
- Assuming the topology can be modified after the stream has started