What are internal repartition and changelog topics in Kafka Streams, and how do they relate to scaling?
answer
- repartition = re-shuffle by new key
- changelog = compacted backup of store
- prefixed by application.id
- repartition partitions = downstream task count
- changelog enables restore/standby
basics
~20 sStreams auto-creates two kinds of internal topics: repartition topics (to re-shuffle data by a new key before stateful ops) and changelog topics (a durable backup of each state store). Their partition counts mirror the input, which also bounds parallelism.
solid answer
~50 sKafka Streams transparently creates internal topics, named with the application.id prefix. Repartition topics appear when an operation needs data grouped by a different key than the input is partitioned on — e.g. after selectKey/map before a groupBy/join. Streams writes re-keyed records to a repartition topic so co-partitioning holds, then reads them back; this splits the topology into sub-topologies. Changelog topics are compacted topics that durably back each stateful store (RocksDB); every store update is logged so the store can be restored or a standby kept warm. For scaling, both matter: repartition topics' partition count determines the task count of the downstream sub-topology, and because they're sized from the source, they propagate the input partition ceiling. Changelog topics are what make state recovery and standby replicas possible, but also add write amplification. You generally don't create these by hand; Streams manages them, and you tune them via topic-level configs or StreamsConfig overrides.
go deeper
Know the two internal topic types and that Streams creates them automatically.
Explain when a repartition topic appears (re-keying) and that changelog topics back state stores.
Connect repartition partition counts to downstream task counts and changelogs to failover/standby mechanics.
Reason about write amplification, topology.optimization, ACLs/pre-creation, and reset-tool hygiene at scale.
## Why internal topics exist Kafka Streams sometimes needs Kafka itself as scratch space. It auto-creates **internal topics**, all prefixed with your `application.id`, in two flavors. ### 1. Repartition topics Many stateful operations require **co-partitioning**: records with the same key must be in the same partition index across all inputs, and grouped/joined by the *current* key. If you change the key mid-topology (e.g. `selectKey`, `map`, `groupBy` on a new field), the data is no longer partitioned by that key. Streams fixes this by inserting a **repartition topic**: it writes the re-keyed records out to this topic (partitioned by the new key), then reads them back in. This: - Restores correct co-partitioning for the downstream operation. - **Splits the topology** at that point into a new **sub-topology** (the upstream and downstream halves run as separate task sets). The repartition topic's **partition count** is derived from the topology requirements (typically matching the source it must co-partition with). That count becomes the **task count of the downstream sub-topology** — so repartition topics carry the partition-count ceiling forward through the pipeline. ### 2. Changelog topics Every **stateful store** (RocksDB-backed) has a corresponding **changelog topic**: a **log-compacted** Kafka topic where every key/value update to the store is also appended. It is the **durable source of truth** for the store. Its roles: - **Restoration**: a task moving to a fresh instance replays the changelog to rebuild its store. - **Standby replicas / warmup replicas**: these consume the changelog to keep hot copies in sync. Changelog topics for KTables built directly from a source topic can use the source as their changelog (the **'source topic optimization'**, `topology.optimization`), avoiding a separate topic. ## How this ties to scaling - **Repartition topic partition count = downstream parallelism ceiling.** You can't get more tasks downstream than the repartition topic has partitions, and that count is propagated from the upstream input. So increasing parallelism deep in a pipeline still traces back to the original input partition count. - **Changelog topics enable safe scaling/failover.** Without them, moving a stateful task between instances during a rebalance or scaling event would lose state. They are the backbone of cooperative rebalancing's state locality, standby replicas, and warmup replicas. ## Operational notes / edge cases - Naming: `<application.id>-<storeOrOperatorName>-repartition` / `-changelog`. - They're auto-created by Streams (the app needs topic-create ACLs, or you pre-create them with matching partitions). **Mismatched partition counts** cause startup failures. - Changelog topics add **write amplification** — every state update is an extra Kafka write — so heavy aggregations cost broker I/O. - Deleting these topics out-of-band corrupts the app; use the **application reset tool** (`kafka-streams-application-reset`) to clean them safely. - `topology.optimization=all` can reduce the number of repartition topics by reusing them.
- Why does a selectKey followed by groupByKey trigger a repartition topic?selectKey changes the record key, so the data is no longer partitioned by that new key. groupByKey needs all records for a key in one partition, so Streams inserts a repartition topic to re-shuffle by the new key, restoring co-partitioning and splitting the topology into a new sub-topology.
- What happens if a changelog topic's partition count doesn't match the store's task count?Streams refuses to start (or behaves incorrectly) — internal topic partition counts must match the topology's expectation. You fix it via the application reset tool or by recreating the topic with the correct partition count.
saying these in an interview costs you the question
- Saying you should manually create/delete internal topics during normal operation
- Claiming changelog topics are for throughput rather than durability/recovery
- Thinking repartition topics let you exceed the input partition count's parallelism ceiling