Part 8 · 1 chapters · ~8 min

Kafka Streams and Flink Basics

Stateless and stateful processing, local state stores and changelog topics, KStream and KTable, windowing (tumbling, hopping, sliding, session), event time, watermarks and late data, stream joins, exactly-once processing, and choosing between Kafka Streams and Flink.

10

Processing streams

code
// Kafka Streams: count transfers per account in 10-minute windows, flag bursts
StreamsBuilder b = new StreamsBuilder();
b.stream("transfers", Consumed.with(Serdes.String(), transferSerde))
 .groupByKey()
 .windowedBy(SlidingWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(10)))
 .count(Materialized.as("transfer-counts"))
 .toStream()
 .filter((windowedAccount, count) -> count > 5)
 .map((w, count) -> KeyValue.pair(w.key(), new BurstAlert(w.key(), count, w.window().end())))
 .to("fraud-alerts");
// processing.guarantee = exactly_once_v2
STREAM PROCESSING CONCEPTS
Kafka Streams and Flink share the same ideas
statelessfilter, map, branch: each recordon its own.statefulcounts, aggregates, joins: statekept per key in a local storebacked by a changelog topic.windowstumbling (fixed), hopping(overlapping), sliding, session(gaps close them).event timeUse when the event happened, notwhen it arrived; watermarks decidewhen a window is complete.joinsstream-stream (within a window),stream-table (enrich with thelatest value).Streams vs FlinkKafka Streams: a library in yourapp. Flink: a separate clusterwith richer time handling.
swipe the figure sideways, or tap expand for full screen
1/6
stateless
Filtering and mapping need no memory of earlier records: easy to scale and restart.
per-record transformstrivially scalable