Course Content
System Design Interview
31 sections · 71 lessons
Message Queue: delivery guarantees, replication and failure modes
The design so far has partitions, segment files and consumer offsets. This lesson is about what happens when something fails: a consumer crashes mid-message, a network drops an acknowledgement, a broker dies. These are the questions interviewers spend the last fifteen minutes on.
Start with delivery semantics. Three phrases get used loosely. Each one is defined by what happens at one specific failure point, so define them that way.
The three guarantees, by failure
Consider a consumer that reads a message, processes it — charges a card, sends an email — and then commits its offset.
At-most-once: commit the offset first, then process. If the consumer crashes between the two steps, the offset has moved and the message is never processed. Zero duplicates, occasional loss. Right for high-volume metrics where one lost sample is invisible.
At-least-once: process first, then commit. If the consumer crashes after processing but before committing, restart re-reads the message and processes it again. Zero loss, occasional duplicates. This is the default and the right default.
Exactly-once: every message affects the world once. There is no failure point at which it is lost or repeated.
Why exactly-once is mostly a lie
The producer sends a message and the network drops the acknowledgement. The producer does not know whether the broker wrote it. It can retry, risking a duplicate, or not retry, risking a loss. There is no third option — this is the two-generals problem, and it is a proof, not an engineering gap.
Brokers narrow it with a producer sequence number: each producer gets an ID, numbers its messages, and the broker rejects a number it has already seen. That removes duplicates from the retry path into the log.
Some systems then add transactions that atomically commit a consumer offset and a set of output messages, giving exactly-once within the system — read from topic A, write to topic B, commit offset, all or nothing.
What no broker can give you is exactly-once out of the system. If the consumer's side effect is charging a card at an external processor, the broker has no way to make that external call atomic with the offset commit.
Offset commit timing, concretely
- Auto-commit every 5 seconds is the convenient default and gives you at-least-once with a window of up to 5 seconds of replay on crash. At 1,000 messages per second per partition, that is up to 5,000 duplicates per crash.
- Commit after each message shrinks the window to one message and adds a network round trip per message, capping throughput at roughly the inverse of the round-trip time — a 1 ms round trip caps you near 1,000 messages per second.
- Commit in batches of N after processing is the usual compromise: batch of 500 at 1,000 messages per second means committing twice a second and replaying at most 500 on crash.
Pick the batch size from how expensive a duplicate is, and say so.
Replication and coordination
Everything so far assumed brokers do not die. They do.
Leader, followers, and the in-sync set
Each partition has one leader replica and some number of followers. All reads and writes go to the leader. Followers continuously fetch from the leader and append what they receive.
A follower that is caught up — within a bounded lag, typically a few seconds — is in the in-sync replica set, usually shortened to ISR. A follower that falls behind is dropped from the ISR and stops counting towards durability until it catches up.
When a leader dies, a controller elects a new leader from the ISR. Because every ISR member has all acknowledged messages, no acknowledged message is lost.
Acknowledgement levels: the durability dial
The producer chooses when the broker replies "written".
| Setting | Broker replies after | Durability | Round-trip cost |
|---|---|---|---|
acks=0 | Sending, without waiting | None — a broker crash loses the batch | ~0 ms added |
acks=1 | The leader's own append | Survives follower loss, not leader loss | one round trip |
acks=all | Every ISR member has it | Survives leader loss | one round trip plus slowest follower |
acks=all alone is not sufficient. If the ISR has shrunk to one member — the leader — then "all" means "the leader", and you have acks=1 while believing otherwise. Pair it with a minimum in-sync replicas setting of 2, which makes the broker reject writes rather than accept them unsafely. With replication factor 3 and minimum ISR 2, you tolerate one broker failure with no loss and no write outage, and refuse writes after two — availability sacrificed for durability, deliberately.
Recommendation: acks=all with replication factor 3 and minimum ISR 2 for anything that matters; acks=1 for logs and metrics where a rare loss is acceptable and the extra millisecond is not.
Consumer groups and rebalancing
A consumer group is a set of consumers sharing a group ID, between which partitions are divided. One partition is assigned to exactly one consumer in the group. When a consumer joins or leaves, a rebalance reassigns partitions.
The naive rebalance protocol stops every consumer in the group, recomputes the assignment, and restarts them. For a group of 100 consumers, one restarting pod pauses all 100 for the duration — often several seconds. Cooperative rebalancing revokes only the partitions that actually move, which is the modern default and worth naming.
Backpressure when consumers fall behind
The broker does not push; consumers pull at their own rate. So a slow consumer does not overwhelm anything — it accumulates lag, the offset distance between the head of the partition and the consumer's position.
Lag is the metric to alert on. Lag growing linearly means the consumer is permanently too slow and needs more instances, which needs more partitions. Lag that exceeds the retention window means messages expire unread, which is silent data loss — alert well before that point, and say what "well before" means: if retention is seven days and lag is one day and climbing by six hours a day, you have eight days.