skip to content

How does Kafka Streams split a topology into sub-topologies, and what role do repartition topics play?

level: seniorimportance: should knowfreq 55%

answer

  1. Sub-topology boundary = repartition topic
  2. Key-changing op + stateful step triggers repartition
  3. groupBy repartitions; groupByKey does not
  4. mapValues/filter key-preserving (no repartition)
  5. topology.describe() shows Sub-topology blocks + repartition topics

basics

~20 s

A topology is split into sub-topologies at points where data must be re-shuffled by key — repartition topics. Operations that change the key (selectKey, map, groupBy) force a repartition; the topology breaks there into independently-scheduled sub-topologies.

solid answer

~50 s

Kafka Streams compiles your DSL into a `Topology` made of processor nodes connected by streams. Wherever data must be **re-distributed across partitions by a new key**, the runtime inserts an internal **repartition topic** (auto-created, named `<app-id>-<name>-repartition`) and splits the graph into separate **sub-topologies** on either side of it. Triggers are **key-changing operations** followed by a stateful/keyed step: `selectKey`, `map`, `flatMap`, `transform`, or `groupBy` before an aggregation or join. The repartition writes records back to Kafka keyed by the new key so they land in the correct partition for the downstream task. Each sub-topology is **scheduled independently** as its own set of tasks (one task per partition), connected only through Kafka topics — not in-process. You can see all of this with `topology.describe()`, which lists each `Sub-topology:` block and the source/sink topics that join them. Minimizing unnecessary repartitions (e.g. preferring `groupByKey` over `groupBy` when the key is unchanged) reduces extra topics, network, and latency.

go deeper

for a junior

Know that a repartition (re-shuffle by key) splits the topology and that changing the key causes it.

for a middle

Name which operators trigger repartitions and that sub-topologies are linked by internal topics.

for a senior

Explain lazy repartition flagging, co-partitioning for joins, and how to read describe() output to spot extra topics.

for a principal

Optimize a whole pipeline's repartition count and partitioning for throughput/latency/cost, and decide where to force explicit repartitions to rescale stages.

## What a topology is The DSL is a builder for a **`Topology`**: a directed graph of **processor nodes** (source → stateless/stateful processors → sink). When you call `builder.build()`, the DSL operators are compiled into this graph. ## Why sub-topologies exist Kafka Streams parallelism is **partition-based**: a task owns one partition of each of its input topics. For a **stateful or keyed** operation (aggregation, join) to be correct, all records for a key must reach the **same task**. If an upstream operation **changes the key**, records keyed by the new key may belong to a different partition, so they must be **re-shuffled** through Kafka. The runtime does this by inserting an **internal repartition topic**: - A sink node writes records (re-keyed) to topic `<application.id>-<name>-repartition`. - A source node downstream reads them back, now correctly partitioned by the new key. This **break** in the graph is a **sub-topology boundary**. The graph upstream of the repartition is one sub-topology; downstream is another. They communicate **only through Kafka topics**, never in-process — which is what lets them be scheduled and scaled independently (each sub-topology gets one task per partition of its inputs). ## What triggers a repartition / new sub-topology - **Key-changing operators** that mark the stream as needing repartition: `selectKey`, `map` (can change key), `flatMap`, `transform`/`transformValues` is value-only so it does NOT, but `transform` (full) can. - `groupBy(...)` (changes grouping key) → repartition before the aggregation. By contrast `groupByKey()` keeps the existing key and does **NOT** repartition if the key is unchanged. - A **join** whose inputs aren't co-partitioned will insert repartition topics to align them. - Note: the repartition is created **lazily** — only if a downstream operation actually needs the re-keyed data to be correctly partitioned (the DSL tracks a 'repartition required' flag). A `selectKey` with no downstream stateful op may not materialize a repartition. ## Reading it: Topology#describe() `topology.describe()` returns a `TopologyDescription` whose `toString()` prints blocks like: ``` Sub-topology: 0 Source: KSTREAM-SOURCE-0000 (topics: [input]) Processor: KSTREAM-KEY-SELECT-0001 ... Sink: ...-repartition-sink (topic: app-...-repartition) Sub-topology: 1 Source: ...-repartition-source (topics: [app-...-repartition]) Processor: KSTREAM-AGGREGATE-... (stores: [...]) ``` You paste this into the **Kafka Streams TopologyViz / kafka-streams-viz** tool to render the DAG. ## Performance / design implications - Each repartition adds a **round-trip to Kafka**: extra produce + consume, extra internal topic, more end-to-end latency, and another place to size partitions/retention. - **Avoid needless repartitions**: use `groupByKey` when the key is unchanged; use `mapValues`/`transformValues`/`filter` (value-/predicate-only, key-preserving) instead of `map`/`transform` when you don't actually change the key, because key-preserving operators do **not** set the repartition flag. - Repartition topics are **internal, auto-created, and managed** by Streams (named by `application.id`); their partition count matches the source to preserve co-partitioning. ## Edge cases - Multiple `KStream` sources from `builder.stream()` that are never joined produce **independent sub-topologies** (sub-topology numbering reflects this), even without any repartition. - `through()`/`repartition()` operators let you force an explicit repartition (and thus a sub-topology split) deliberately — useful to scale a downstream stage to a different partition count via `Repartitioned.numberOfPartitions(...)`.

  • How would you reduce the number of repartition topics in a topology?
    Prefer key-preserving operators (mapValues, transformValues, filter) over key-changing ones (map, transform); use groupByKey instead of groupBy when the key is already correct; ensure joined topics are co-partitioned up front so Streams doesn't insert alignment repartitions; and only selectKey when a downstream stateful op truly needs the new key.
  • Do two separate builder.stream() sources that are never joined run in the same sub-topology?
    No — disconnected parts of the graph become separate sub-topologies, each scheduled independently with its own tasks, even though there's no repartition between them.

saying these in an interview costs you the question

  • Claiming sub-topologies communicate in-process — they communicate only via Kafka topics.
  • Saying groupByKey causes a repartition — it does not when the key is unchanged; groupBy does.
  • Believing mapValues/filter trigger repartitions — they are key-preserving.
  • Thinking repartition topics are something you must create manually — Streams auto-creates and manages them.

context