System Design Interview

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.

One topic, three partitions, two consumer groupsTOPIC: ordersP0P1P2order within a partition is guaranteed; across partitions it is notpartitions are spread across brokersBroker 1Broker 2Broker 3GROUP: billingconsumer 1consumer 2consumer 3one consumer per partition — themaximum useful parallelismGROUP: analyticsconsumer 1consumer 2a second group reads the same recordsindependently, at its own offsetPartition count caps parallelism within a group, and the partition key decides what stays ordered. Both are chosen once and are painful to change.
Adding a fourth consumer to the billing group buys nothing — three partitions is three consumers' worth of parallelism.

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.

Why the disk is not the bottleneckRandom access• A seek costs roughly 10 ms• Hundreds of operations per second• This is the instinct that misleadsSequential append• Hundreds of megabytes a second• The page cache absorbs the reads• Zero-copy sends skip user space
The log is append-only precisely so the disk is used in the one access pattern at which it is genuinely fast.

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.