Lesson 20 / 25

Stream Processing: Kafka Streams, Flink and ksqlDB

Choose a stream-processing approach for filtering, joining and aggregating events.

Compute on the stream

Stream processors read topics, transform events and write results to other topics. Kafka Streams is a Java library that runs inside your application, giving stateful operations (aggregations, joins, windows) with local state backed by changelog topics. ksqlDB offers a SQL-like layer. Apache Flink is a separate, powerful engine for large, complex streaming jobs. Core ideas: event time vs processing time, windows (tumbling, hopping, session), state and late data. Start with the simplest tool that fits: a plain consumer for simple transforms, Streams or Flink when you need joins, windows and state.

A windowed count, run

I ran a plain-Python model of a tumbling one-minute window counting events per window start. Real engines add state stores, late-data handling and fault tolerance.

from collections import Counter
events = [5, 20, 59, 61, 75, 130]            # event times in seconds
window = 60
counts = Counter((t // window) * window for t in events)
print(dict(sorted(counts.items())))

Output:

{0: 3, 60: 2, 120: 1}

Plan for late and out-of-order events

Networks and retries mean events arrive late. Decide how long you will wait (a grace period) before closing a window, and what to do with later arrivals.

Quick check: What does a tumbling window do?

  • Drops all events
  • Groups events into fixed, non-overlapping time intervals
  • Overlaps every window by 99%
  • Sorts topics alphabetically
Answer

Groups events into fixed, non-overlapping time intervals — Tumbling windows split time into back-to-back fixed buckets.