Kafka Streams Topology Viewer

Paste the output of topology.describe() and see the shape of it. The number that matters is the sub-topology count: each boundary is a full round trip through the broker, and the description states it only as a heading.

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 Draw it. 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.

A filter topology

Sub-topologies, sources and sinks read out of describe() output

Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [input])
      --> KSTREAM-FILTER-0000000001
    Processor: KSTREAM-FILTER-0000000001 (stores: [])
      --> KSTREAM-SINK-0000000002
      <-- KSTREAM-SOURCE-0000000000
    Sink: KSTREAM-SINK-0000000002 (topic: output)
      <-- KSTREAM-FILTER-0000000001

A stateful topology

An aggregation, its state store, and the changelog topic it creates on the cluster

Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [orders])
      --> KSTREAM-AGGREGATE-0000000001
    Processor: KSTREAM-AGGREGATE-0000000001 (stores: [order-totals])
      <-- KSTREAM-SOURCE-0000000000

Common mistakes

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

  1. Assuming sub-topologies share a thread

    Each sub-topology is scheduled independently, and tasks are assigned per sub-topology. A repartition splits one topology into two that no longer share state or ordering.

    Instead:Read the sub-topology boundaries as the real parallelism unit.

  2. Ignoring the repartition topics

    Any key-changing operation followed by an aggregation creates an internal repartition topic, which is real network traffic, real disk and real retention on the cluster.

    Instead:Check what the topology creates before deploying it. The topic predictor on this site lists them.

  3. Treating a state store as ephemeral

    Stores are backed by changelog topics that are compacted and kept forever by default. Deleting the application without cleanup leaves them on the cluster.

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

The sub-topology count is the most expensive number in your application

Everything else in a topology description is detail. The number of sub-topologies is how many times each record is written back to Kafka and read again, and it is stated only as a heading number nobody reads as a cost.

Every sub-topology boundary is a full round trip through the broker

Streams splits a topology wherever the key changes, because the data has to be redistributed before the next stage can group or join on it. That split is a repartition: the record is serialised, produced to an internal topic, replicated, fetched back and deserialised. Two sub-topologies means one round trip per record; four means three. It is almost always the largest cost in a Streams application and it is invisible in the description unless you already know to count the headings.

The usual cause is map where mapValues would do

selectKey, map, groupBy and a join on a different key all tell Streams the key may have changed, and Streams believes them. mapValues and flatMapValues promise the key is untouched, so no repartition is inserted. Swapping one for the other is often a one-line change that removes an entire round trip, and it is the first thing to look for when this page shows more sub-topologies than you expected.

// forces a repartition
stream.map((k, v) -> KeyValue.pair(k, f(v)))

// does not
stream.mapValues(v -> f(v))

Task count is the ceiling on useful parallelism

Streams creates one task per sub-topology per input partition, and a task is the unit of assignment. Multiply the sub-topology count by the partition count and that is the maximum number of threads that can ever do work. Instances beyond it hold no tasks and sit idle, which looks like a scaling problem and is a partitioning one.

A state store is a changelog topic, disk, and restore time

Every store is backed by a compacted changelog unless it was explicitly built with logging disabled. That is broker disk proportional to the state, and it is what has to be replayed to rebuild the store when a task moves to another instance. Long stalls after a rebalance are usually restore, and standby replicas trade memory for cutting that time rather than cutting the disk.

What this cannot see

Partition counts, so it can state the task-count formula and not the number. Throughput, so it can say a boundary is expensive and not how expensive. Whether a store has logging disabled, because the description does not record it. It also cannot see what your processors do: a custom Processor that writes to an external system is a black box here and shows as a node with no successors.

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 Streams Internal Topic Predictor Create them before Streams does 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 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