Galewatch is a company that exists only in this course: it takes a reading every two seconds from each of the wind turbines on the farms it monitors, and sells the owners dashboards. At two minutes past eleven on a Thursday morning in July, one of its engineers deployed a small process that was meant to watch readings go past. For the next fourteen minutes it ate them instead.
The watcher was there to answer the oldest complaint in the product. Lesson 003 has an engineer reading a green tile for a turbine that had been feathered since eleven that morning, and everybody wanted a process that sees every reading as it arrives and speaks up when one turbine sits at zero output for a minute while its neighbours turn.
Three weeks earlier the company had finally built what lesson 017 argued for and nobody had got round to. Between the four ingest boxes and the nineteen Postgres machines there was now a broker, and the shape they chose was a log: a stream called readings, cut into twenty four partitions, one farm per partition, seven days of retention, drained by six writer processes that insert each reading into the shard lesson 014's routing table gives it, plus a second stream for replayed readings, which gets a section of its own below.
They wired up the metric 017 asked for. A log gives you two versions of it: the count of records between a reader's position and the end of the partition, which is 017's depth under a new name, and the age of the record sitting at that position, which is the one 017 told you to alert on. Galewatch graphs both and calls the pair lag.
The watcher's configuration was copied from the writer's, because that is where a working configuration lives. One line came with it: group.id.
So the watcher was not a second reader of the stream. It was two more pairs of hands on the team that was already reading it, and every partition belongs to exactly one pair of hands.
At 11:02 the group reshuffled. Eight members now, twenty four partitions, three each, and the watcher's two processes were handed six of them, every one with a farm on it: six farms, 360 turbines between them, 180 readings a second. Those readings went to a process whose entire job was to look at them and decide nothing was wrong. The writer never saw them. Postgres never heard about them.
Nothing alerted. The lag graph was beautiful. Every partition's reader was within a second of the newest record, because there was a reader on every partition and it was keeping up easily. The group was healthy. It just had the wrong members.
At 11:09 the engineer at Kilmore Sands, a hundred and forty turbines and the largest farm that still fits on one machine, rang to ask whether the panel was broken. The power curve had stopped advancing seven minutes earlier, and the only reason it took a person seven minutes to say so is that the staleness number lesson 012 argued for, and lesson 003 found missing from this exact dashboard, is still not on it.
At 11:14 the on-call listed the group's members and found eight where there should have been six, two of them carrying a client id with the word watcher in it. At 11:16 the watcher was stopped. The group reshuffled again, the writer took all twenty four partitions back, and it resumed from the positions the group had committed, which the watcher had been diligently advancing all along. Every panel went current inside a minute. The fourteen minutes stayed missing.
Here is the part that was not possible three weeks earlier.
At 11:21 they moved the group's position on those six partitions backwards. A reset like that is done by time rather than by counting records, and the time on a record is the one the ingest box stamped when it accepted the reading, which is the honest clock in this system: lesson 034 established that the four ingest boxes share a time server and agree with each other to a millisecond or two. Had the turbines been producing into the log themselves, the stamp would have been a turbine's own clock, and 034 found one forty one seconds out.
They reset to 11:00, twenty one minutes of history rather than the fourteen they needed, because overlap is free and being approximately right in the cheap direction takes no thought. Reading forward, the writer put away 151,200 readings that had never reached a database and 75,600 that already had, and lesson 019's unique constraint on (turbine_id, recorded_at) refused every one of the second set without comment. It read those partitions at about five times the rate the farms were filling them, which is a dial they chose rather than a limit they hit, so twenty one minutes of records took a little over five, and the panels walked forward from 11:00 while it caught up. That is lesson 017's own warning arriving as a feature: drain a backlog in order and somebody watching the glass sees the morning replay in fast forward. At 11:26 all twenty four partitions were current and complete.
On the queue 017 actually described, those 151,200 readings would be gone. The ingest box had answered 200, so the turbine had already deleted its own copy, which is the dial 017 named and the only-copy problem lesson 018 built a section on. A queue hands a message out once and forgets it the moment somebody says they are done.
Fourteen minutes of a configuration file being right about everything except one string. The bill was six farms whose panels could not be trusted from 11:02 to 11:26, fifteen of those minutes showing nothing new at all, and not one reading lost.
Nothing leaves when you read it
A log here is a file you only ever append to. Writes go on the end, nothing in the middle is edited, and every record gets an offset, which is its position counted from the start of the file. A topic is a named stream of records. A partition is one of the logs a topic is cut into, with its own offsets starting at zero. Lesson 034 used the word offset for how wrong a clock is, and this is a different one; the industry uses both and so does the course.
Lesson 017's queue and this are not two settings on one machine. They keep different things.
A queue keeps state per message. Lesson 036 opened it: a message handed out is hidden, with a timer on the hiding, and a count of how many times it has been handed out, and a flag for whether somebody has acknowledged it. Marlow Books, the four person bookshop, runs about one message a minute through that machinery, so the cost never comes up. Put Galewatch's stream through it. Lesson 036 did the division: 1,240 turbines at a reading every two seconds is 620 a second, and with the thirty second hold 036 named, up to 18,600 messages can be in flight at once, each with a timer the broker has to fire on time. That ceiling is only reached when the handlers are slow, which is the case 036 spent a whole lesson on, so it is the number a broker has to be built for.
A log keeps one number per reader per partition. Galewatch has two readers of that stream today and twenty four partitions, so the whole of its bookkeeping is forty eight integers, updated a few times a second. That is the trade in one line: a queue remembers what it has given you, a log remembers where you have got to.
| Lesson 017's queue | Today's log | |
|---|---|---|
| kept per message | a hold, a timer, a count | nothing |
| kept per reader | nothing | one offset per partition |
| what a read does | hides it, then deletes it | nothing at all |
| a second reader costs | the producer writing twice | a different group id |
| where a backlog sits | in the broker | on disk, where it already was |
The writing is cheap for the same reason. An append is sequential, and lesson 002's ladder has an SSD streaming about two gigabytes a second against a hundred microseconds for a random read. Galewatch's whole fleet at 200 bytes a reading is 124 kilobytes a second, one part in sixteen thousand of what one drive will stream, and that is a bound from the hardware rather than a measurement of anybody's broker. Stagefront, the ticketing service where two hundred thousand people press the same button at ten in the morning, offers 3,333 requests a second at the on-sale minute by lesson 017's count, and at the few hundred bytes 017 sized a message at, appending those is about a megabyte a second. Appending has never been Stagefront's problem. Two buyers wanting one seat is, and lessons 015 and 035 are where that lives.
Falling behind is a disk question
Retention decides how long a record stays. Set it by time, where a week is the usual starting point, or by size, and the size one is per partition, which is where people get their capacity arithmetic wrong by a factor of however many partitions they have.
Records are not deleted one at a time. A partition on disk is a run of segment files, and retention closes a segment and later deletes the whole thing, so you always have a little more history than you asked for and never less.
Price a week for Galewatch. 620 readings a second at 200 bytes is 10.7 gigabytes a day, which lesson 034 printed as under eleven, and seven of those is 75 gigabytes. Lesson 013 sized a Postgres shard at 720 gigabytes, so a week of every reading from every turbine in the fleet is a tenth of one of the nineteen machines the company already runs. The fourteen days of the same readings sitting in the hot tables is twice that before a single index.
Retention is how long you are allowed to be wrong for. The Thursday was fourteen minutes against seven days and was never close to the edge, which is the entire reason it reads as an anecdote rather than an incident.
That changes what a backlog is. Lesson 017 worked out the peak depth of a four hour replay from Tarrow Ridge, a sixty turbine farm on a weather-beaten link, and got 259,000 messages and 52 megabytes sitting in a broker, the whole argument being whether a broker can hold what you hand it. In a log the question dissolves: the records are on disk for seven days whether anybody reads them or not. A consumer four hours behind has not made the broker do anything. It has a smaller number than the writer does.
The same property has a sharp edge, and the edge is worth more than the saving.
The offset is one number, so you can only ever be done with a prefix. There is no way to tell a log that you handled record 48,112 and not 48,110.
Lesson 018 has a Friday where forty messages written by an old producer raised a KeyError in a new consumer, and 036 has already been back to it, so take only the part that is about today. Those forty blocked the queue, and what unblocked it was the broker moving them into the drawer 036 named, which it can do because it is holding each one individually. A log holds nothing. The record stays where it was appended, in the order it was appended, and a consumer that cannot deal with the record in front of it either stops, taking everything behind it in that partition, which at Galewatch is one farm's dashboard going flat, or steps over it and loses it on purpose. The third option is to write the record to a topic of your own and carry on, which is building 036's dead letter queue by hand, with the same two sets to keep apart afterwards.
Your group id is the whole decision
A consumer group is a set of processes that share one group id and split the partitions of a topic between them. Each partition goes to exactly one member. The broker keeps a committed offset per group per partition, which is why two groups reading the same topic never notice each other.
Two group ids are two independent readers of every record. One group id is one reader with more hands.
That is the whole of the July Thursday. The watcher went back out on the Monday with four characters changed.
The independence cuts the other way too, and it is the part of the Thursday nobody wrote down. The readings went back into Postgres because the writer's group had its offsets moved. The watcher is a different group with a position of its own, and a group with a new id starts at the end of the log by default, so nothing ever re-read those fourteen minutes looking for a turbine sitting at zero. If one of those 360 turbines was feathered at 11:08, the process built to speak up about exactly that has still not noticed. Replaying a log replays it into one group's consumers, not into all of them.
The same rule is the ceiling on how fast a group can go. Twenty four partitions means at most twenty four useful processes in it; the twenty fifth joins, gets nothing, and sits there costing money.
When a member joins or leaves, the group has to decide who owns what again, and that is a rebalance. The classic version stops everybody: all members give up their partitions, wait for an assignment, and start again. Price it on Galewatch's writer. Six processes restarted one at a time is twelve rebalances, since each one leaves and then rejoins. Call a rebalance two seconds, which is about what a clean shutdown and a fresh join cost when nobody is parked in a long poll; kill a process without warning instead and the group waits out its session timeout, which is tens of seconds. So an ordinary deploy stops the writer for roughly twenty four seconds and builds about fifteen thousand readings of lag.
Nothing is lost; the lag just drains.
What it costs is the product. Galewatch sells a dashboard whose value is that it is current, which lesson 008 called the read that cannot be cached at any TTL, and the shape above makes every panel in the product half a minute stale during every deploy. Both mitigations are things you turn on rather than things you get: static membership, which gives a process a stable instance id so that a quick restart reclaims its own partitions instead of triggering a reassignment, and a cooperative assignor, which moves only the partitions that have to move rather than all of them. Kafka's 4.0 line moves the coordination into the broker and reassigns incrementally, which is the real fix and is new enough that plenty of running clusters have never met it.
One thing to collect rather than rediscover. Lesson 036 said that in Kafka the committed offsets are themselves a topic, and that is literally true: there is an internal topic holding them, compacted, cut into fifty partitions out of the box. A commit is a record appended to a log. That is the whole reason Kafka can put a commit inside a producer transaction, which is the exactly-once processing 036 priced.
It also means 036's question arrives here in new clothes. Where you commit the offset is where you put the acknowledgement. Commit before you do the work and a failure in the gap loses it; commit after and a failure in the gap repeats it. Same two placements, same absence of a third, one line of code.
A partition key is a shard key with different failure modes
Records carrying the same key go to the same partition and stay in the order they were appended. Across partitions there is no order at all, for the reason lesson 034 gave: nothing talks, so nothing is ordered. Choosing the key is therefore choosing whose order you keep and where your load lands, which is lessons 013 and 014 arriving again at a different layer.
Galewatch keys by farm, one farm to a partition, and the partition comes from an explicit map, a column alongside the routing table lesson 033 put on three machines, rather than from a remainder. Lesson 014 had already made the farm layer a table, and lesson 031 had already said that a twenty key ring is boulders. Twenty four partitions for twenty one farms leaves three spare, and lesson 005 established that this company grows when sales signs a contract, so three spare is three contracts of notice.
Skew comes with it, as 014 promised, and it bites harder here than in the database. The fleet average is 620 over twenty one farms, which is 29.5 a second. Kilmore Sands at a hundred and forty turbines is seventy, 014's own figure, so 2.4 times the average. Ardnave Point is the one to look at, though. Lesson 014 widened the database's key to the pair (farm, turbine) precisely so that farm could be split, and 031 split it three ways; the log's key is the farm alone, so three hundred turbines are one partition at a hundred and fifty readings a second, five times the average, handed to a writer that then fans them out across three machines. Neither number is anywhere near the 5,000 writes a second lesson 005 derived from the hundred connection limit and lesson 017 put on a single shard.
The part nobody sells you on is that lag concentrates. A writer process that gets slow makes one farm's dashboard stale while the other twenty are perfect, instead of making the whole fleet slightly late. That is lesson 027's partial outage arriving by choice rather than by accident. It is also a much easier phone call, because the farm whose panel is wrong is the only customer who needs telling.
Now the failure mode a shard key does not have. Resharding a database moves data, and lesson 031 priced Galewatch's own: 665 gigabytes and fifty five minutes to move 133 turbines when a divisor went from 2 to 3. Adding a partition to a topic moves nothing, which sounds like the better deal until you look at what it breaks. If the producer picks the partition with hash(key) % n, then the day n changes, that key's records start landing somewhere else, and the old ones do not follow, because an append-only log cannot be rewritten. The key's history is now split across two partitions with its order broken at the seam, for good. Lesson 031's ring does not rescue you here, because what you lose is not location but order. And you cannot take a partition away afterwards either.
Galewatch walks past all of that, and the reason is small and worth keeping: its partition comes from an explicit map, so the twenty fifth farm takes a partition with no history in it to split.
Order has one more bill to present. Lesson 017's best idea was two lanes, a live one and a backfill one, so that four hours of a turbine's replayed flash does not sit in front of the reading that just arrived. A partition is strictly ordered, so inside one there is no such thing as serving the newer record first, and a second consumer group does not help because it reads the same order. Two lanes have to be two topics. The ingest box decides which, using the two columns lesson 034 built: a reading whose recorded_at is more than a minute behind its received_at came out of a turbine's flash, so it goes to the backfill topic. The writer subscribes to both and pauses the backfill partitions whenever the live ones hand back a full batch.
A log gives you fan-out for nothing and takes priority away.
Where the log stops being a database
Set a topic's cleanup policy to compact and the broker's job changes. Instead of deleting old records by age, it keeps the most recent record for every key and throws the superseded ones away, and a record whose value is null is a tombstone meaning this key is gone. What is left is a table: every key, its current value, rebuildable by anyone who reads the topic from offset zero and keeps a map.
Galewatch's routing table is exactly that shape. A key per farm, a value saying which machine, and four ingest boxes that reload the whole table every five seconds because lesson 033 shrank a twelve minute window of disagreement down to five seconds by polling harder. As a compacted topic each box reads it once on start and then tails it, and both the poll and the five second window go away. Nobody at Galewatch has done this, and lesson 048 owns getting configuration to the boxes that need it, so I will leave it there.
The slogan hides three limits and one precondition.
A log has one access path: read a partition from an offset. That is the whole menu. Ask it what turbine DM-11 did in the last hour and the answer is a scan of a week, because the index lesson 010 built to make that question cost ten milliseconds is a Postgres index and a log has nothing of the kind. Replay is not a query.
Compaction keeps the latest value per key, which means it throws history away by design. A compacted topic is a snapshot that happens to be shaped like a log. Ask it what the routing table said last Tuesday and it has no idea, and lessons 059 and 073 are where wanting that answer leads.
Replaying takes the time your consumers take, not the time the disk takes. Galewatch's week is 75 gigabytes, which an SSD streams in under a minute at lesson 002's rate. The same week is 375 million readings, and Kilmore Sands' share of it is 42 million through a writer whose shard tops out at the 5,000 writes a second lesson 005 derived, which is a little over two hours and twenty minutes. The log was never the slow part.
Then the precondition, which is the one I would put on the wall. Replay is a property of the log. Being able to use it is a property of your consumer. Galewatch's writer can be replayed all afternoon because lesson 019's constraint is on the reading's own identity and a second copy is refused, which is the second pile. Point a replay at a consumer that sends email or charges a card and you are in lesson 036's third pile, where every record you reread becomes a thing that happens to a customer twice.
Marlow should not run any of this, and saying why is the useful half. Its confirmation queue carries a fifth of a message a second at Christmas by 017's count, and 036 put ninety minutes of an ordinary March evening at a hundred and eight of them. A log sells replay, fan-out and retention. The shop has one consumer, has never been asked for a second reader, and its effect lives at an email provider, so replaying the log means sending the email again, which is precisely the bug it spent a week fixing in 036. Against that, partitions it cannot take back, a group to understand, retention to size, and a cluster with a quorum of its own to keep alive. Lesson 017 gave the shop a queue, and 017 was right.
Recap
Nothing leaves when you read it. A queue hands a message out and forgets it on an acknowledgement; a log appends a record and hands every reader a number. At Galewatch that is forty eight integers of bookkeeping where a queue has to be ready for 18,600 live timers, and a reader falling behind stops being the broker's problem.
Falling behind is a disk question, not a memory question. Retention is how long you are allowed to be wrong for. A week of Galewatch's entire fleet is 75 gigabytes, a tenth of one machine lesson 013 sized, which is why fourteen minutes of a bad group id cost six farms a morning's worth of untrustworthy panels instead of 151,200 readings.
The offset is one number, so you can only be done with a prefix. The property that makes a log cheap is the property that makes a poison record worse here than in a queue. There is no drawer unless you build one, and then you own both of lesson 036's two sets.
Your group id is the whole decision. Two group ids are two independent readers of every record; one group id is one reader with more hands. A rebalance loses nothing and makes the product briefly wrong, and how much that matters is how much of your product is freshness.
A partition key is a shard key whose failure mode is order rather than location. Adding partitions moves no data and splits a key's history at the seam, and an append-only log cannot be rewritten to tidy that up. An explicit map instead of a remainder is why Galewatch's twenty fifth farm is free.
Compaction is where the log stops being a history and becomes a table. Latest value per key, tombstones for deletes, rebuildable from offset zero, and no answer at all to what the table said last Tuesday.
Replay is a property of the log, and being able to use it is a property of your consumer. Nothing is deleted, so you can always go back. Whether going back is safe is lesson 019's question with lesson 036's answer, and no broker setting changes it.
Check your understanding
Take the July Thursday and change one thing: the topic keeps three days instead of seven, and the engineer who deployed the watcher is on leave until Monday. Say what Monday morning costs, and say which decision made in June you would go back and argue about.
A colleague wants to replace a team's queue with a log, and the reason given is that replay will let them recover from bad deploys. The consumer calls a payment provider. Say what you would ask before agreeing, and what would have to be true of that consumer first.
You own a topic with twelve partitions keyed by customer id, and one customer has grown into a quarter of all traffic. Lay out the options you actually have, say which one you would refuse, and say what you would measure before choosing.
A consumer group in production rebalances every few minutes and nobody knows why. Say what you would look at, in order, and name the earlier lesson each thing belongs to.
Galewatch wants a second reader that produces a daily per turbine rollup and nothing else. Argue for either a new consumer group on the live topic or a nightly job over Postgres, using this lesson's figures, and say what would change your mind.
Somebody proposes making the compacted routing topic the source of truth for which machine owns which farm, and retiring the Postgres table lesson 033 built. Say what breaks the first time an engineer asks what the table said last Tuesday, and say what you would keep alongside it.
Next lesson
038 Event-Driven Architecture and Its Failure Modes. Today put a durable, replayable log in the middle of a system and made two readers of it independent; next lesson asks what happens to a whole architecture when services stop calling each other and start publishing events instead, and which of its failures are new.