Kafka Streams¶
π Kafka Streams Documentation - Official Kafka Streams documentation π Confluent Kafka Streams Guide - Confluent developer guide
Kafka Streams Architecture¶
Overview¶
- Client library for building stream processing applications
- No separate cluster required - runs in your application process
- Leverages Kafka for fault tolerance, scalability, and exactly-once processing
- Uses consumer groups for parallel processing and partition assignment
- Input and output are Kafka topics
Topology¶
- A topology is a directed acyclic graph (DAG) of stream processors
- Source processors - Read from Kafka topics
- Stream processors - Transform data (filter, map, join, aggregate)
- Sink processors - Write to Kafka topics
- Topologies are built using the StreamsBuilder DSL or Processor API
π Topology - Understanding processor topologies
Tasks and Threads¶
- A topology is divided into tasks based on input partition count
- Each task processes data from one or more partitions
num.stream.threadscontrols the number of processing threads per instance- Tasks are distributed across threads and application instances
- Maximum parallelism = number of input partitions
State Stores¶
- Used by stateful operations (aggregation, joins, windowing)
- Backed by RocksDB by default (can use in-memory stores)
- Each state store has a changelog topic for fault tolerance
- Changelog topics enable state recovery on failure
- Interactive queries allow external access to state stores
π State Stores - State store types and configuration
KStream and KTable¶
KStream (Record Stream)¶
- Represents an unbounded stream of records
- Each record is an independent, immutable event
- Insert semantics - every record is a new event
- Analogous to a database INSERT log
- Created from a topic:
builder.stream("topic")
KTable (Changelog Stream)¶
- Represents a table of latest values per key
- Upsert semantics - new records update or insert by key
- Records with null values are treated as deletes (tombstones)
- Analogous to a database table with primary key
- Created from a topic:
builder.table("topic") - Backed by a state store
Stream-Table Duality¶
- Every stream can be viewed as a table (aggregating all changes)
- Every table can be viewed as a stream (changelog of updates)
KStream.toTable()converts stream to tableKTable.toStream()converts table to stream
GlobalKTable¶
- Fully replicated table - all partitions on every instance
- Used for enrichment joins without repartitioning
- Created from a topic:
builder.globalTable("topic") - Supports join with KStream on any key (not just the partition key)
- Higher memory usage since all data is on each instance
- Not partitioned - reads all partitions regardless of assignment
Stateless Transformations¶
π Stateless Transformations - DSL reference for stateless operations
filter / filterNot¶
KStream<String, Long> filtered = stream.filter(
(key, value) -> value > 100
);
map / mapValues¶
// map - can change key and value (triggers repartition if key changes)
KStream<String, String> mapped = stream.map(
(key, value) -> KeyValue.pair(key.toUpperCase(), value.toString())
);
// mapValues - only changes value (no repartition)
KStream<String, String> mappedValues = stream.mapValues(
value -> value.toString()
);
map() that changes the key triggers repartitioning - Prefer mapValues() when only transforming values to avoid repartition flatMap / flatMapValues¶
KStream<String, String> flatMapped = stream.flatMap(
(key, value) -> {
List<KeyValue<String, String>> result = new ArrayList<>();
for (String word : value.split(" ")) {
result.add(KeyValue.pair(word, word));
}
return result;
}
);
flatMap() can change key (triggers repartition) - flatMapValues() only changes values (no repartition) branch (split)¶
// Kafka Streams 2.8+ uses split()
Map<String, KStream<String, Long>> branches = stream.split(Named.as("split-"))
.branch((key, value) -> value > 100, Branched.as("high"))
.branch((key, value) -> value > 10, Branched.as("medium"))
.defaultBranch(Branched.as("low"));
KStream<String, Long> highStream = branches.get("split-high");
selectKey¶
KStream<String, Order> rekeyed = stream.selectKey(
(key, value) -> value.getCustomerId()
);
merge¶
KStream<String, String> merged = stream1.merge(stream2);
Stateful Transformations¶
π Stateful Transformations - DSL reference for stateful operations
groupByKey / groupBy¶
// groupByKey - uses existing key (no repartition)
KGroupedStream<String, Long> grouped = stream.groupByKey();
// groupBy - new key (triggers repartition)
KGroupedStream<String, Long> regrouped = stream.groupBy(
(key, value) -> value.getCategory()
);
groupByKey() avoids repartitioning - groupBy() triggers repartitioning count¶
KTable<String, Long> counts = grouped.count(
Materialized.as("count-store")
);
aggregate¶
KTable<String, Double> aggregated = grouped.aggregate(
() -> 0.0, // Initializer
(key, value, aggregate) -> aggregate + value, // Aggregator
Materialized.as("aggregate-store")
);
reduce¶
KTable<String, Long> reduced = grouped.reduce(
(value1, value2) -> value1 + value2,
Materialized.as("reduce-store")
);
Windowing¶
π Windowing - Window types and configuration
Tumbling Windows¶
KTable<Windowed<String>, Long> counts = grouped
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
.count();
Hopping Windows¶
KTable<Windowed<String>, Long> counts = grouped
.windowedBy(TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(5),
Duration.ofMinutes(1))
.advanceBy(Duration.ofMinutes(1)))
.count();
Sliding Windows (Join Windows)¶
KStream<String, String> joined = stream1.join(
stream2,
(value1, value2) -> value1 + "-" + value2,
JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(5))
);
Session Windows¶
KTable<Windowed<String>, Long> counts = grouped
.windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(5)))
.count();
Grace Period¶
- Time after window closes to accept late-arriving records
- Records arriving after the grace period are dropped
TimeWindows.ofSizeAndGrace(size, grace)- Default grace period: 24 hours (in older versions)
- Use
WithNoGracevariants for zero grace period
Joins¶
π Joins - Join semantics and types
Join Types Summary¶
| Join | Left | Right | Windowed | Output |
|---|---|---|---|---|
| KStream-KStream (inner) | KStream | KStream | Yes (required) | KStream |
| KStream-KStream (left) | KStream | KStream | Yes (required) | KStream |
| KStream-KStream (outer) | KStream | KStream | Yes (required) | KStream |
| KStream-KTable (inner) | KStream | KTable | No | KStream |
| KStream-KTable (left) | KStream | KTable | No | KStream |
| KTable-KTable (inner) | KTable | KTable | No | KTable |
| KTable-KTable (left) | KTable | KTable | No | KTable |
| KTable-KTable (outer) | KTable | KTable | No | KTable |
| KStream-GlobalKTable | KStream | GlobalKTable | No | KStream |
Key Requirements¶
- KStream-KStream: Both sides must have the same key
- KStream-KTable: Both sides must have the same key (co-partitioned)
- KTable-KTable: Both sides must have the same key (co-partitioned)
- KStream-GlobalKTable: Custom key mapper allows joining on any field
- Co-partitioned topics must have the same number of partitions
Co-Partitioning Requirements¶
- Same number of partitions on both sides
- Same partitioning strategy
- If not co-partitioned, use
selectKey()+repartition()or use GlobalKTable
Interactive Queries¶
π Interactive Queries - Querying state stores
ReadOnlyKeyValueStore<String, Long> store =
streams.store(StoreQueryParameters.fromNameAndType(
"count-store", QueryableStoreTypes.keyValueStore()));
Long count = store.get("my-key");
- Access state stores from outside the stream processing topology
- Read-only access - cannot modify state
- Local queries only return data from partitions assigned to the instance
- For global queries, use RPC between instances or GlobalKTable
Configuration¶
Key Kafka Streams Settings: - application.id - Unique ID for the streams application (used as consumer group.id) - bootstrap.servers - Kafka cluster connection - num.stream.threads (default: 1) - Processing threads per instance - state.dir (default: /tmp/kafka-streams) - Directory for state stores - processing.guarantee - at_least_once (default) or exactly_once_v2 - default.key.serde / default.value.serde - Default serialization - cache.max.bytes.buffering (default: 10485760 / 10 MB) - Record cache size - commit.interval.ms (default: 30000 with at_least_once, 100 with exactly_once) - Commit frequency
π Kafka Streams Configuration - All configurable settings