How do you inspect a Kafka Streams topology, and why does naming DSL operators and stateful stores matter for topology compatibility?
answer
- describe() -> TopologyDescription: sub-topologies, sources, processors+stores, sinks
- Auto node names = position index -> fragile
- Internal topics derive names: -repartition, -changelog
- Name via Named/Materialized/Repartitioned/Grouped/Joined
- Snapshot describe() in a test to catch breaking changes
- Naming = rolling-upgrade / topology compatibility
basics
~10 sCall topology.describe() to print the processor graph (sub-topologies, sources, sinks, stores). Naming operators (via Named/Materialized/Repartitioned/Joined) gives stable node and internal-topic names, so refactoring the DSL doesn't change changelog/repartition topic names and break state restore.
solid answer
~50 s`builder.build().describe()` returns a `TopologyDescription` whose `toString()` lists each **Sub-topology** with its source topics, processor nodes, attached **state stores**, and sink topics — you paste it into the kafka-streams-viz tool to render the DAG. The crucial operational reason to **name** operators is **topology stability across deployments**. By default, Streams auto-generates node names like `KSTREAM-AGGREGATE-0000000003` from the operator's **position index** in the graph. Internal **repartition** and **changelog** topic names are derived from those auto-names. If you insert or reorder a DSL operator, the indices shift, the internal topic names change, and on restart Streams can't find the existing changelog/repartition topics — state restore breaks or reprocessing occurs. Giving explicit names via `Named.as(...)`, `Materialized.as(...)`, `Repartitioned.as(...)`, `Grouped.as(...)`, and `Joined.as(...)` **pins** node and internal-topic names so you can evolve the DSL while staying **topology-compatible** with already-deployed state. This is essential for rolling upgrades of stateful apps.
go deeper
Know topology.describe() prints the processor graph (sources, processors, sinks).
Read describe() output, identify sub-topologies and stores, and know stores back changelog topics.
Explain why auto-names are position-derived and how naming pins repartition/changelog topic names.
Establish a team practice: name all operators/stores, snapshot-test describe(), and manage topology compatibility for rolling upgrades and optimization toggles.
## Inspecting a topology ```java Topology topology = builder.build(); System.out.println(topology.describe()); ``` `Topology#describe()` returns a **`TopologyDescription`**. Its string form groups the graph into **`Sub-topology: N`** blocks, each listing: - **Source** nodes (the input topics, including internal repartition topics), - **Processor** nodes (with the **state stores** they attach, in `stores: [...]`), - **Sink** nodes (output topics, including internal repartition topics). Disconnected sub-topologies and repartition boundaries are visible here. The community **kafka-streams-viz** / Streams TopologyViz renders this text into a DAG diagram — invaluable for reviewing where repartitions and stores live. ## The naming problem (the principal-level point) Kafka Streams **auto-names** every processor node using a **monotonic index based on its position** in the DSL graph, e.g.: ``` KSTREAM-SOURCE-0000000000 KSTREAM-KEY-SELECT-0000000001 KSTREAM-AGGREGATE-0000000003 ``` Two classes of **internal topics** derive their names from these node names: 1. **Repartition topics**: `<application.id>-<node-name>-repartition` 2. **Changelog topics** (backing state stores): `<application.id>-<store-name>-changelog` ### Why auto-names are fragile The index is assigned by **traversal order**. If you later **insert a new operator**, **reorder branches**, or **upgrade a library version** that changes graph construction, the indices for *downstream* nodes **shift**. That changes the derived **internal topic names**. On the next deployment: - Streams looks for changelog/repartition topics under the **new** names, finds them empty/absent, and **can't restore** the existing state — it may reprocess from scratch or, worse, leave the old topics orphaned holding the real state. - This is a classic cause of broken **rolling upgrades** of stateful Streams apps. ### The fix: name your operators and stores Give **explicit, stable names** so node and internal-topic names don't depend on position: - `Named.as("...")` — on stateless operators (filter, map, peek, branch, etc.) and as an overload on many operators. - `Materialized.as("store-name")` — names a state store (and thus its `-changelog` topic). - `Repartitioned.as("...")` — names an explicit repartition topic. - `Grouped.as("...")` — names the repartition topic created by `groupBy`. - `Joined.as("...")` — names join repartition topics/stores. - `StreamJoined.as("...")` — names the join's state stores. With everything named, you can refactor the DSL (insert filters, reorder stages) and remain **topology-compatible** with state already on disk and in changelog topics. ### Verifying compatibility Because `describe()` is deterministic for a given graph, teams **snapshot the describe() output in a test** and diff it across commits — a changed internal-topic name in the diff flags a breaking, non-rolling-upgradeable change before it ships. ### Related: topology optimization `builder.build(props)` with `topology.optimization=all` (a.k.a. `StreamsConfig.TOPOLOGY_OPTIMIZATION_CONFIG`) can **reuse the source topic as the changelog** and collapse redundant repartition/through topics — which also changes the graph, so it must be set consistently and is another reason to verify `describe()` across changes. ### Edge cases - Names must be **unique** within the topology; duplicates throw at build time. - Naming the **store** fixes the changelog name but you should also name the **repartition** producers (via Grouped/Joined/Repartitioned) since those are separately derived. - Changing `application.id` changes ALL internal topic names (it's the prefix) — that's a full reset, distinct from the position-index problem.
- Why can inserting a new filter operator in the middle of a topology break state restore on the next deploy?Auto-generated node names use a position index. Inserting an operator shifts the indices of downstream nodes, which changes the derived names of their changelog and repartition topics. On restart Streams looks for state under the new topic names, can't find the existing data, and reprocesses or loses access to the old state — unless those operators/stores were explicitly named.
- How can a team detect a topology-incompatible change before deploying it?Snapshot topology.describe() output in a unit test and diff it across commits; any change to internal (repartition/changelog) topic names appears in the diff. Combined with explicitly naming operators and stores, this catches non-rolling-upgradeable changes in CI.
- Which builder objects let you name internal topics/stores?Named (stateless ops), Materialized (state stores / changelog), Repartitioned (explicit repartition topic), Grouped (groupBy repartition), Joined / StreamJoined (join repartition topics and stores).
saying these in an interview costs you the question
- Saying topology naming is purely cosmetic — it determines internal topic names and thus state restore.
- Claiming changing application.id is equivalent to naming operators — application.id changes ALL topic prefixes and resets state.
- Believing describe() executes the topology or reads from Kafka — it only prints the static graph.
- Assuming auto-generated names are stable across DSL refactors — they shift with operator position.