Kafka Streams: Embedded Stream Processing on Kafka
Kafka Streams is a client library for building stream processing applications that read from and write to Kafka. It is not a separate cluster or a processing engine you deploy remotely; it is a Java library that runs inside your application process. You add it as a dependency, write a topology (a directed acyclic graph of operators like map, filter, join, aggregate), and run it. The library handles partitioning, task assignment, state management, fault tolerance, and rebalancing. The key insight is that Kafka Streams is embedded: there is no scheduler, no resource manager, and no separate cluster to operate. Scaling is done by running more instances of your application, and Kafka Streams distributes the work using the consumer group protocol. This is fundamentally different from Flink or Spark Streaming, which require a cluster and a job manager.
The mechanism that makes this work is the stream thread and task model. Each Kafka Streams instance runs one or more stream threads. Each stream thread processes one or more tasks, where a task corresponds to a partition of the input topics. The number of tasks is determined by the number of input partitions; you cannot have more tasks than partitions, which is why partition count is a key scalability limit. Kafka Streams uses the consumer group protocol to assign partitions to instances, so when an instance joins or leaves, a rebalance occurs and tasks are reassigned. State stores are local to the instance and are backed by changelog topics in Kafka, so if a task moves to another instance, its state can be restored from the changelog. This is what gives Kafka Streams fault tolerance without a separate storage cluster.
When would you use it? Kafka Streams is a good fit when your data is already in Kafka, your processing logic is expressible as a topology of transformations, and you want to avoid operating a separate stream processing cluster. It is particularly strong for stateful operations like joins, aggregations, and windowing, because the state store and changelog abstraction hides a lot of complexity. It is not a good fit when you need to call external systems synchronously in the hot path, when you need complex event-time processing with advanced watermarks beyond what the DSL offers, or when your team is not a JVM shop. The trade-off compared to Flink is that Flink has a richer event-time model, better support for batch and stream unification, and a separate cluster that can be scaled independently. Kafka Streams is simpler to operate but its scaling is tied to Kafka partitions and it runs in your application's JVM, which means your processing shares resources with your application. Version note: the exactly-once processing guarantee is controlled by processing.guarantee; use exactly_once_v2 (KIP-447, Kafka 2.5+) for better scalability. Kafka Streams also has a Processor API for lower-level control when the DSL is not enough.
Kafka Streams is an embedded Java library, not a separate cluster.
You write a topology of operators and run it in your application process.
Scaling is via more instances using the consumer group protocol; tasks equal input partitions.
State stores are local and backed by changelog topics for fault tolerance.
Good fit: data already in Kafka, stateful joins/aggregations, want to avoid a separate cluster.
Not a good fit: synchronous external calls in the hot path, non-JVM teams, or needing Flink-level event-time features.
Use exactly_once_v2 for scalable EOS; Processor API for low-level control.
0-2 years experience
2-5 years experience
5-8 years experience
8+ years experience