System Design Interview

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.

Each guarantee is one failure pointBefore processingThe message is lostAfter processingThe message repeatsAtomically,with outputHolds insideone systemOffset committedA crash then meansAt-most-onceAt-least-onceExactly-once
Exactly-once holds only where offset and result commit together; across a network boundary it becomes idempotency.

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.

A write, and when it is acknowledgedProducersends to leaderLeaderappends its logISRfollowerspull itacks=allwaits for ISRAck returnsto produceracks=1 returns sooner and loses data if the leader dies first.
The acknowledgement level is a durability dial, and every setting on it is a choice somebody has shipped.

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".

SettingBroker replies afterDurabilityRound-trip cost
acks=0Sending, without waitingNone — a broker crash loses the batch~0 ms added
acks=1The leader's own appendSurvives follower loss, not leader lossone round trip
acks=allEvery ISR member has itSurvives leader lossone 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.