14.2 Stream processing concepts
Examples used in the chapter: filter, count, group-by, left-join.
2.1 Topology
"A processing topology starts with one or more SOURCE STREAMS that are passed through a graph of STREAM PROCESSORS connected through event streams, until results are written to one or more SINK STREAMS. Each stream processor is a computational step applied to the stream of events in order to transform the events."
Examples used in the chapter: filter, count, group-by, left-join.
2.2 ⚠️ Time — "probably the most important concept in stream processing and often the most confusing"
"having a common notion of time is critical because most stream applications perform operations on TIME WINDOWS. For example, our stream application might calculate a moving five-minute average of stock prices. In that case, we need to know what to do when one of our producers goes offline for two hours due to network issues and RETURNS WITH TWO HOURS' WORTH OF DATA — most of the data will be relevant for five-minute time windows THAT HAVE LONG PASSED AND FOR WHICH THE RESULT WAS ALREADY CALCULATED AND STORED."
(Recommended reading: Justin Sheehy's paper "There Is No Now.")
The three notions of time:
| Event time — “usually the time that matters most” | Log append time (ingestion time) | Processing time — “highly unreliable and best avoided” | |
|---|---|---|---|
| What it is | “the time the events we are tracking occurred and the record was created — the time a measurement was taken, an item was sold, a user viewed a page.” | “the time the event arrived at the Kafka broker and was stored.” | “the time at which a stream processing application received the event in order to perform some calculation. This time can be milliseconds, hours, or days after the event occurred.” |
| Who sets it | “In version 0.10.0 and later, Kafka automatically adds the current time to producer records at the time they are created.” | “brokers will automatically add this time if Kafka is configured to do so or if the records arrive from older producers and contain no timestamps.” | Whichever application happens to read the event. |
| The catch | ⚠ “If this does not match the application's notion of event time — such as in cases where the Kafka record is created based on a database record sometime after the event occurred — then we recommend adding the event time as a field in the record itself so that both timestamps will be available for later processing.” | “Typically less relevant… if we calculate the number of devices produced per day, we want to count devices that were actually produced on that day, even if network issues meant the event only arrived the following day.” | ⚠ “assigns different timestamps to the same event depending on exactly when each stream processing application happened to read the event. It can even differ for two threads in the same application!” |
| When it's fine | Almost always the one you want. | 💡 “in cases where the real event time was not recorded, log append time can still be used consistently because it does not change after the record was created, and assuming no delays in the pipeline, it can be a reasonable approximation of event time.” | Never, if you can help it. |
Kafka Streams' mechanism: "assigns time to each event based on the TimestampExtractor interface. Developers can use different implementations, which can use either of the three time semantics or A COMPLETELY DIFFERENT CHOICE OF TIMESTAMP, INCLUDING EXTRACTING A TIMESTAMP FROM THE CONTENTS OF THE EVENT ITSELF."
💡 The four output-timestamp rules
"When Kafka Streams writes output to a Kafka topic, it assigns a timestamp to each event based on the following rules:"
- Output maps directly to an input record → “the output record will use the same timestamp as the input”.
- Output is the result of an aggregation → “the timestamp will be the maximum timestamp used in the aggregation”.
- Output is the result of joining two streams → “the timestamp is the largest of the two records being joined.” “When a stream and a table are joined, the timestamp from the stream record is used.”
- Output generated on a schedule regardless of input (e.g.
punctuate()) → “the output timestamp will depend on the current internal times of the stream processing app”.
"When using the lower-level Processor API rather than the DSL, Kafka Streams includes APIs for MANIPULATING THE TIMESTAMPS OF RECORDS DIRECTLY, so developers can implement timestamp semantics that match the required business logic."
⚠️ MIND THE TIME ZONE
"THE ENTIRE DATA PIPELINE SHOULD STANDARDIZE ON A SINGLE TIME ZONE; OTHERWISE, RESULTS OF STREAM OPERATIONS WILL BE CONFUSING AND OFTEN MEANINGLESS. If you must handle data streams with different time zones, you need to make sure you can convert events to a single time zone BEFORE performing operations on time windows. Often this means STORING THE TIME ZONE IN THE RECORD ITSELF."
2.3 State
When you don't need it: "if all we need to do is read a stream of online shopping transactions, find the transactions over $10,000, and email the relevant salesperson, we can probably write this in just a few lines of code using a Kafka consumer and SMTP library."
When you do: "Stream processing becomes really interesting when we have operations that involve MULTIPLE EVENTS: counting the number of events by type, moving averages, joining two streams... we need to keep track of more information... We call this information a STATE."
⚠️ The local-variable trap
"It is often tempting to store the state in variables that are LOCAL to the stream processing app, such as a simple hash table. In fact, WE DID JUST THAT IN MANY EXAMPLES IN THIS BOOK. HOWEVER, THIS IS NOT A RELIABLE APPROACH because WHEN THE APPLICATION IS STOPPED OR CRASHES, THE STATE IS LOST, WHICH CHANGES THE RESULTS."
(That's the honest self-critique of Ch. 4's moving-average and word-count examples — and Ch. 7 §5.2's "consumers may need to maintain state.")
Two kinds of state:
| Local / internal | External | |
|---|---|---|
| Where | "an embedded, in-memory database running WITHIN the application" | "an external data store, often a NoSQL system like Cassandra" |
| Visibility | "accessible ONLY by a specific INSTANCE" | "can be accessed from multiple instances or even different applications" |
| ✅ | "extremely fast" | "virtually unlimited size" |
| ❌ | "we are limited to the amount of memory available" | "the extra LATENCY and COMPLEXITY introduced with an additional system, as well as AVAILABILITY — the application needs to handle the possibility that the external system is not available" |
💡 "As a result, MANY OF THE DESIGN PATTERNS IN STREAM PROCESSING FOCUS ON WAYS TO PARTITION THE DATA INTO SUBSTREAMS THAT CAN BE PROCESSED USING A LIMITED AMOUNT OF LOCAL STATE."
"Most stream processing apps try to AVOID having to deal with an external store, or at least limit the latency overhead by caching information in the local state and communicating with the external store AS RARELY AS POSSIBLE. THIS USUALLY INTRODUCES CHALLENGES WITH MAINTAINING CONSISTENCY BETWEEN THE INTERNAL AND EXTERNAL STATE."
💡 2.4 Stream-table duality — the central idea
*"A TABLE is a collection of records, each identified by its primary key... Table records are MUTABLE. Querying a table allows checking the state of the data AT A SPECIFIC POINT IN TIME. ... Unless the table was specifically designed to include history, WE WILL NOT FIND THEIR PAST CONTACTS in the table.
Unlike tables, STREAMS CONTAIN A HISTORY OF CHANGES. A stream is a string of events wherein each event CAUSED a change. A table contains a CURRENT STATE of the world, which is THE RESULT of many changes.
From this description, it is clear that STREAMS AND TABLES ARE TWO SIDES OF THE SAME COIN — the world always changes, and SOMETIMES WE ARE INTERESTED IN THE EVENTS THAT CAUSED THOSE CHANGES, whereas OTHER TIMES WE ARE INTERESTED IN THE CURRENT STATE of the world. SYSTEMS THAT ALLOW US TO TRANSITION BACK AND FORTH BETWEEN THE TWO WAYS OF LOOKING AT DATA ARE MORE POWERFUL THAN SYSTEMS THAT SUPPORT JUST ONE."*
Table → stream: “capture the changes that modify the table. Take all those INSERT, UPDATE, and DELETE events and store them in a stream. Most databases offer change data capture (CDC) solutions… and there are many Kafka connectors that can pipe those changes into Kafka.”
Stream → table is “apply all the changes” = materializing: “We create a table, either in memory, in an internal state store, or in an external database, and start going over all the events in the stream from beginning to end, changing the state as we go.”
The shoe-store example:
- The table answers: “what our inventory contains right now”, “how much money we made”.
- The stream answers: “how busy the store is” (four customer events today), “why the blue shoes were returned”.
- ► Neither view can answer the other's questions. You need both.
2.5 Time windows — the three dimensions nobody thinks about
"Most operations on streams are WINDOWED operations, operating on slices of time: moving averages, top products sold this week, 99th percentile load... Join operations on two streams are ALSO windowed. VERY FEW PEOPLE STOP AND THINK ABOUT THE TYPE OF WINDOW THEY WANT."
| Dimension | The question | What to know |
|---|---|---|
| ① Size of the window | “Every five-minute window? Every 15-minute? Or the entire day?” | ⚠ Trade-off: “Larger windows are smoother but they lag more — if the price increases, it will take longer to notice than with a smaller window.” |
| ② Advance interval — how often the window moves | “Five-minute averages can update every minute, every second, or every time there is a new event.” | “Windows for which the size is a fixed time interval are called hopping windows. When the advance interval is equal to the window size, it is called a tumbling window.” |
| ③ Grace period — how long the window remains updatable | “Our five-minute moving average calculated the average for the 00:00–00:05 window. Now, an hour later, we are getting a few more input records with their event time showing 00:02. Do we update the result? Or do we let bygones be bygones?” | “Ideally, we'll be able to define a certain time period during which events will get added to their respective time slice. For example, if the events were delayed up to four hours, we should recalculate and update. If events arrive later than that, we can ignore them.” |
💡 A session window is the special case where “the size of the window is defined by a period of inactivity. The developer defines a session gap, and all events that arrive continuously with gaps smaller than the defined session gap belong to the same session. A gap in arrivals will define a new session.”
Aligned vs unaligned windows:
Those hopping windows are aligned to clock time — a five-minute window advancing every minute, starting on the minute. Unaligned windows “simply start whenever the app started”: 03:17–03:22 │ 03:18–03:23 │ …
2.6 Processing guarantees
"A KEY REQUIREMENT for stream processing applications is the ability to process each record EXACTLY ONCE, regardless of failures. WITHOUT EXACTLY-ONCE GUARANTEES, STREAM PROCESSING CAN'T BE USED IN CASES WHERE ACCURATE RESULTS ARE NEEDED."
processing.guarantee = exactly_once
processing.guarantee = exactly_once_betaexactly_once_beta is “a more efficient exactly-once implementation” — it needs Kafka Streams 2.6+ and brokers 2.5+.
(Full mechanism: Ch. 8. exactly_once_beta is what allows one transactional producer to handle many partitions — Ch. 8 §4.1.)