System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Notification System: reliable delivery, preferences and scaling


The design from Notification System: scope, scale and high-level design gets messages into a per-channel log and out to providers. This lesson makes it correct and then makes it survive scale: first reliability, then the preference and throttling rules that decide whether a message should be sent at all, then what breaks first when a broadcast hits.

The reliability deep dive is the one genuinely new idea this section owns: reliable fan-out with deduplication. Later sections in Part III reuse fan-out and caching; none of them re-derive this.

Reliability: exactly-once as a goal

At-least-once, plus a dedup keyEvent with aunique keyCheck thededup storeEnqueueper channelWorkersends, then acksFive fails:dead letterA worker that dies after sending but before acking still sends twice.
Exactly-once delivery is a lie; exactly-once effect is achievable only if the receiver keys on your identifier.

The honest position on exactly-once

Exactly-once delivery, as a network property, does not exist across a third party. What you can build is at-least-once delivery plus deduplication, which produces exactly-once as the user experiences it. Say that sentence in the interview. Candidates who claim exactly-once delivery outright get probed until they concede it.

Where duplicates actually come from

Three places, and they are worth enumerating because the fix has to cover all three.

  1. The caller retries. The order service posts a notification, the response times out, it does not know whether we accepted it, so it posts again.
  2. The worker crashes mid-flight. The worker calls the push provider, the provider accepts, the worker dies before committing its offset in the log. Another worker picks up the same message and sends it again. This is the classic at-least-once consumer, covered in Message queues and event streaming.
  3. Failover between providers. The primary SMS provider times out, the worker fails over to a secondary, and the primary had actually delivered.

The mechanism

The dedupe_key is supplied by the caller and describes the event, not the attempt: order-8891-shipped. Two posts of the same real-world event carry the same key.

At the last possible moment before the provider call, the worker performs an atomic set-if-absent against a fast store — a Redis SET key value NX EX 259200 gives a three-day window:

  • Key absent → the set succeeds, this worker owns the send, call the provider.
  • Key present → another attempt already owns it, acknowledge and drop.

Atomicity matters. A read-then-write pair lets two workers both see "absent" in the gap between the read and the write, and both send. Making it distributed, in Section 6, makes the same point about distributed rate limiting, and the fix is the same: one atomic operation, not two.

The hole in that plan, stated honestly

Claiming the key before the provider call means that if the provider call then fails permanently, the key is consumed and the retry is suppressed — a lost notification. Claiming it after means a crash between send and claim produces a duplicate.

You must pick which error you prefer, per category:

CategoryPreferenceClaim timing
Two-factor codeNever lose it; a duplicate is tolerableClaim after send
Marketing pushNever duplicate; a loss is tolerableClaim before send

The better answer, where the provider supports it, is to hand the provider your own message identifier and let the provider deduplicate. Several email and push providers accept a client-supplied idempotency key. When one does, use it — the problem moves to where it can actually be solved.

Retry, backoff, and the dead-letter queue

Retry with exponential backoff and jitter (Reliability patterns): 1 s, 2 s, 4 s, 8 s, each multiplied by a random factor between 0.5 and 1.5. Without jitter, ten thousand workers that failed at the same instant retry at the same instant and reproduce the outage.

Distinguish retryable from terminal failures. A 500 or a timeout is retryable. "Invalid device token" is terminal — retrying it forever is a slow leak that eventually consumes a worker pool. Terminal failures delete the token; exhausted retries go to a dead-letter queue, a separate topic holding messages a human or a repair job must look at.

Preferences and throttling

Delivery reliability makes messages arrive. Preferences and throttling decide whether they should have been sent at all — and this is where notification systems actually lose users.

Where the preference check belongsEvent producedPreference check hereThrottle, digest windowChannel worker sends
Checking before the queue makes an opted-out user cost one lookup instead of a whole delivery attempt.

Where the preference check belongs

The tempting place is at the API, before enqueueing. That is wrong, and the reason is a timing argument.

A marketing broadcast enqueues 20 million messages. At the provider's 10,000/s ceiling, the tail of that queue is drained 2,000 seconds — about 33 minutes — after the head. A user who unsubscribes at minute five has their unsubscribe honoured only if the check happens when the message is sent, not when it was queued.

So: check at the API to fail fast and avoid enqueueing obvious waste, and check again in the worker, immediately before the provider call. Two checks, and the second one is the one that counts.

The preference model

Three dimensions: user, category, channel.

usercategorychannelallowed
u_4412order_updatespushyes
u_4412order_updatesemailyes
u_4412promotionspushno
u_4412promotionsemailyes

A global "pause everything" switch sits on top. Transactional categories — password reset, security alert, two-factor code — are usually not opt-outable, and that is a product decision you should state rather than assume.

At 33,000 sends/s during a broadcast, this table is read 33,000 times a second. It is small (100M users × ~20 rows × ~40 bytes ≈ 80 GB, and far less in practice because most users never change a default), changes rarely, and is read constantly — a textbook cache-aside case (Caching strategies), with the cache invalidated on preference update.

Digesting a burst

A post gets 5,000 likes in ten minutes. Without intervention that is 5,000 push messages to one person, and that person uninstalls the app.

The fix is a per-user, per-category collection window. The first event schedules a delivery for now plus 5 minutes and opens a counter. Subsequent events in the window increment the counter instead of sending. At the end of the window, one message: "Your post got 5,000 likes."

The window is a product trade-off, not an engineering one. Five minutes on a social feed; zero on a direct message.

Quiet hours and per-user rate limits

Quiet hours need the user's timezone, not the server's. A message generated at 03:00 in the user's local time is either held until 08:00 or dropped, depending on whether it is still useful in the morning. Transactional messages ignore quiet hours.

Per-user rate limits cap the total: at most three promotional pushes a day, at most one an hour. A token bucket per user and category does this in one operation — the same algorithm built in The five algorithms, compared.

Scaling: what breaks first

Everything above works at steady state. Here is what breaks first, and the questions this problem always attracts.

One API call, twenty million devicesOnecampaign requestExpand in abatch jobShard by user IDPriority lanesRate-limitedto providerA transactional alert must not queue behind a marketing blast.
The fan-out has to be a background job, because no synchronous request can enumerate twenty million rows.

Fan-out to 20 million on one API call

A campaign service posts one request targeting an audience of 20 million. If the notification API tries to expand that synchronously, the request runs for minutes, holds a connection, and loses everything if the process restarts.

Two-stage fan-out instead:

  1. The API writes a single campaign record and returns 202 immediately.
  2. A fan-out coordinator splits the audience into batches by user ID range — say 1,000 users per batch, giving 20,000 batch tasks — and writes those tasks to a queue.
  3. Fan-out workers consume a batch, resolve those 1,000 users to devices, check preferences, and write individual messages to the channel topic.

The work is now restartable at batch granularity. A worker that dies re-processes 1,000 users, not 20 million, and the deduplication key from the reliability section above makes that re-processing safe.

Priority lanes

Recall the arithmetic: 20 million marketing messages at 10,000/s is 2,000 seconds of backlog. A password-reset message that lands behind them arrives 33 minutes late, and the user has already given up.

The fix is not a priority field on a message — a distributed log delivers in order and does not reorder for you. The fix is separate topics and separate worker pools:

LaneTopicWorkersProvider quota share
Transactionalpush.transactional4030% reserved
Standardpush.standard6050%
Bulkpush.bulk2020%, and it may be paused

Reserving quota is the part candidates miss. Separate queues without a reserved share of the provider's rate limit still starve the transactional lane, because both lanes contend for the same 10,000/s.

Templating and localisation

A template identifier plus parameters plus a locale, rendered at send time in the worker. Rendering late means a template fix reaches queued messages, and the message is stored as ~100 bytes of parameters rather than a full rendered body. Fall back to a default locale when a translation is missing — and log the miss, because a silent fallback to English for a market is a bug nobody reports.

Delivery tracking

Providers report status asynchronously through webhooks: delivered, bounced, rejected. A webhook receiver writes those into the delivery log, keyed by your message identifier.

Be honest about the limits: for mobile push, "the provider accepted it" is the strongest signal you get. Whether the device displayed it, and whether a human saw it, is not reported by the push service. Open tracking requires the app to report back, and that report is missing for anyone who never opens the app — which is exactly the population you wanted to measure.