ContentsThe library

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.

FIG 1A sort asked to produce its first row
steptimearrived so farsmallest seencan the first row be emittedwhat happened
110:0041, 17, 8817noSeventeen is the smallest so far. The next arrival could be smaller, so emitting it now could be wrong.
210:0141, 17, 88, 99noAnd there it is. Emitting 17 one second ago would already have been an error.
314:00four million values2noFour hours and four million values later, the situation is unchanged in kind. Nothing about having more data makes the next arrival predictable.
4neverstill arrivingwhatever it isnoThis 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.
4 steps
The fourth row is the one that changes how you think about the rest of the course. An operation whose first output depends on the final input is not an operation you can run on a stream at any speed, on any hardware, with any budget. It needs replacing rather than optimising.

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.

FIG 2Familiar operations against what survives on a stream
works as isneeds a windowexact within the windowwants an approximation
count1000
sum1000
sort0110
median0111
exact distinct count0101
maximum0110
filter1000
The first column is the operations that are already incremental: each arrival updates the answer and nothing has to be revisited. Everything else needs a window around it to have an end again, and the two marked rows need something more, because even inside a window they cost memory proportional to the data unless an approximation is accepted.

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.

FIG 3Choosing the route for an operation
Three of the four leaves are the three escapes, and the fourth is the trap. Running this decision over each operation in a pipeline takes a few minutes and reliably finds the one nobody had thought about, which is usually something keeping per-key state.

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.

FIG 4The memory a per-key aggregate needs
memory held
keys seen so far
the rate at which new keys appear
how long the job has been running
bytes per key
The term that matters is the one multiplying t. A job holding sixty bytes per key with a hundred new keys an hour is holding an extra five megabytes a month, which is invisible for a quarter and is an out of memory failure in the second year. The failure arrives long after the code was written and is attributed to load.

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

NextTwo Different Clocks →

The rest of this course

  1. 01Sort This. There Is No Last Row.you are here
  2. 02The Spike at Nine Was Six Hours of Traffic Arriving at Onceopening only
  3. 03Three Shapes, and the Question Tells You Which Oneopening only
  4. 04Nothing Can Tell You That Nothing Else Is Comingopening only
  5. 05Somebody Has Already Seen the Number You Are About to Changeopening only
  6. 06An Average Needs Two Numbers, a Median Needs All of Themopening only
  7. 07Twelve Kilobytes Will Count a Billion Different Thingsopening only
  8. 08The Windows Came Back, and So Did Nine Minutes of Totalsopening only

Read alongside