Topic 11.1
KStream, KTable, GlobalKTable, and Stateless Operations
In one line
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.
Think of it like this
A bank. The stream of transactions (deposits, withdrawals) is a KStream; the current balance per account is a KTable; the list of branch addresses every teller keeps a full copy of is a GlobalKTable.
Key ideas
- 01
Stream–table duality: a table is the result of applying a stream of updates; a table's changes form a stream.
stream.toTable()andtable.toStream()convert between them. - 02
Stateless:
filter,filterNot,map(changes key, which marks the stream for repartitioning),mapValues(keeps key, no repartition),flatMap,branch/split,merge,peek,selectKey. - 03
Topology: a DAG of processors from sources (input topics) to sinks (output topics). Kafka Streams splits it into tasks, one per input partition (or partition group), and distributes tasks across application instances and threads; scaling means starting more instances.
- 04
Configuration:
application.id(also the consumer group and prefix for internal topics),processing.guarantee(at_least_onceorexactly_once_v2),num.stream.threads,num.standby.replicas. - 05
The Processor API gives low-level control (custom state access, punctuators for timers) when the DSL isn't enough.
Code & diagrams
StreamsBuilder b = new StreamsBuilder();
KStream<String, Order> orders = b.stream("commerce.orders", Consumed.with(Serdes.String(), orderSerde));
Map<String, KStream<String, Order>> branches = orders
.filter((k, o) -> o.total() > 0)
.split(Named.as("orders-"))
.branch((k, o) -> o.country().equals("IN"), Branched.as("india"))
.branch((k, o) -> o.total() > 100_000, Branched.as("high-value"))
.defaultBranch(Branched.as("rest"));
branches.get("orders-india").mapValues(Order::toInvoice).to("invoices.in");
branches.get("orders-high-value").to("fraud.review");
KTable<String, Customer> customers = b.table("customers.state"); // latest per customerId
GlobalKTable<String, Country> countries = b.globalTable("ref.countries");Explain it without notes
Explain the difference between KStream and KTable with the same topic.
Practice
Write a topology that filters paid orders, re-keys them by customerId and writes them to a new topic.
Trade-offs
- ↔
Kafka Streams is a library (no cluster to run) tightly tied to Kafka and the JVM; Flink is a separate cluster with more connectors, SQL and richer event-time features.
Done when you can
I can explain KStream, KTable, GlobalKTable, stream-table duality and stateless operators.