High-Throughput Streaming
Picture a giant school where every little thing that happens is written on a note: a snack was bought, a book was borrowed, a goal was scored. The notes ride along conveyor belts, one after another, all day and all night. A note is not thrown away after someone reads it. It stays on the belt for a week, so the lunch staff, the librarian and the coach can each walk along and read every note at their own speed. To keep up with millions of notes, the school runs many belts side by side and hires teams of helpers to read them. This page is about splitting notes across belts, sharing the reading work, doing math while the notes fly past, and delivering one post to everyone who wants it.
Event streams
A conveyor belt of notes that stays put for a while
An event is a small note about something that already happened, like this one: Mia bought an apple at 12:03. An event stream (Kafka is a well known example) is an endless line of these notes, kept in the order they arrived. New notes are only ever added at the end. Old notes are never changed or shuffled around.
This is different from a queue. A queue is like the teacher's in-tray: a helper takes a note out, does the job, and the note is gone. A stream is more like the class diary. Reading a page does not tear it out. Many different readers can read the same notes, each at their own pace, and they can even go back and read old pages again, for example to redo some math after fixing a mistake.
Each note has a number that shows its place in line, called its offset. Every reader keeps a bookmark with the number of the next note it will read. Notes are kept for a set time (like 7 days) or until the belt reaches a size limit. Then the oldest notes are deleted, whether or not anyone read them.
The price is space. Keeping every note for days takes a lot of disk. And because the belt does not track who handled what, each reader must look after its own bookmark.
- Notes are added only at the end, and are never changed.
- Reading a note does not remove it.
- Each reader keeps its own bookmark.
- Old notes are deleted after a set time or size, read or not.
Remember
A stream is a diary that many readers share, not an in-tray that empties.
Partitioning
Splitting one busy belt into several belts
One belt can only move so many notes each second, and one helper can only read so fast. So we split a topic (one kind of note, like snack sales) into several belts called partitions. Now many belts move at once, and many helpers can read at once.
Which belt does a note go on? Each note can carry a key, like a user name. The key goes through a hash, which is a fixed recipe that always turns the same name into the same number. The remainder after dividing that number by the number of belts picks the belt. So every note about Mia lands on the same belt, and Mia's notes stay in the order they happened. Notes without a key are simply spread around to balance the load.
Order is only kept inside one belt. If Mia's note is on belt 0 and Leo's note is on belt 2, nobody promises which one is read first. Usually that is fine, because we mostly care about the order of one person's notes.
There are weak spots. If one key is extremely busy, like a famous user, its belt gets crowded while the others sit quiet. This is called a hot key. The number of belts also limits how many helpers in one team can work at once. And if you add belts later, the recipe sends keys to new places, so a key's new notes may land on a different belt than its old ones, and their order across that change is no longer guaranteed. So pick the number of belts with room to grow.
Remember
Same key, same belt, same order. Different belts, no promise about order.
Consumer groups
Teams of helpers that share the reading work
A consumer is a helper program that reads notes and does something with them. A consumer group is a team of helpers doing the same job, like the billing team. The team splits the belts between its members, and each belt is read by exactly one helper in the team at a time. That keeps each belt's notes in order, and the team does not handle the same note twice in normal running.
One helper can look after several belts. But two helpers in the same team never share a belt. So if a team has more helpers than belts, the extra helpers sit idle, like kids on the bench waiting in case a player gets tired. With 3 belts, at most 3 helpers in one team can be busy.
Each team keeps its own bookmarks. The billing team and the analytics team both read every note, each at its own speed, without getting in each other's way. How far a team's bookmark is behind the newest note is called consumer lag. Lag that keeps growing means the team cannot keep up and needs more helpers, or more belts.
When a helper joins, leaves or crashes, the team hands the belts out again. This is called rebalancing, and it can pause reading for a short time. Also, a helper usually moves its bookmark after it finishes the work. If it crashes in between, the next helper redoes those few notes. So the work should be safe to do twice.
Remember
Inside a team, one belt has one reader. Different teams each read everything, with their own bookmarks.
Stream processing
Doing math while the notes fly past
Instead of saving all the notes and counting them tomorrow, a stream processor works on them as they pass, like a kid at the cafeteria door clicking a counter for every student who walks in. It can filter (keep only the chocolate milk sales), count (how many each minute), enrich (add each student's class from a lookup list) and join (match each order with its payment).
To count, the processor must remember things between notes, like the running total for each kid. This memory is called state, and it is kept separately for each key. Since a key always lands on the same belt, one helper sees all of that key's notes and can keep its total.
What if the helper crashes? Every so often it takes a checkpoint. It saves its state and its bookmark together, at the same moment, like writing your score and your page number on the same sticky note. After a crash it reloads that pair and carries on from there, so the totals come out right. If the totals and the bookmark were saved at different times, a crash could count some notes twice or skip some.
The cost is extra work: state needs space, checkpoints take time, and anything the helper does to the outside world, like sending an email, may still happen twice after a crash unless that receiver can spot repeats.
Remember
Save the running totals and the bookmark together, so a restart picks up exactly where it left off.
Windows, event time and watermarks
Counting in time buckets, even when notes arrive late
A stream never ends, so you cannot ask how many sales there were in total. Instead you count inside windows, which are buckets of time. A tumbling window is like back to back recess periods: 12:00 to 12:05, then 12:05 to 12:10. They never overlap, so each note lands in exactly one bucket. A sliding (or hopping) window has a fixed size but starts more often, like a 5 minute bucket that starts every minute, so buckets overlap and one note can count in several. A session window has no fixed size. It groups a burst of activity and closes after a quiet gap, like a play session that ends when nobody has touched the toy for 10 minutes.
Every note has two times. Event time is when the thing really happened. Processing time is when our helper finally reads it. They can be far apart: a phone on a school bus with no signal may send its notes an hour late. If we bucket by processing time, those notes land in the wrong bucket. So careful systems bucket by event time.
But then, when is a bucket finished? A watermark is the processor's best guess, saying: we believe every note from before 3:00 has now arrived. When the watermark passes the end of a window, the window closes and its answer is sent out. A note that shows up after that is late. It can be dropped, put in a separate pile, or used to fix the answer if we agreed to wait a little longer.
This is a real trade-off. Waiting longer gives more complete answers, but slower ones. Closing early gives fast answers that may miss a few late notes.
Remember
Count in time buckets by when things happened, and let a watermark decide when a bucket is done.
Fan-out on write vs. fan-out on read
Delivering one post to many followers: copy early or gather late
When someone posts, all their followers should see it in their feed. Fan-out on write (also called push) does the work right away. The moment Mia posts, a copy (usually just the post's ID number) is dropped into every follower's ready-made feed, like a teacher slipping a flyer into every cubby. When Leo opens the app, Leo's feed is already waiting, so reading is fast and simple.
Fan-out on read (also called pull) does the work later. Mia's post is saved once. When Leo opens the app, the system visits every account Leo follows, grabs their newest posts and sorts them together, like walking to each friend's desk to see what is new. Posting is cheap, but every read does a lot of gathering, so reads are slower and cost more.
Push breaks down for stars. An account with 50 million followers would need 50 million copies for every single post, and the last copies would land long after the first. Pull breaks down for readers who follow thousands of accounts, because every refresh gathers from all of them.
So most big social apps use a hybrid. Posts from normal accounts are pushed into feeds ahead of time. Posts from the few huge accounts are not copied. They are pulled when someone reads, and merged into the ready-made feed. Most of the work is done early, and nobody has to make 50 million copies.
Remember
Push makes reads fast, pull makes posting cheap, and the hybrid pushes for most people but pulls for stars.
Quick recap
- A stream is an append-only line of notes that is kept for a while, so many readers can read and re-read it.
- Partitions split a topic into several belts. The key picks the belt, so one key's notes stay in order, but there is no order across belts.
- In a consumer group each belt has exactly one reader at a time. Extra readers sit idle, and each group keeps its own bookmarks.
- Stream processing keeps state per key and saves it together with the bookmark, so a crash does not double count.
- Windows count in time buckets. Event time and watermarks handle notes that arrive late.
- Fan-out on write is fast to read, fan-out on read is cheap to write, and big feeds mix both.
Grown-up words
and what they mean in plain words
- Event
- A small note about something that already happened.
- Topic
- One named stream for one kind of note, like snack sales.
- Partition
- One of the belts a topic is split into. Order is kept inside it.
- Offset
- A note's place number on its belt. A bookmark stores one.
- Consumer group
- A team of helpers sharing a topic's belts, each belt read by one helper.
- Consumer lag
- How far a team's bookmark is behind the newest note.
- Rebalancing
- Handing the belts out again when helpers join or leave.
- Window
- A bucket of time that notes are counted in.
- Watermark
- A marker that says notes from before this time have probably all arrived.
- Fan-out
- Spreading one post out to many followers' feeds.