Kafka Streams Internal Topic Predictor

The internal topics a Streams application will create, derived from its application.id. Create them in advance and you choose their durability rather than inheriting the broker defaults.

Every internal topic name starts with this. It is also the consumer group id, the state directory name and the prefix your ACLs are granted on, which makes it a security boundary as well as an identity.

Paste below, or drop a file anywhere on this panel

Or drop a file anywhere on this panel. Nothing is uploaded: the analysis runs in this tab.

The answer appears here

Paste on the left and press Predict. Nothing leaves this tab.

Examples

Real input you can load into the tool above. Each one shows a different thing going wrong, because that is what the tool is for.

Internal topics

The changelog and repartition topics a Streams application creates on the cluster

order-totals
KSTREAM-AGGREGATE-STATE-STORE-0000000003

A repartition

A key-changing operation before an aggregation, and the topic it forces

KSTREAM-KEY-SELECT-0000000002-repartition

Common mistakes

These are the ones that fail silently. The config is accepted, nothing raises an error, and the consequence arrives later.

  1. Assuming internal topics are cleaned up

    Changelog and repartition topics survive the application. Deleting the app leaves them on the cluster with their retention.

    Instead:Run the application reset tool, then delete the internal topics deliberately.

  2. Ignoring changelog retention

    A changelog is compacted, so it keeps the latest value per key forever. On a large keyspace that is real and permanent disk.

    Instead:Size for it, and set a retention where the store is a windowed one.

  3. Changing the application id casually

    It is part of every internal topic name, so changing it orphans the old topics and rebuilds every store from scratch.

    Instead:Treat it as permanent.

Streams creates topics at runtime, with the broker defaults

The names are fully determined by application.id, so they can be created in advance instead. That is the difference between choosing their durability and inheriting it.

The names are a formula

A changelog is {application.id}-{store}-changelog and a repartition topic is {application.id}-{name}-repartition. Nothing about either is random or version dependent, which means they can be created before the application starts, with the replication factor, cleanup policy and ACLs you chose rather than whatever the broker defaults to.

A changelog is the durable copy of your state

Every write to a state store goes to its changelog topic first, and that topic is what a new instance replays to rebuild the store after a failover. So a changelog at replication factor 1 means the state has one copy, and losing that broker loses it permanently. No number of standby replicas helps, because every standby reads from the same changelog. This is the reason to create them yourself rather than letting the defaults decide.

The two cleanup policies are opposite

A changelog must be compacted, because it has to keep the latest value for every key indefinitely: created with cleanup.policy=delete instead, the state silently loses whatever aged out and a rebuild produces a store missing older keys. A repartition topic must not be compacted, because it is a transient shuffle read once and then worthless, and compacting one keeps every key forever for no reason at all.

The 249 character limit is on the derived name

Not on application.id. A 230 character application.id is comfortably legal and still produces a changelog topic name over the limit, and Streams fails when it tries to create it, at startup, after the application has already connected successfully. The failure names the topic rather than explaining that the prefix is too long.

What this cannot see

It does not read your topology, so you have to name the state stores and repartition points yourself: one per line, with a trailing exclamation mark for a repartition rather than a store. Which operations create a repartition is a topology question, and the short version is that any key-changing operation followed by an aggregation or a join does. For the ACLs these topics need, the ACL generator on this site emits the prefixed grant a Streams application requires.

More kafka tools

Kafka Confluent Wire Format Decoder The five junk bytes in front of your payload Kafka Key to Partition Mapper Which partition does this key land on? Kafka Topic Name Validator Legal, risky, or 249 characters too long? Kafka Replication Safety Checker How many brokers can you lose Kafka Producer Config Linter Will it start, and will it lose a record? Kafka Message Payload Decoder The first five bytes are usually not data Kafka Connect Source Connector Generator tasks.max is a ceiling, not a count Kafka Connect Sink Connector Generator A dead letter queue with no context headers is a pile of records Kafka Connect SMT Chain Builder The order is the transforms list Kafka MirrorMaker 2 Config Generator It renames every topic by default Kafka Partition Reassignment Generator The throttle is not optional Strimzi Kafka Resource Generator Without the cluster label, nothing happens Kafka mTLS Config Generator The certificate is the identity Kafka Schema Registry Config Generator The compatibility direction is your deployment order Kafka Exactly-Once Config Generator Half of it is worse than none Kafka Broker and KRaft Config Generator The internal topics that break a one-broker cluster Kafka Quota Generator Byte rates are per broker, not per cluster Kafka Streams Config Generator application.id is four things at once Kafka Connect Worker Config Generator Security three times, or the tasks fail Kafka Retention and Unit Converter log.retention.hours does not take milliseconds Kafka Timestamp Converter Two sentinels and two meanings Kafka .properties to YAML Converter Dotted keys stay flat Kafka ACL Generator The grant you forgot is on another resource type Kafka Topic Config Generator min.insync.replicas is the one that matters Kafka client.properties Generator The file every CLI tool asks for Kafka Producer Config Generator No password field, on purpose Kafka Consumer Config Generator The commit mode decides the semantics Kafka Disk and Retention Calculator retention.bytes is per partition Kafka Partition Count Calculator The number you can never reduce Kafka Cluster Sizing Calculator The traffic no client metric shows Kafka Consumer Lag Catch-Up Calculator Whether it ever clears, not just when Kafka Producer Batching Calculator linger.ms=0 still batches Kafka Segment and Index Sizing Why retention.ms is a lower bound Kafka Rebalance Duration Estimator What a rolling restart really costs Kafka Cost Estimator Your rates, so nothing goes stale Kafka Config Explorer by Version The answer depends on the release Kafka Default Config Reference What moved under a config you never edited Kafka OAuth Bearer Token Decoder Will Kafka accept it, and can it refresh Kafka Record Header Viewer Headers are a list, not a map Kafka Topic Regex Subscription Tester Kafka matches the whole name Kafka ACL Permission Matrix Viewer DENY beats every ALLOW Kafka Connect Config Validator The mistakes that raise no error Kafka Consumer Group Id Validator Which broker coordinates the group Kafka Partition Assignment Visualizer Leadership is the load, not replicas Kafka Consumer Assignment Visualizer The three assignors disagree Kafka ZooKeeper to KRaft Config Converter The authorizer class nobody changes Kafka Config to Strimzi Half of it belongs elsewhere Kafka Docker Compose Generator (KRaft) Reachable from inside and outside Kafka JAAS Config Decoder The line that stops SASL working Kafka CRC32C Calculator Which CRC, over which bytes Kafka Config Upgrade Checker What breaks when you upgrade Kafka Kafka Config Diff Which change actually changed something Kafka Consumer Config Linter Why the group rebalances, and where the records went Kafka Avro Schema Validator The defaults Avro accepts and rejects Kafka Schema Compatibility Checker What the registry will say, before you ask it Kafka Avro Schema Diff Which direction each change breaks Kafka Compression Comparison Measured on your bytes Kafka Delivery Semantics Exactly-once has a consumer half Kafka ksqlDB Query Builder It looks like SQL and the rules are not Kafka Connect SMT Predicate Tester negate reads backwards Kafka Streams Topology Viewer Count the repartitions Kafka Connect Pipeline Visualizer The order things really run in Kafka Protobuf Binary Decoder Works without the .proto Kafka Protobuf JSON Converter Why your JSON does not round-trip Kafka Protobuf to Avro Schema What does not survive the conversion Kafka Avro Binary Decoder Wrong schema, no error Kafka Avro JSON Converter Why the console producer rejects your line Kafka Avro Sample Data Generator Records that actually serialize Kafka JSON to Avro Schema What JSON cannot tell you Kafka JSON Schema to Avro What does not survive the conversion Kafka SASL JAAS Generator One login module, four syntaxes Kafka CLI Command Builder kcat is librdkafka, not Kafka

Elsewhere on the site