ContentsThe library

Counting Things as They Arrive

The Windows Came Back, and So Did Nine Minutes of Totals

Last timeClose Enough, in Far Less Room

Everything in this course lives in memory on a machine that will be restarted. Surviving that means saving state and position together, and it means the output system has to cooperate.

Three things, not one

Everything built in this course so far is in memory: the open windows and their

accumulators, the running totals, the sketch registers, the believed point, the

timers waiting to close windows. All of it sits on a machine that will be

restarted this week, for a deploy or a crash or a rescheduling you did not ask

for. What has to survive is more than the accumulators, and the list is short

enough to memorise.

The window state, meaning every open window's accumulator for every key. The

position in every input partition, meaning exactly how far along each one has

been read. And the pending timers, meaning the scheduled closings that have not

fired yet, because a window that was due to close at 10:05 must still close

after a restart at 10:03.

Any two of those without the third is not resumable. Accumulators without

positions cannot be continued at all, because nothing says where to start

reading. Positions without timers lose the closing of every window already due,

which shows up as a few windows that are never reported and never missed.

One snapshot, not two writes

Suppose you save the state and the position as two separate writes. This is the

obvious implementation and it is wrong in a way that has only two flavours.

FIG 1Two writes, two failures, and the fix
Note that both failure branches are silent. Neither produces an error, a failed task or a log line that says anything is wrong, and the symptom in both cases is a number that is a little off in a window somebody has probably already read. This is the same shape as every failure in this course: the pipeline looks healthy and the figure is wrong.

The correct form is a single snapshot that contains the state and the positions

together, so that restoring it puts the job in a condition where resuming from

those positions reproduces exactly those accumulators.

Doing that across a distributed job without stopping the stream is the one

genuinely clever mechanism in this course. A marker is injected at the sources

and flows along with the events. Each operator, on seeing the marker on all of

its inputs, writes its own state and forwards the marker onward. The set of

states written for one marker is a consistent cut of the whole job: everything

before the marker is included everywhere, everything after it is included

nowhere. The job never pauses, and the resulting snapshot is as good as if it

had.

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 Themopening only
  7. 07Twelve Kilobytes Will Count a Billion Different Thingsopening only
  8. 08The Windows Came Back, and So Did Nine Minutes of Totalsyou are here

Read alongside