14.3 Stream processing design patterns
Stream joins
- 240 MB
- local state
No window needed. The table is materialised locally from a compacted changelog topic, so each event is enriched by a local lookup — no network call per record and no remote database to overload. The cost is holding that state on every instance and rebuilding it on failover.
3.1 Single-event processing (map/filter)
"The most basic pattern... also known as a map/filter pattern because it is commonly used to filter unnecessary events from the stream or transform each event. (The term map is based on the map/reduce pattern in which the map stage transforms events and the reduce stage aggregates them.)"
- “reads log messages and writes ERROR events into a high-priority stream and the rest into a low-priority stream”
- “reads events and modifies them from JSON to Avro”
💡 Why it's easy: “Such applications do not need to maintain state because each event can be handled independently. This means that recovering from app failures or load-balancing is incredibly easy as there is no need to recover state; we can simply hand off the events to another instance.” ► It “can be easily handled with a simple producer and consumer”.
3.2 Processing with local state
"Most stream processing applications are concerned with AGGREGATING information, especially WINDOW aggregation." Example: min/max stock prices per day, moving average.
Why local state suffices — the partitioning argument:
“All this can be done using local state (rather than a shared state) because each operation in our example is a group-by aggregate. That is, we perform the aggregation per stock symbol, not on the entire stock market in general.”
② is “a Kafka consumer guarantee”: “each instance of the application will get all the events from the partitions that are assigned to it.”
Key-based partitioning (Ch. 3 §9) + consumer group ownership (Ch. 4 §1) = each instance's local state is complete for its own keys. This is why local state works at all. It is the whole trick.
Three problems local state creates:
- Memory usage. The state “ideally fits into the memory available. Some local stores allow spilling to disk, but this has significant performance impact.”
- Persistence. “we need to make sure the state is not lost when an instance shuts down”. 💡 How Kafka Streams does it:
- “local state is stored in-memory using embedded RocksDB, which also persists the data to disk for quick recovery after restarts”
- “but all the changes to the local state are also sent to a Kafka topic. If a stream's node goes down, the local state is not lost — it can be easily re-created by rereading the events from the Kafka topic.” For example, “current minimum for IBM = 167.19” is stored in Kafka.
- “Kafka uses log compaction for these topics to make sure they don't grow endlessly and that re-creating the state is always feasible.”
- Rebalancing. “Partitions sometimes get reassigned. When this happens, the instance that loses the partition must store the last good state, and the instance that receives the partition must know to recover the correct state.”
(Log-compacted changelog topics are the Ch. 1 §3.6 / Ch. 6 §7 compaction feature doing exactly the job it was designed for: "changelog-type data, where only the last update is interesting.")
3.3 Multiphase processing / repartitioning
"Local state is great if we need a group-by type of aggregate. But what if we need a result that uses ALL AVAILABLE INFORMATION?"
The example: "suppose we want to publish the top 10 stocks each day — the 10 stocks that gained the most from opening to closing. Obviously, NOTHING WE DO LOCALLY ON EACH APPLICATION INSTANCE IS ENOUGH because ALL THE TOP 10 STOCKS COULD BE IN PARTITIONS ASSIGNED TO OTHER INSTANCES."
The summary topic is “obviously much smaller with significantly less traffic than the topics that contain the trades themselves, and therefore it can be processed by a single instance.”
💡 "This type of multiphase processing is very familiar to those who write MapReduce code, where you often have to resort to multiple reduce phases. If you've ever written map-reduce code, you'll remember that YOU NEEDED A SEPARATE APP FOR EACH REDUCE STEP. UNLIKE MAPREDUCE, MOST STREAM PROCESSING FRAMEWORKS ALLOW INCLUDING ALL STEPS IN A SINGLE APP, with the framework handling the details of which application instance will run each step."
3.4 Stream-table join (processing with external lookup)
The obvious approach, and why it fails:
The obvious idea: “for every click event, Look up the user in the profile database and write an event that includes the original click Plus the user age and gender.”
⚠ Three problems:
- Latency: “an external lookup adds Significant latency to the processing of Every record — Usually between 5 and 15 milliseconds. In many cases, this is not feasible.”
- Throughput mismatch: “stream processing systems can often handle 100K–500k events per second, but the database can only handle perhaps 10k events per second at reasonable performance.”
- Availability: “our application will need to handle situations when the external DB Is not available.”
The caching dilemma:
"To get good performance and availability, we need to CACHE the information... Managing this cache can be challenging though — HOW DO WE PREVENT THE INFORMATION IN THE CACHE FROM GETTING STALE? If we refresh events TOO OFTEN, WE ARE STILL HAMMERING THE DATABASE, and the cache isn't helping much. If we wait TOO LONG to get new events, WE ARE DOING STREAM PROCESSING WITH STALE INFORMATION."
The resolution — CDC:
“But if we can capture all the changes that happen to the database table in a stream of events, we can have our stream processing job listen to this stream and update the cache based on database change events.”
💡 “because we are using a local state, this scales a lot better and will not affect the database and other apps using it.” ► “We refer to this as a stream-table join because one of the streams represents changes to a locally cached table.”
(Ch. 9's Debezium recommendation is exactly the tool for the CDC leg.)
3.5 Table-table join
*"There is no reason why we can't have those materialized tables in BOTH SIDES of the join operation.
Joining two tables is ALWAYS NONWINDOWED and joins THE CURRENT STATE of both tables at the time the operation is performed."*
| Equi-join | Foreign-key join |
|---|---|
| “both tables have the same key that is partitioned in the same way, and therefore the join operation can be efficiently distributed between a large number of application instances and machines.” | “the key of one stream or table is joined with an arbitrary field from another stream or table.” |
On foreign-key joins, see “Crossing the Streams” (Kafka Summit 2020).
3.6 Streaming join (windowed join)
"What makes a stream 'real'? ... When we use a stream to represent a TABLE, we can IGNORE MOST OF THE HISTORY because we only care about the current state. BUT WHEN WE JOIN TWO STREAMS, WE ARE JOINING THE ENTIRE HISTORY, trying to match events in one stream with events in the other that have THE SAME KEY and HAPPENED IN THE SAME TIME WINDOWS. THIS IS WHY A STREAMING JOIN IS ALSO CALLED A WINDOWED JOIN."
The example: search queries ⋈ clicks on search results.
"We want to match search queries with the results they clicked on so that we will know which result is most popular for which query. Obviously, we want to match results based on the search term BUT ONLY MATCH THEM WITHIN A CERTAIN TIME WINDOW. We assume the result is clicked SECONDS AFTER the query was entered. So we keep a small, few-seconds-long window on each stream and match the results from each window."
How Kafka Streams makes it work:
“Kafka Streams supports equi-joins, where streams, queries, and clicks are partitioned on the same keys, which are also the join keys.”
“Kafka Streams then makes sure that partition 5 of both topics is assigned to the same task. So this task sees all the relevant events for user_id:42. It maintains the join window for both topics in its embedded RocksDB state store, and this is how it can perform the join.”
3.7 Out-of-sequence events
"a challenge not just in stream processing but ALSO in traditional ETL systems. Out-of-sequence events happen QUITE FREQUENTLY AND EXPECTEDLY in IoT scenarios."
Three real causes:
- “a mobile device loses WiFi signal for a few hours and sends a few hours' worth of events when it reconnects”
- “monitoring network equipment (a faulty switch doesn't send diagnostics signals until it is repaired)”
- “manufacturing (network connectivity in plants is notoriously unreliable, especially in developing countries)”
The four things an application must do:
- Recognize that an event is out of sequence — “requires that the application Examines the event time and discovers that it is Older than the current time”
- Define a reconciliation period — “Perhaps a Three-hour delay should be reconciled, and events over Three weeks old can be Thrown away.”
Have an In-band capability to reconcile the event
💡 “This is the main difference between streaming apps and batch jobs. If we have a daily batch job and a few events arrived after the job completed, We can usually just rerun yesterday's job. with stream processing, there is no ‘rerun yesterday's job’ — the same continuous process needs to handle both old and new events at any given moment.”
- Be able to Update results — “If the results are written into a Database, a PUT or UPDATE is enough. If the stream app sends results by email, updates may be trickier.”
How frameworks support it:
"Several frameworks, including Google's Dataflow and Kafka Streams, have built-in support for the notion of event time INDEPENDENT of the processing time... typically done by maintaining MULTIPLE AGGREGATION WINDOWS available for update in the local state and giving developers the ability to configure how long to keep those window aggregates available. ⚠ Of course, THE LONGER THE AGGREGATION WINDOWS ARE KEPT AVAILABLE FOR UPDATES, THE MORE MEMORY IS REQUIRED to maintain the local state."
And the compaction trick that makes updates work:
"The Kafka Streams API always writes aggregation results to result topics. Those are usually COMPACTED TOPICS, which means that ONLY THE LATEST VALUE FOR EACH KEY IS PRESERVED. In case the results of an aggregation window need to be updated as a result of a late event, Kafka Streams will simply WRITE A NEW RESULT for this aggregation window, WHICH WILL EFFECTIVELY REPLACE THE PREVIOUS RESULT."
Late-event correction is implemented by log compaction. Compaction retains only the latest value for each key, so no mutation is needed anywhere.
💡 3.8 Reprocessing — two variants, one strong recommendation
| Variant 1 — a new version of the app (A/B, comparison) — easy | Variant 2 — the existing app is buggy; reprocess in place — hard |
|---|---|
“made simple by the fact that Apache Kafka stores the event streams in their entirety for long periods of time in a scalable data store.” Requires only:
| “More challenging — it requires:
|
► The recommendation: “While Kafka Streams has a tool for resetting state, our recommendation is to try to use the first method whenever sufficient capacity exists to run two copies and generate two result streams. The first method is much safer — it allows switching back and forth between multiple versions and comparing results, and doesn't risk losing critical data or introducing errors during the cleanup process.”
(This is Ch. 5 §6.4's shoe-counting warning, resolved architecturally: don't reset offsets and state — run a second app.)
3.9 Interactive queries
"Most of the time the users of stream processing applications get the results by reading them from an output topic. In some cases, however, IT IS DESIRABLE TO TAKE A SHORTCUT AND READ THE RESULTS FROM THE STATE STORE ITSELF. This is common when the result is a TABLE (e.g., the top 10 best-selling books) and the stream of results is really a STREAM OF UPDATES to this table — it is MUCH FASTER AND EASIER to just read the table directly from the stream processing application state."