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_v2STREAM PROCESSING CONCEPTS
Kafka Streams and Flink share the same ideas
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