Command Palette

Search for a command to run...

Hectal
Phase 11Advanced12 of 18 in Apache Kafka

Kafka Streams

KStream, KTable and GlobalKTable; stateless and stateful operations; RocksDB state stores, changelogs, standby replicas and recovery; windowing; event time and late data; joins, co-partitioning and repartitioning.

Kafka Streams is a Java library for building stream processors that read from and write to Kafka, with local state that's fault-tolerant because it's backed by Kafka topics. There's no separate cluster to run: your application instances are the processing cluster.

This phase covers the model precisely enough to design real pipelines, and notes when a dedicated engine like Apache Flink is a better fit.

0/5 · 0%
5 topics ~40 min 10 code blocks & diagrams
Start with the first topic
1
11.1

KStream, KTable, GlobalKTable, and Stateless Operations

A KStream is an unbounded sequence of events; a KTable is the latest value per key (a changelog viewed as a table); a GlobalKTable is a table fully replicated to every instance. Stateless operations (filter, map, flatMap, branch, merge) transform records one at a time; stateful ones (aggregate, count, join, windows) keep state.

7 min 1 diagram 1 code practice

2
11.2

State Stores, Changelogs, and Recovery

Stateful operations keep local state in RocksDB (or in-memory) stores, and every change is also written to a compacted changelog topic. When a task moves to another instance, its state is restored from the changelog; standby replicas keep warm copies to make failover fast. Interactive queries expose state stores for reads.

8 min 1 diagram 1 code practice

3
11.3

Windowing: Tumbling, Hopping, Sliding, Session

Windows group events by time for aggregation. Tumbling windows are fixed and non-overlapping, hopping windows are fixed and overlapping, sliding windows are defined by the time difference between records, and session windows are separated by inactivity gaps. Each answers a different business question.

8 min 2 code practice

4
11.4

Event Time, Late Events, and Grace Periods

Stream processors should usually group by event time (when it happened), not processing time (when it arrived). Events arrive late and out of order; a grace period keeps windows open for late arrivals, and events later than the grace are dropped (and counted) or handled by a correction path.

8 min 1 diagram 1 code practice

5
11.5

Joins, Co-Partitioning, and Repartitioning

Streams can join stream–stream (within a time window), stream–table (enrich events with latest state), table–table (combine latest states), and stream–GlobalKTable (enrich without co-partitioning). Non-global joins require co-partitioning: same key, same partition count, same partitioner. Changing keys triggers repartition topics, which cost network, storage and latency.

9 min 1 diagram 1 code practice