Learn Labs
12. Stream Processing

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.