ContentsThe library

Counting Things as They Arrive

An Average Needs Two Numbers, a Median Needs All of Them

Last timeThe One That Arrives Afterwards

Some aggregates can be kept as a tiny fixed state and updated one event at a time. Others need every value you have ever seen. The difference is structural and it decides your design.

The four pieces

A window has to be aggregated, and the aggregate has to be computed without

keeping the events. The test for whether that is possible is mechanical.

An aggregate streams if you can write down four things. A state, which is some

fixed number of values. An update, which folds one new event into the state. A

merge, which combines two states into one. And a finish, which turns a state

into the answer somebody asked for.

If all four exist and the state does not grow with the number of events, the

aggregate fits in bounded memory no matter how long the stream runs. If the

state has to grow, it does not fit, and this is a property of the aggregate

rather than a limitation of your implementation.

FIG 1The same four pieces, for an aggregate that fits and one that does not
plaintext
state: a pair, sum and count
init:    s = 0, n = 0
update:  s = s + x, n = n + 1
merge:   s = s1 + s2, n = n1 + n2
finish:  mean = s / n

state: every value ever seen
init:    values = empty list
update:  add x to values
merge:   concatenate the two lists
finish:  sort, then take the middle one
Both blocks have all four pieces, so the existence of the pieces is not the test. The test is the first line of each. The top state is two numbers whether the window held ten events or ten billion. The bottom state is the whole input, and its merge step concatenates rather than combines, which is the signature of an aggregate that does not stream. Reading a proposed aggregate this way takes a minute and settles the design question before any code is written.

The merge step is the one people leave out, and it is the one that matters most

in practice. Without it the aggregate cannot be computed in parallel, cannot be

split across partitions, and cannot reuse work between overlapping windows.

The ones that fit in three numbers

Count is one number. Sum is one number. Minimum and maximum are one number each.

Those are obvious and they are also the only aggregates most people are

confident about.

A mean is two numbers, a sum and a count, and the common error is to keep a

running mean instead. A running mean updates fine and merges incorrectly: given

a mean of 80 from one partition and 60 from another, there is no way to combine

them without knowing how many events each came from. Keep the sum and the count

and the merge is addition.

The lesson stops here

5 more paragraphs to go

You have read the opening. The rest of the argument, the problems that check whether it landed, and the lines worth keeping at the end all come with a plan.

The first lesson of every course in the library reads the whole way through, free, so you can see exactly what the rest of them are.

See the planThe contents

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

The rest of this course

  1. 01Sort This. There Is No Last Row.
  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 Themyou are here
  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