14.5 Kafka Streams architecture
5.1 Building a topology
"Topology (also called DAG, or directed acyclic graph) is a set of operations and transitions that EVERY EVENT MOVES THROUGH from input to output. ... Even a simple app has a NONTRIVIAL topology."
The three processor kinds:
- Source processors
- “consume data from a topic and pass it on”
- processors
- “implement an operation of the data — filter, map, aggregate, etc.”
- Sink processors
- “take data from earlier processors and produce it to a topic”
⇒ “A topology Always starts with one or more Source processors and Finishes with one or more Sink processors.”
5.2 💡 Optimizing a topology — three steps, and why step 2 matters
"By DEFAULT, Kafka Streams executes applications built with the DSL API by MAPPING EACH DSL METHOD INDEPENDENTLY to a lower-level equivalent. By evaluating each DSL method independently, OPPORTUNITIES TO OPTIMIZE THE OVERALL RESULTING TOPOLOGY WERE MISSED."
- “The Logical topology is defined by creating
KStreamandKTableobjects and performing DSL operations.” “
StreamsBuilder.build()generates a Physical topology from the logical topology.”⇒ “The second step… is where overall optimizations to the plan can be applied.”
- “
KafkaStreams.start()Executes the topology — this is where data is consumed, processed, and produced.”
Enabling it:
StreamsConfig.TOPOLOGY_OPTIMIZATION = StreamsConfig.OPTIMIZE
builder.build(props) // ⚠ MUST pass the config⚠️ "If you only call
build()WITHOUT PASSING THE CONFIG, OPTIMIZATION IS STILL DISABLED.""Currently, Apache Kafka only contains a few optimizations, mostly around REUSING TOPICS where possible. It is recommended to test applications with AND without optimizations and to COMPARE EXECUTION TIMES AND VOLUMES OF DATA WRITTEN TO KAFKA, and of course, VALIDATE THAT THE RESULTS ARE IDENTICAL in various known scenarios."
5.3 Testing a topology
The primary tool: TopologyTestDriver
"Since its introduction in version 1.1.0, its API has undergone significant improvements, and versions since 2.4 are convenient and easy to use. These tests look like normal unit tests. We define input data, produce it to MOCK INPUT TOPICS, run the topology with the test driver, read the results from MOCK OUTPUT TOPICS, and validate."
⚠️ The gap in
TopologyTestDriver"since it DOES NOT SIMULATE KAFKA STREAMS CACHING BEHAVIOR (an optimization... entirely unrelated to the state store itself, which IS simulated), THERE ARE ENTIRE CLASSES OF ERRORS THAT IT WILL NOT DETECT."
Two integration test frameworks:
| Framework | How | Verdict |
|---|---|---|
| EmbeddedKafkaCluster | "runs Kafka brokers inside the JVM that runs the tests" | |
| Testcontainers | "runs Docker containers with Kafka brokers (and many other components, as needed)" | ✅ Recommended — "since by using Docker it FULLY ISOLATES Kafka, its dependencies, and its RESOURCE USAGE from the application we are trying to test" |
Further reading: the "Testing Kafka Streams — A Deep Dive" blog post.
5.4 💡 Scaling a topology — tasks are the unit of parallelism
"Kafka Streams scales by allowing MULTIPLE THREADS of executions within one instance AND by supporting LOAD BALANCING between DISTRIBUTED INSTANCES. We can run on one machine with multiple threads or on multiple machines; in either case, ALL ACTIVE THREADS WILL BALANCE THE WORK."
“The Streams engine parallelizes execution of a topology by Splitting it into tasks. the number of tasks is determined by the streams engine and depends on the number of partitions in the topics the application processes.”
“Each Task is responsible for A subset of the partitions: the task will Subscribe to those partitions and Consume events from them. For every event it consumes, the task will Execute all the processing steps that apply to this partition In order before eventually writing to the sink.”
⇒ “Those tasks are The basic unit of parallelism in Kafka Streams, because Each task can execute independently of others.”
The scaling recipe:
“we will have as many tasks as we have partitions in the topics we are processing.”
- If we want to process faster → “add more threads”
- If we run out of resources → “start another instance on another server”
“Kafka will automatically coordinate work — it will assign each task its own subset of partitions, and each task will independently process events and maintain its own local state.”
(Partition count is once again the ceiling on parallelism — Ch. 2, Ch. 4, and now Ch. 14.)
Task dependency case 1: joins
*"if we join two streams... we need data from a partition in EACH stream before we can emit a result. Kafka Streams handles this by ASSIGNING ALL THE PARTITIONS NEEDED FOR ONE JOIN TO THE SAME TASK so that the task can consume from all the relevant partitions and perform the join independently.
⚠️ THIS IS WHY KAFKA STREAMS CURRENTLY REQUIRES THAT ALL TOPICS THAT PARTICIPATE IN A JOIN OPERATION HAVE THE SAME NUMBER OF PARTITIONS AND BE PARTITIONED BASED ON THE JOIN KEY."*
That's a hard, checkable precondition — the most common cause of "my Kafka Streams join won't start."
Task dependency case 2: repartitioning (the shuffle)
*"in the ClickStream example, all our events are keyed by user ID. But what if we want to generate statistics per PAGE? Or per ZIP CODE? Kafka Streams will REPARTITION the data by zip code and run an aggregation with the new partitions. If task 1... reaches a processor that repartitions the data (a
groupByoperation), it will need to SHUFFLE, or send events to other tasks.💡 UNLIKE OTHER STREAM PROCESSOR FRAMEWORKS, KAFKA STREAMS REPARTITIONS BY WRITING THE EVENTS TO A NEW TOPIC WITH NEW KEYS AND PARTITIONS. Then ANOTHER SET OF TASKS reads events from the new topic and continues processing."*
💡 Why this design is good: “The second set of tasks depends on the first… however, the first and second sets of tasks can still run independently and in parallel because the first set writes data into a topic at its own rate and the second set consumes from the topic and processes the events on its own.”
“There is no communication and no shared resources between the tasks, and they don't need to run on the same threads or servers. This is one of the more useful things Kafka does — reduce dependencies between different parts of a pipeline.”
The shuffle becomes a Kafka topic. That single decision is what makes Kafka Streams a library rather than a cluster framework: there is no shuffle service, no inter-worker network protocol, no coordinated barrier — just producers and consumers.
5.5 Surviving failures
- “Kafka is Highly available, and therefore the data we persist to Kafka is also highly available. So if the application fails and needs to restart, it can Look up its last position in the stream from kafka and continue from the Last offset it committed.”
- “if the Local state store is lost (e.g., because we needed to Replace the server), the streams application Can always re-create it from the change log it stores in kafka.”
“Kafka Streams also Leverages kafka's consumer coordination to provide high availability for Tasks. If a task failed but there are threads or other instances active, The task will restart on one of the available threads.”
⇒ “Kafka Streams Benefited from improvements in Kafka's consumer group coordination protocol, such as Static group membership and Cooperative rebalancing (Ch. 4), as well as improvements to Kafka's Exactly-once semantics (Ch. 8).”
⚠️ 5.6 The real problem: recovery speed
"While the high-availability methods described here work well in theory, REALITY INTRODUCES SOME COMPLEXITY. ONE IMPORTANT CONCERN IS THE SPEED OF RECOVERY. When a thread has to start processing a task that used to run on a failed thread, it FIRST NEEDS TO RECOVER ITS SAVED STATE — the current aggregation windows, for instance. Often this is done by REREADING INTERNAL TOPICS from Kafka in order to WARM UP Kafka Streams state stores. DURING THE TIME IT TAKES TO RECOVER THE STATE OF A FAILED TASK, THE STREAM PROCESSING JOB WILL NOT MAKE PROGRESS ON THAT SUBSET OF ITS DATA, LEADING TO REDUCED AVAILABILITY AND STALE DATA."
Two techniques, both specific and actionable:
- Aggressive compaction on all Kafka Streams topics. “A key technique is to make sure all Kafka Streams topics are configured for aggressive compaction — by setting a low
min.compaction.lag.msand configuring the segment size to 100 MB instead of the default 1 GB (recall that the last segment in each partition, the active segment, is not compacted).”► Why: recovery time ∝ the size of the changelog you must replay. Smaller segments → more of the log is eligible for compaction → the compacted changelog is smaller → faster warm-up. (Ch. 6 §7.6: the active segment is never compacted. With 1 GB segments, up to 1 GB per partition is uncompacted and must be fully replayed.)
- Standby replicas. “those are tasks that simply shadow active tasks and keep the current state warm on a different server. When failover occurs, they already have the most current state and are ready to continue processing with almost no downtime.”
This is the single most valuable operational paragraph in the chapter. Kafka Streams HA is not about whether it recovers — it always does. It's about how long the affected keys are stale. Two knobs:
segment.bytes= 100 MB + lowmin.compaction.lag.ms— smaller replaynum.standby.replicas> 0 — no replay