Counting Things as They Arrive
Sort This. There Is No Last Row.
Nearly everything you know about working with data assumes a final row. Remove it and half the familiar operations become impossible rather than slow.
What unbounded actually means
Every query anybody writes for the first several years assumes the data stopped.
It is such a deep assumption that it is almost invisible, so it is worth
surfacing before anything else.
Sorting assumes an end. You cannot emit the first row of a sorted result until
you have seen every row, because the last one might belong at the front. An
average assumes an end: it is a sum over a count, and both are over a set that
is finished. A maximum assumes an end. Order by with a limit assumes an end.
Every one of those is taught as an operation on data, and every one of them is
an operation on a finished collection of data.
Remove the end and something stronger than slowness happens.
The property that matters is not size. A very large file is finished and
sorting it is an engineering problem with known answers. The property is that
there is no moment at which you have all of it. Any answer computed at any
instant is an answer about a prefix, and the rest of the input, which has not
arrived, can change it.
| step | time | arrived so far | smallest seen | can the first row be emitted | what happened |
|---|---|---|---|---|---|
| 1 | 10:00 | 41, 17, 88 | 17 | no | Seventeen is the smallest so far. The next arrival could be smaller, so emitting it now could be wrong. |
| 2 | 10:01 | 41, 17, 88, 9 | 9 | no | And there it is. Emitting 17 one second ago would already have been an error. |
| 3 | 14:00 | four million values | 2 | no | Four hours and four million values later, the situation is unchanged in kind. Nothing about having more data makes the next arrival predictable. |
| 4 | never | still arriving | whatever it is | no | This is the point. The answer is not slow or expensive. It does not exist, because the condition under which it would be correct never occurs. |
The operations that need an end
The test is a single question: can the first output be produced before the last
input has been seen? If no, the operation needs an end.
| works as is | needs a window | exact within the window | wants an approximation | |
|---|---|---|---|---|
| count | 1 | 0 | 0 | 0 |
| sum | 1 | 0 | 0 | 0 |
| sort | 0 | 1 | 1 | 0 |
| median | 0 | 1 | 1 | 1 |
| exact distinct count | 0 | 1 | 0 | 1 |
| maximum | 0 | 1 | 1 | 0 |
| filter | 1 | 0 | 0 | 0 |
Filter and count survive untouched. A filter looks at one record and decides;
nothing about later records changes that decision. A count goes up by one.
These are the operations that were never secretly about a finished set.
Sorting, the median, the maximum and an exact distinct count all fail the test,
and they fail it for two different reasons that are worth keeping apart.
Sorting and the maximum fail on correctness: the answer can be invalidated by a
later arrival. The median and an exact distinct count fail on memory: the
median needs every value kept until the end, and an exact distinct count needs
every distinct value kept forever. A stream has no end, so both grow without
limit.
Joins deserve their own mention because they fail in a way people do not expect.
Joining two streams means holding records from one side waiting for a match on
the other, and the waiting has no natural end: a record might match something
that arrives tomorrow. Every streaming join therefore carries a time bound,
stated or implied, and the records outside it are silently unmatched.
Bound it, update it, or approximate it
There are three ways out and every streaming system is some arrangement of
them.
A window is the most common answer and the rest of this course is largely about
it. The idea is simple and the consequences are not: by putting a boundary
around a range of time, you make the data inside it finished, which makes every
familiar operation available again inside it. What you give up is any answer
about all of history, and what you take on is the problem of deciding when a
window is actually over, which turns out to be the hard middle of the course.
The state that grows forever
The last concept is the one that catches working systems rather than exercises,
so it is worth more attention than its size suggests.
An operation can be perfectly incremental and still grow without bound, if it
keeps state per key. A running total per customer updates on each arrival and
needs no window. Its memory is one entry per customer, which is bounded if the
set of customers is bounded, and unbounded if customers keep appearing.
- memory held
- keys seen so far
- the rate at which new keys appear
- how long the job has been running
- bytes per key
Session identifiers, request identifiers, device identifiers, anything derived
from a URL or a user agent: all of these look like reasonable keys and all of
them are unbounded over a long enough period. A job that works for six months
and then dies on a Sunday night is usually this.
Every piece of per-key state on a stream needs an explicit answer to one
question: what removes an entry. The available answers are a time to live after
which an untouched key is dropped, an explicit end-of-session event, a window
that closes and takes its state with it, or a bound on the number of keys with
the least recently used ones evicted. Each of those is a decision with
consequences for correctness, and none of them is a default.
The next lesson introduces the distinction that everything after it depends on.
Every event carries two times, the moment it happened and the moment it reached
you, and almost every confusing streaming result comes from a system that used
the wrong one.
Recap
- An unbounded input has no last row, which makes any operation whose first output depends on the final input impossible rather than merely expensive.
- The three escape routes are the same every time: bound it by a window, make it incremental, or accept an approximate answer in fixed memory.
- Anything that keeps per-key state on an unbounded input grows without limit unless something expires it, and that expiry is a design decision rather than a tuning parameter.
This is the reading half
Starting the course gives you your own copy of it. Every idea on every page has problems standing under it, marked with a reason rather than a tick, and any sentence you do not believe can be opened and argued with. None of that can happen on a page nobody owns.
The contents