12. Stream Processing
Chapter 12 of Designing Data-Intensive Applications — 12 sections.
"A complex system that works is invariably found to have evolved from a simple system that works. The inverse proposition also appears to be true: A complex system designed from scratch never works and cannot be made to work." — John Gall
The assumption Ch 11 quietly made, now removed:
Batch processing assumed the input is BOUNDED — of a known and finite size — so the process knows when it has finished reading. (MapReduce's central sort MUST read its entire input before producing output, because THE VERY LAST INPUT RECORD COULD BE THE ONE WITH THE LOWEST KEY that needs to be the very first output record.)
In reality, a lot of data is UNBOUNDED because it arrives gradually over time. Users produced data yesterday and today, and will produce more tomorrow. UNLESS YOU GO OUT OF BUSINESS, THIS PROCESS NEVER ENDS, SO THE DATASET IS NEVER "COMPLETE" IN ANY MEANINGFUL WAY.
Batch artificially divides data into fixed-duration chunks — a day's worth processed at the end of every day. The problem is that changes are reflected a day later, which is “too slow for many impatient users”.
You can run the job more frequently: a second's worth every second. Push that idea to its limit and you arrive at stream processing — abandon fixed time slices entirely and process every event as it happens.
A stream = data incrementally made available over time. The concept appears in Unix stdin/stdout, lazy lists, FileInputStream, TCP connections, audio and video over the internet.
Sections
- 12.13Transmitting Event Streams
- 12.22Log-Based Message Brokers
- 12.34Databases and Streams
- 12.45Processing StreamsExamples: 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…
- 12.5Fault Tolerance
- 12.6Deep divesTechnology deep dives
- 12.7Failure catalogProduction failure catalog for this chapter
- 12.8Decision sheet1Decision cheat sheetYes if consumers must be able to bootstrap from the log without a separate snapshot — which is the whole point of "rebuild a derived system from scratch." No if events are intent-…
- 12.9Worked examplesWorked examplesStorage = 200e6 × 86,400 × 7 ≈ 121 TB (× replication factor 3 = 363 TB).
- 12.10Self-testSelf-test
- 12.11TerminologyTerminology introduced here
- 12.12Forward linksForward links