Course Content
System Design Interview
31 sections · 71 lessons
Message Queue deep dive: partitions, ordering and storage on disk
The first lesson settled on a retained, replayable log taking 100,000 messages a second at 1 KB each. One log cannot absorb 100 MB/s, because one log means one machine's disk and one writer. Partitioning is how a topic becomes many logs, and it is where ordering gets decided.
The second half of this lesson looks inside a single partition: how the bytes sit on disk, and why a system doing 100 MB/s does not need to keep its data in memory.
How a message picks a partition
The producer computes a partition key from the message, hashes it, and takes the result modulo the partition count. Same key, same partition, therefore same log, therefore ordered.
With key = order_id, every event for order 4471 — created, paid, shipped, delivered — lands in one partition and arrives in sequence. Events for different orders may interleave arbitrarily, and nobody cares.
With no key, the producer round-robins, which balances load perfectly and provides no ordering at all.
Choosing the partition count
Two forces pull in opposite directions.
More partitions means more consumer parallelism and more brokers sharing the load. At 100,000 messages per second and a consumer that handles 1,000 per second, you need at least 100 consumers, therefore at least 100 partitions.
Fewer partitions means fewer open files, fewer replication streams, and faster leader election when a broker dies. A cluster with 200,000 partitions takes noticeably longer to recover than one with 2,000, because the controller has to reassign every leader.
Start at roughly two to four partitions per broker per topic and revise. For twenty brokers, 40–80 partitions is a sensible opening answer; the consumer arithmetic above pushes it to 100. Increasing partition count later is possible but breaks key-to-partition stability — hash(key) % 100 and hash(key) % 200 send the same key to different places, so ordering is violated across the change. Say this unprompted.
The hot partition
Partition by country_code and one country holds 40% of traffic. That partition's broker receives 40 MB/s while others receive 3 MB/s. Consumer lag grows on one partition only.
Fixes, in order of preference: pick a higher-cardinality key such as user ID; add a random suffix to the hot key and accept losing ordering for that key alone; or give the hot key its own dedicated topic. This is the same hot-key problem as the ring imbalance in Section 7 (Design Consistent Hashing) and the celebrity fan-out in Section 13 (Design a News Feed System) — it recurs in this course more than any other issue.
Storage on disk
Now inside one partition. The instinct is that a system doing 100 MB/s must keep data in memory, because disks are slow. That instinct is wrong here, and the reason is worth understanding properly.
Sequential disk writes are not slow
A 7,200 RPM spinning disk manages roughly 100–150 random input/output operations per second. At 4 KB per operation that is about 0.5 MB/s. Genuinely slow.
The same disk streaming sequentially, with the head not moving, sustains roughly 150 MB/s — around 300 times more. The disk is not slow; seeking is slow. Solid-state drives shrink the gap but do not close it, because sequential writes also avoid write amplification inside the drive's own translation layer.
A message log has the one access pattern that is purely sequential: append at the end, read forwards. So the design is: never seek on the write path.
Segments, indexes, and retention
Each partition is a directory of segment files, each capped at, say, 1 GB. Writes go to the active segment; when it fills, it is closed and a new one opens. Closed segments are immutable, which makes them cheap to replicate, cache, and delete.
Alongside each segment sits a sparse index mapping offsets to byte positions — one entry every few thousand messages, not one per message. To find offset 4,096,881: binary-search the index to the nearest earlier entry, seek once, then scan forward a few kilobytes.
Retention deletes whole segments. Deleting a 1 GB file is one filesystem operation; deleting a million individual messages is a million. This is why "seven days" is cheap and "delete this one record" is expensive.
The page cache does the caching
The broker deliberately does not maintain its own message cache. It writes to the operating system's file cache and lets the kernel decide what stays in memory.
Two benefits. First, a consumer reading recent messages — which nearly all consumers do — is served from memory without the broker knowing or caring. Second, a broker restart does not cold-start the cache, because the page cache belongs to the kernel, not the process.
On a broker with 64 GB of RAM and 10 TB of disk, the cache holds roughly the last 100 minutes of that broker's traffic at 10 MB/s per broker. Consumers within 100 minutes of the head read at memory speed; a consumer replaying from three days ago reads from disk and runs slower. That is the observable, explainable behaviour to describe in the interview.
Zero-copy sends
Serving a message normally means: disk → kernel buffer → application memory → socket buffer → network. Four copies and two context switches.
The sendfile system call collapses this to disk → kernel buffer → network interface. The data never enters the broker's own memory. For a broker pushing 300 MB/s of consumer traffic, removing two copies per byte is a large fraction of its CPU budget.