Learn Labs
12. Stream Processing

12.4 Processing Streams

Examples: rate of a certain event type · rolling average over a time period · comparing current statistics to previous intervals (to detect trends or alert on metrics unusually hi…

Windowing

window type
5
windows
events
e1e2e3
window
w1w2w3w4w5
Safe

Tumbling windows are fixed and non-overlapping: every event belongs to exactly one. The simplest to reason about and to recover after a failure.

A stream has no end, so aggregation needs a window — and a window needs a rule for events that arrive after it closed. That rule is the whole difference between a correct and an approximate answer.

Three things you can do with a stream:

  1. Write it to a database, cache, search index — "the streaming equivalent of Ch 11's batch use cases"
  2. Push events to users — email alerts, push notifications, a real-time dashboard. "In this case, a human is the ultimate consumer"
  3. Process input streams to produce output streams — a pipeline of stages ← this section

A piece of code that processes streams is an OPERATOR or a JOB. It's closely related to Unix processes and MapReduce jobs, and the dataflow pattern is similar: a stream processor CONSUMES INPUT STREAMS READ-ONLY and WRITES ITS OUTPUT APPEND-ONLY. Sharding and parallelization patterns are also VERY SIMILAR.

THE ONE CRUCIAL DIFFERENCE: A STREAM NEVER ENDS.

⇒ Sorting doesn't make sense on unbounded data, so SORT-MERGE JOINS CANNOT BE USED. ⇒ Fault tolerance must change: "with a batch job running for a few minutes, a failed task can simply be RESTARTED FROM THE BEGINNING; with a stream job THAT HAS BEEN RUNNING FOR SEVERAL YEARS, restarting from the beginning after a crash MAY NOT BE A VIABLE OPTION."

4.1 Five uses of stream processing

The classic four monitoring applications: fraud detection (unexpected credit card usage patterns → block the card) · trading systems (examine price changes, execute per rules) · manufacturing (monitor machine status, quickly identify malfunctions) · military and intelligence (track a potential aggressor, raise the alarm at signs of attack).

(a) Complex event processing (CEP)

Developed in the 1990s. "Similarly to the way a REGULAR EXPRESSION lets you search for patterns of characters in a string, CEP lets you specify rules to search for CERTAIN PATTERNS OF EVENTS in a stream."

CEP engines maintain A STATE MACHINE that performs the matching; on a match, they EMIT A COMPLEX EVENT (hence the name).

THE RELATIONSHIP BETWEEN QUERIES AND DATA IS REVERSED. "Usually a database STORES DATA PERSISTENTLY AND TREATS QUERIES AS TRANSIENT — it searches for data matching the query and FORGETS ABOUT THE QUERY when finished. CEP ENGINES REVERSE THESE ROLES: QUERIES ARE STORED LONG-TERM; AS EACH EVENT ARRIVES, THE ENGINE CHECKS WHETHER IT HAS NOW SEEN A PATTERN MATCHING ANY OF ITS STANDING QUERIES."

(Implementations: Esper, Apama, TIBCO StreamBase; Flink and Spark Streaming also have SQL support for declarative stream queries.)

(b) Stream analytics

"The boundary between CEP and stream analytics is blurry, but as a general rule, stream analytics is LESS focused on detecting specific event SEQUENCES and MORE oriented toward AGGREGATIONS AND STATISTICAL METRICS over large volumes of events."

Examples: rate of a certain event type · rolling average over a time period · comparing current statistics to previous intervals (to detect trends or alert on metrics unusually high or low compared to the same time last week).

Averaging over a few minutes SMOOTHS OUT IRRELEVANT FLUCTUATIONS from one second to the next, while still giving a TIMELY picture of changes.

Probabilistic algorithms: Bloom filters (set membership), HyperLogLog (cardinality estimation), percentile estimation (Ch 2).

⚠️ "This use of approximation algorithms SOMETIMES LEADS PEOPLE TO BELIEVE THAT STREAM PROCESSING SYSTEMS ARE ALWAYS LOSSY AND INEXACT, BUT THAT IS WRONG. THERE IS NOTHING INHERENTLY APPROXIMATE ABOUT STREAM PROCESSING, AND USING PROBABILISTIC ALGORITHMS IS MERELY AN OPTIMIZATION."

(c) Maintaining materialized views

Unlike stream analytics, "considering only events within a certain TIME WINDOW is usually NOT SUFFICIENT. Building the materialized view potentially requires ALL EVENTS OVER AN ARBITRARY TIME PERIOD, apart from any obsolete events discarded by log compaction. IN EFFECT, YOU NEED A WINDOW THAT STRETCHES ALL THE WAY BACK TO THE BEGINNING OF TIME."

"In principle any stream processor could be used, ALTHOUGH THE NEED TO MAINTAIN EVENTS FOREVER RUNS COUNTER TO THE ASSUMPTIONS OF SOME ANALYTICS-ORIENTED FRAMEWORKS that mostly operate on windows of limited duration." (Kafka Streams and ksqlDB support it, building on log compaction.)

Incremental view maintenance (IVM) — why databases aren't enough:

Databases refresh materialized views with periodic batch jobs or on demand (PostgreSQL's REFRESH MATERIALIZED VIEW) — not on every update. That has two drawbacks:

DrawbackWhy
Poor efficiency“All data is reprocessed every time, though most of the data likely remains unchanged.”
Data freshness“Changes are not reflected until the query is run again, during its next scheduled update.”

Triggers work “when the data is easily partitioned and the computation is naturally incremental” — total sales revenue per day means updating one row per sale. But “bespoke solutions work in a few cases, but many SQL queries can't be easily or efficiently converted to incremental computation.”

Incremental view maintenance instead “convert[s] queries written in SQL into operators capable of incremental computations. Rather than processing entire datasets, IVM algorithms recompute and update only data that has changed.” Materialize, RisingWave, ClickHouse and Feldera do this: “recent events are buffered in memory and periodically used to update on-disk materialized views. Reads combine the recent events and the materialized data to provide a single real-time view.”

(d) Search on streams

"Conventional search engines FIRST INDEX THE DOCUMENTS AND THEN RUN QUERIES OVER THE INDEX. By contrast, searching a stream TURNS THE PROCESSING ON ITS HEAD: THE QUERIES ARE STORED, AND THE DOCUMENTS ARE EVALUATED AGAINST THEM, as in CEP."

In the simplest case, test every document against every query — "although this can be SLOW if you have a LARGE NUMBER OF QUERIES. To optimize, IT IS POSSIBLE TO INDEX THE QUERIES AS WELL AS THE DOCUMENTS and thus narrow the set of queries that may match." (Elasticsearch's percolator.)

Examples: media monitoring services searching news feeds for mentions of companies/products/topics; real estate sites notifying users when a matching property appears.

(e) Why actor frameworks are NOT stream processors
ActorsStream processors
PurposePrimarily a mechanism for MANAGING CONCURRENCY and distributed execution of communicating modulesPrimarily a DATA MANAGEMENT technique
CommunicationOften EPHEMERAL and ONE-TO-ONEEvent logs are DURABLE and MULTI-SUBSCRIBER
TopologyArbitrary, INCLUDING CYCLIC request/response patternsUsually ACYCLIC PIPELINES where every stream is the output of one particular job and derived from a well-defined set of inputs

(Some crossover: Storm's distributed RPC farms user queries out to nodes that also process event streams, interleaving queries with events. And you can process streams with actor frameworks — but many don't guarantee message delivery on crashes, so THE PROCESSING IS NOT FAULT-TOLERANT unless you add retry logic.)

4.2 Reasoning About Time — the hardest part

Why batch is easy and streaming is not:

A batch process must look at THE TIMESTAMP EMBEDDED IN EACH EVENT. "THERE IS NO POINT IN LOOKING AT THE SYSTEM CLOCK OF THE MACHINE RUNNING THE PROCESS, because the time at which it is run HAS NOTHING TO DO with the time at which the events actually occurred. A batch process may read A YEAR'S WORTH of historical events WITHIN A FEW MINUTES."

Moreover, using event timestamps makes the processing DETERMINISTIC: RUNNING THE SAME PROCESS AGAIN ON THE SAME INPUT YIELDS THE SAME RESULT.

"Many stream processing frameworks use THE LOCAL SYSTEM CLOCK (the PROCESSING TIME) to determine windowing. This is SIMPLE, and reasonable IF THE DELAY BETWEEN EVENT CREATION AND PROCESSING IS NEGLIGIBLY SHORT. HOWEVER, IT BREAKS DOWN WITH ANY SIGNIFICANT PROCESSING LAG."

Six causes of processing lag: queueing · network faults · a performance issue causing contention in the broker or processor · a restart of the stream consumer · reprocessing past events while recovering from a fault · or after fixing a bug.

Out-of-order arrival: "A user makes request 1 (handled by server A), then request 2 (handled by server B). B's event REACHES THE BROKER BEFORE A's. Stream processors see B then A, EVEN THOUGH THEY OCCURRED IN THE OPPOSITE ORDER."

The Star Wars analogy:

Episode IV was released in 1977, V in 1980, VI in 1983, then I, II, III in 1999, 2002, 2005, and VII, VIII, IX in 2015, 2017, 2019. "If you watched them in the order they came out, THE ORDER IN WHICH YOU PROCESSED THEM IS INCONSISTENT WITH THE ORDER OF THEIR NARRATIVE. (THE EPISODE NUMBER IS LIKE THE EVENT TIMESTAMP, AND THE DATE YOU WATCHED IS THE PROCESSING TIME.) As humans we cope with such discontinuities, BUT STREAM PROCESSING ALGORITHMS NEED TO BE SPECIFICALLY WRITTEN TO ACCOMMODATE THEM."

The concrete damage:

You measure requests per second. You redeploy the stream processor: it is down for a minute, then processes the backlog when it comes back.

downtimeanomalous spikereq/sectime

The actual request rate was steady the whole time. Measured by processing time, the same traffic reads as a gap followed by a sudden anomalous spike that never actually happened.

Figure 12.4.2The concrete damage
Straggler events

"You can never be sure whether you have received ALL the events for a particular window or SOME ARE STILL TO COME." You've counted events in minute 37; time has moved on to 38 and 39. When do you declare minute 37 finished?

You can TIME OUT after seeing no new events for a while. "However, SOME EVENTS COULD BE BUFFERED ON ANOTHER MACHINE SOMEWHERE, DELAYED BY A NETWORK INTERRUPTION."

Two options:

  1. IGNORE the stragglers — "they are probably a small percentage in normal circumstances. TRACK THE NUMBER OF DROPPED EVENTS AS A METRIC AND ALERT IF YOU START DROPPING A SIGNIFICANT AMOUNT OF DATA"
  2. PUBLISH A CORRECTION — an updated value for the window with stragglers included. "You may also need to RETRACT the previous output"

(A third mechanism: a special message meaning "from now on there will be no more messages with a timestamp earlier than t" — a watermark. "However, if several producers on different machines are generating events, each with its own minimum threshold, THE CONSUMERS NEED TO TRACK EACH PRODUCER INDIVIDUALLY. ADDING AND REMOVING PRODUCERS IS TRICKIER IN THIS CASE.")

Whose clock are you using?

A mobile app reports usage metrics. It may be used OFFLINE, buffering events locally and sending them WHEN A CONNECTION IS NEXT AVAILABLE — WHICH MAY BE HOURS OR EVEN DAYS LATER. "To any consumers of this stream, THE EVENTS WILL APPEAR AS EXTREMELY DELAYED STRAGGLERS."

The timestamp SHOULD really be the time the user interaction occurred, per the DEVICE's local clock. HOWEVER, THE CLOCK ON A USER-CONTROLLED DEVICE OFTEN CANNOT BE TRUSTED — it may be accidentally or DELIBERATELY set wrong. The time the server received it is more likely ACCURATE (the server is under your control) but LESS MEANINGFUL in terms of describing the user interaction.

The three-timestamp trick:

①
time the event Occurred, per the Device clock
②
time the event was Sent to the server, per the Device clock
③
time the event was Received by the server, per the Server clock

③ − ② = THE OFFSET between device and server clocks — assuming network delay is negligible relative to required accuracy.

TRUE EVENT TIME ≈ ① + (③ − ②) — assuming the device clock offset did not change between ① and ②.

"This problem is NOT UNIQUE TO STREAM PROCESSING; BATCH PROCESSING SUFFERS FROM EXACTLY THE SAME ISSUES. It is just MORE NOTICEABLE in a streaming context, where we are more aware of the passage of time."

The four window types
WindowShapeImplement by
TumblingFixed length, every event in exactly one window: 10:03:00–10:03:59 | 10:04:00–10:04:59 | 10:05:00–10:05:59Rounding the timestamp down to the nearest minute.
HoppingFixed length, with overlap for smoothing — a 5-minute window on a 1-minute hop: 10:03:00–10:07:59, 10:04:00–10:08:59, 10:05:00–10:09:59First computing 1-minute tumbling windows, then aggregating several adjacent ones.
SlidingAll events within a certain interval of each other: events at 10:03:39 and 10:08:12 are in the same 5-minute sliding window, which tumbling and hopping windows with fixed boundaries would not have grouped.A buffer sorted by time, removing old events as they expire.
SessionNo fixed duration: group all events for the same user that occur closely together; the window ends when the user has been inactive for some time, for example 30 minutes.“Sessionization is a common requirement for website analytics.”

State cost varies enormously: "a COUNTING operation will have ONLY ONE COUNTER regardless of window size or event count. On the other hand, SLIDING WINDOWS OR STREAM JOINS REQUIRE THAT EVENTS BE BUFFERED UNTIL THE WINDOW FINISHES. Therefore, LARGE WINDOW SIZES OR HIGH-THROUGHPUT STREAMS CAN CAUSE STREAM PROCESSORS TO KEEP A LOT OF TEMPORARY STATE."

4.3 Stream Joins — three kinds

Why it's harder than batch: "the fact that NEW EVENTS CAN APPEAR AT ANY TIME makes joins on streams more challenging."

(a) Stream–stream join (window join)

The example: search click-through rate.

search eventsclick eventsjoin on session IDwithin a windowclick-through rate

“The click may never come if the user abandons their search, and even if it comes, the time between search and click may be highly variable — a few seconds, or as long as days or weeks (a user runs a search, forgets about that browser tab, then returns and clicks later). Because of variable network delays, the click event may even arrive before the search event.”

Implementation: maintain state — all events in the last hour, indexed by session ID. On each event, add it to its index and check the other index for a match.

  • Match found → emit “this search result was clicked”.
  • Search event expires → emit “these search results were not clicked”.

Why not just embed the search details in the click event? “That would tell you only about the cases where the user clicked, not about the searches where the user did not click any result. To measure search quality you need accurate click-through rates, for which you need both kinds of event.”

Figure 12.4.4The example: search click-through rate.
(b) Stream–table join (stream enrichment)
activity eventsuser_id, …joinenriched events+ profile infolocal copy of the profile DBin-memory hash table or local-disk indexCDC on the profile databasekeeps the copy up to date

The lookup could be a remote database query; “however, such remote queries are likely to be slow and risk overloading the database.” Better to “load a copy of the database into the stream processor so it can be queried locally without a network round trip. This is a hash join — the local copy might be an in-memory hash table if small enough, or an index on the local disk.”

The difference from batch: “a batch job uses a point-in-time snapshot, whereas a stream processor is long-running and the database changes over time, so the local copy needs to be kept up to date. This is solved by CDC.” “Thus we obtain a join between two streams: the activity events and the profile updates.”

Against a stream–stream join, “the biggest difference is that for the table changelog stream, the join uses a window that reaches back to the beginning of time (a conceptually infinite window), with newer versions overwriting older ones. For the stream input, the join might not maintain a window at all.”

Figure 12.4.5(b) Stream–table join (stream enrichment)
(c) Table–table join (materialized view maintenance)

The social network timeline, as a join:

SELECT follows.follower_id AS timeline_id,
  array_agg(posts.* ORDER BY posts.timestamp DESC)
FROM posts
JOIN follows ON follows.followee_id = posts.sender_id
GROUP BY follows.follower_id

The four events the stream processor must handle:

  • User u sends a post → add it to the timeline of every user following u
  • User deletes a post, or their entire account → remove it from all timelines
  • u1 starts following u2 → add u2's recent posts to u1's timeline
  • u1 unfollows u2 → remove u2's posts from u1's timeline

"The join of the streams corresponds DIRECTLY to the join of the tables in this query. THE TIMELINES ARE EFFECTIVELY A CACHE OF THE RESULT OF THE QUERY, UPDATED EVERY TIME THE UNDERLYING TABLES CHANGE."

The calculus aside, which is genuinely illuminating: "If you regard a stream as THE DERIVATIVE OF A TABLE, and regard a join as A PRODUCT of two tables u·v, something interesting happens: THE STREAM OF CHANGES TO THE MATERIALIZED JOIN FOLLOWS THE PRODUCT RULE (u·v)′ = u′v + uv′. ANY CHANGE OF POSTS IS JOINED WITH THE CURRENT FOLLOWERS, AND ANY CHANGE OF FOLLOWS IS JOINED WITH THE CURRENT POSTS."

Time dependence of joins — the subtle killer

All three require the processor to MAINTAIN STATE derived from one input and QUERY THAT STATE when processing the other. THE ORDER OF THE EVENTS THAT MAINTAIN THE STATE IS IMPORTANT — it matters whether you first FOLLOW and then UNFOLLOW, or the other way round.

In a sharded log, ordering within a single partition is preserved, BUT THERE IS TYPICALLY NO ORDERING GUARANTEE ACROSS DIFFERENT STREAMS OR SHARDS.

"If state changes over time, and you join with a state, WHAT POINT IN TIME DO YOU USE FOR THE JOIN?"

The tax rate example: "If you sell things, you need to apply the right tax rate, which depends on country/state, product type, and DATE OF SALE (since tax rates change). WHEN JOINING SALES TO A TABLE OF TAX RATES, YOU PROBABLY WANT THE TAX RATE AT THE TIME OF THE SALE — WHICH MAY DIFFER FROM THE CURRENT RATE IF YOU ARE REPROCESSING HISTORICAL DATA."

"If the ordering across streams is undetermined, THE JOIN BECOMES NONDETERMINISTIC — YOU CANNOT RERUN THE SAME JOB ON THE SAME INPUT AND NECESSARILY GET THE SAME RESULT."

The warehouse solution — slowly changing dimensions (SCD):

Use A UNIQUE IDENTIFIER FOR A PARTICULAR VERSION of the joined record — every time the tax rate changes it gets a new identifier, and the invoice includes the identifier for the rate at the time of sale.

This makes the join DETERMINISTIC, "but it has the consequence that LOG COMPACTION IS NOT POSSIBLE, since ALL VERSIONS of the records need to be retained. Alternatively, you can DENORMALIZE the data and include the applicable tax rate DIRECTLY IN EVERY SALE EVENT."


On this page