System design practice

Notification Fan-out

medium queues async external-providers failover

Queue it, check preferences, fail over between providers.

Solve it in your browser Read the lesson first

The Order Service wants to tell customers when their order ships, is delayed or arrives. Build the notification system it hands these events to: it picks the channel each customer wants and sends the message through an external email, SMS or push provider.

Functional requirements

Use these use case and scenario names exactly: the traffic, requirements and tests in problem.proschi refer to them.

Scale

Constraints

What is given

problem.proschi declares the orders service (the caller), the existing prefs service (each user's channel and opt-outs) and the providers email, sms, smsBackup and push, with their rate limits. It also holds the traffic, requirements and tests. Add the components, the connections and the two use cases.

How your design is checked

You write the design as text in Proschi. Tests run in your browser: a simulation of the traffic above checks latency, availability, cost and what happens when a machine fails. How the simulation works.

Start designing

Lesson · 12 min read

Learn it: Notification Fan-out #

Notification Fan-out: accept fast, deliver patiently #

"Tell the customer their order shipped" sounds like one line of code: call the email API. Then the email provider has a bad afternoon, every order update starts taking seconds, and some of them fail. In this lesson you design a notification system that takes events from the Order Service in a few milliseconds. It then picks the right channel, respects opt-outs and survives a provider outage.

What you'll learn #

The problem, explained #

Who uses it. The Order Service, an internal caller. It emits events like OrderShipped and wants them turned into an email, SMS or push notification, whichever the customer prefers.

Functional requirements.

Non-functional requirements. Handing over an event takes under 50 ms at p99 (99% of hand-overs are faster), and the Order Service never waits for a provider. Accepting events is available 99.95% of the time. An accepted event is never lost, even if every provider is down for a while. Losing any single machine, or the SMS provider, must not stop notifications. At most $2,000 a month, including the given Order Service and Preferences service.

What is given, and why. given.proschi declares the caller (orders, four replicas), the existing Preferences service (prefs, three replicas) and four providers: email, sms, smsBackup and push. The providers carry capacity limits (2k, 500, 500 and 10k requests a second) because real providers rate-limit you. They are external systems: you do not run them, they cost nothing in the model, and they answer in about 200 ms with 99.9% availability.

What the tests check.

Back-of-the-envelope #

Peak is 1,500 events a second, each accepted once and delivered once. The channel mix splits that load: 55% email, 30% push, 11% SMS (10% sent by the primary, 1% failing over to the backup) and 4% opted out.

QuantityArithmeticResult
Queue writes (Notify)1,500 events/s1,500 rps
Preference lookupsevery delivered event1,500 rps
Email sends55% × 1,500825 rps (limit 2k)
Push sends30% × 1,500450 rps (limit 10k)
SMS sends10% × 1,500150 rps, plus 15 calls that time out (limit 500)
SMS failover: backup sends1% × 1,50015 rps (limit 500)
Opted out: nothing sent4% × 1,50060 events/s
Provider calls in flight1,500/s × ~0.2 s (Little's law)about 300 at once
Backlog if all providers are down for 10 minutes1,500/s × 600 s900,000 events

Rate limits. Every channel is under its provider's limit at peak. Email is the closest, at about 40%. That matters: if the Order Service also called a provider inline, the email provider would get that load on top and could go over its limit.

Concurrency. Little's law says the number of requests in flight equals arrival rate times time in the system. With 200 ms provider calls, workers hold about 300 calls open at once. Real workers need asynchronous I/O or a big enough pool. Proschi does not model a worker waiting on a slow dependency (a service replica is one server with a 10 ms service time), so you have to point this out yourself.

Backlog. A ten-minute outage of everything leaves under a million small messages in the queue. Queues keep messages for days (SQS for up to 14), so "never lose an accepted event" is a property you get from the queue, not from heroics in the workers.

Sizing the parts you add. A queue node takes 50k writes a second per replica, so two are lightly loaded and give redundancy. A service replica takes 2k requests a second; the worker receives 1,500 events a second. Divide by your target utilisation, then make sure the tier stays under 100% with one replica lost.

Latency in the model. Notify's path is the Order Service writing the queue, a 5 ms hop. The 50 ms limit leaves lots of room, unless a provider is on the path: one provider call is already 200 ms. Deliver has no latency requirement, and the providers decide its p99. In the model a provider is a single server, so at 41% busy the email provider takes around 340 ms on average because of queueing. That shows why you never want a third party's slow tail on your caller's path.

Availability. Notify's availability is the queue's, and two replicas of a 99.99% queue are effectively always up in the model. Deliver's availability is capped by the providers (99.9% each). That is fine, because the queue holds events until they are delivered.

Cost. The given services cost $700 a month ($100 per service replica). A queue replica is $200. The budget leaves room for a sensible design, not for doubling every tier.

Concepts #

Queue-based load levelling #

Put a durable queue between a fast caller and slow work. The caller's request ends when the queue stores the message. Workers then drain the queue at the pace the providers allow. The queue absorbs bursts (load levelling) and outages (buffering).

Why it works: the queue is a simple, highly available system whose only job is to accept and hold messages. Its availability and latency are much better than any provider's, so the caller inherits those, not the provider's.

Trade-offs: delivery becomes asynchronous, so the caller cannot know whether the email was sent, only that it will be attempted. You also need to monitor queue depth and message age. When not to use it: when the caller truly needs the result now. (A one-time password the user is waiting for may still go through a queue, but with a priority lane and tight alerting.)

title "Accept now, work later"

caller "Caller"   [REST API]        x2
inbox  "Inbox"    [AWS SQS]         x2
worker "Worker"   [Worker]          x2
slow   "Slow API" [Third Party API]

caller -> inbox  : SendMessage
inbox  -> worker : deliver
worker -> slow   : call

usecase "Accept" {
  caller -> inbox  : SEND ReportRequested {"id": 7}
  inbox --> caller : 200 queued
}

usecase "Work" {
  inbox   -> worker : ReportRequested {"id": 7}
  worker  -> slow   : POST /render
  slow   --> worker : 200
  worker --> inbox  : delete
}

Note the shape: the second use case starts at the queue. In Proschi a request to a queue always counts as a write, so a worker that "polls" with a request would look like write load on the queue and its flow would start at the worker. Real SQS consumers do long-poll; the model's convention is that the queue hands the message over, as an SQS-triggered Lambda or a push subscription would.

At-least-once delivery and idempotency keys #

Queues like SQS deliver at least once. A message stays in the queue, hidden for a visibility timeout while a worker handles it, and is deleted only when the worker says so. If the worker crashes, the message reappears and another worker takes it. Occasionally a message is delivered twice even without a crash, so the SQS documentation tells you to make consumers idempotent (safe to run twice).

For notifications, a duplicate means a customer gets the same text twice. The standard defence is an idempotency key: a unique id (the event id) that the worker passes to the provider or records itself, so a repeated send with the same key does nothing. Stripe's write-up on idempotency explains the pattern for APIs in general.

Trade-off: someone must store the keys for a while. Providers that accept an idempotency key do it for you; otherwise a small table or cache of recently sent event ids does it.

Provider failover and circuit breakers #

A third party will fail. For channels with alternatives (SMS has many vendors), the worker tries the primary and, on a timeout or error, sends through a backup. A circuit breaker makes this cheap. After enough failures, it stops calling the primary for a while and goes straight to the backup, so you do not pay a timeout on every message.

In Proschi, a failed call is -x, which costs a timeout (1,000 ms by default) and gets no answer. A scenario that calls a node with -x and still succeeds through another node is a fallback, which is what survive failure of … and handles failure of … look for:

title "Two providers"

jobs    "Address Events"  [AWS SQS]         x2
app     "App"             [Worker]          x2
primary "Geocoder"        [Third Party API]
backup  "Backup Geocoder" [Third Party API]

jobs -> app     : deliver
app  -> primary : lookup
app  -> backup  : lookup

usecase "Geocode" {
  jobs -> app : AddressAdded
  alt "Primary" {
    app      -> primary : GET /geocode
    primary --> app     : 200
  } alt "Failover" when "the primary times out" {
    app     -x primary : GET /geocode
    app     -> backup  : GET /geocode
    backup --> app     : 200
  }
}

When not to fail over: when the two providers would both act (a payment captured twice), or when the backup is much worse and a short delay would be better. Then retrying later from the queue is the safer choice.

Designing it step by step #

1. Scope. Clarify: which channels (email, SMS, push), who decides the channel (the user's preferences), whether one event can produce several notifications (no, at most one), the peak rate, and the promise to the caller ("accepted" means stored, not sent). Ask about opt-outs: for legal and ethical reasons, they must be honoured before anything goes out.

2. High-level design. Two flows, one boundary between them:

Name the alternative and reject it: the Order Service calling providers itself. It couples order processing to the worst provider's latency and availability, and the requirement says the Order Service must never wait for one.

3. Deep dive.

4. Wrap-up. Check: Notify is one queue write; no path from the Order Service to a provider; the SMS outage is covered by the failover scenario; every node you run has two or more replicas; the bill fits. Then extend: priority queues for urgent messages, per-user rate limits so nobody gets twenty texts in a minute, a dead-letter queue for events that keep failing, and delivery receipts from providers to track what actually arrived.

Common mistakes #

A notifier that queues and then sends inline (wrong/notifier-calls-provider-inline). The Order Service calls a notifier, which writes the queue and then calls the email provider before answering. The queue is there, but the caller still waits for a 200 ms third party, and every provider outage becomes an order-processing outage. In the model it is worse: the email provider now gets the inline sends on top of the worker's and goes over its 2k rps limit, so Notify's p99 runs to many seconds. It fails The Order Service never waits for a provider, and also Notify's p99 and availability limits.

A worker that polls the queue (wrong/worker-polls-queue). The Deliver flow starts with the worker sending ReceiveMessage to the queue. That is how SQS consumers are often written, but in this model it reads as the worker writing to the queue, and the flow no longer starts at the queue. It fails Workers take events from the queue. Model the queue handing the event over instead.

Other classic mistakes.

In the interview #

Say the core idea in one sentence: "The caller's promise is accepted, not sent; a durable queue makes that promise cheap and keeps third parties off the request path." Draw the queue as the boundary, then walk through one event per scenario.

Follow-ups you should expect:

Further reading #

Now design it

More system design problems