The difficult part isn't sending an email. It's making sure a slow provider, a viral campaign, or a retry storm cannot take down the services that need to send one.
Imagine an e-commerce platform just completed an order. The order service needs to send a confirmation email, the payment service needs to send a receipt, and a security service may need to send an urgent login alert.
Now imagine the email provider starts taking eight seconds to respond.
If those services call the provider synchronously, an external notification dependency suddenly becomes part of the checkout and login latency budget. If a promotional campaign simultaneously produces millions of messages, critical security alerts can end up behind a mountain of marketing traffic.
A production notification system has to solve a more interesting problem than sending messages: accept notification intent durably, then deliver it independently at the rate each downstream provider can sustain.
Let's build that architecture from first principles, and examine where seemingly sensible implementations fail.
1. Requirements: What Are We Designing?
Our service accepts requests from many producer applications and delivers notifications through push, email, and SMS.
It should support:
- Transactional messages such as order confirmations and payment updates.
- Time-sensitive messages such as security alerts and one-time passcodes.
- Promotional campaigns, potentially targeting millions of recipients.
- Recipient preferences, opt-outs, quiet hours, and valid destination checks.
- Delivery status, retries, and investigation of failed messages.
The important non-functional requirements are fast producer acknowledgment, durable processing, channel isolation, preferential latency for urgent traffic, and observability.
One distinction defines our API contract:
ACCEPTED does not mean DELIVERED. It means the system has durably recorded the notification for asynchronous processing. Delivery may later succeed, fail permanently, expire, or be suppressed by policy.
That contract is the foundation of the Notification System Complete Design.
2. Capacity Estimation: Ingestion and Delivery Are Different Rates
Consider an illustrative traffic burst:
| Metric | Example workload |
|---|---|
| Incoming notifications | 5,000/second |
| Provider's sustainable send capacity | 800/second |
| Excess arriving during the burst | 4,200/second |
| Backlog after a 60-second burst | 252,000 notifications |
If incoming traffic drops to zero and the provider continues delivering at 800/second, that backlog alone takes 315 seconds, or about 5 minutes 15 seconds, to drain.
In reality, new traffic may continue arriving, reducing the spare capacity available to clear old messages.
This leads to a rule that is easy to overlook:
A queue absorbs a temporary difference between arrival and delivery rates. It does not manufacture provider capacity.
If 5,000 messages arrive every second indefinitely and the provider can only deliver 800, the backlog grows indefinitely. We eventually need more provider capacity, slower campaign ingestion, traffic shedding, or an explicitly longer delivery window for less urgent notifications.
3. Start With the Simplest Architecture — and Watch It Break
The first implementation might look like this:
Order Service
|
v
Notification API
|
v
Email / SMS / Push Provider
|
v
Response to Order Service
It's simple, but the calling service is now waiting for a provider it does not control.
If the provider slows to an eight-second p99 response time, the producer's request may inherit that delay. If the provider is unavailable, notifications may fail alongside otherwise healthy business operations.
The first meaningful architectural change is not adding ten microservices. It's putting a durable asynchronous boundary between notification creation and delivery.
Producer
|
v
Notification API
|
v
Durable Queue -----> Delivery Workers -----> Provider
^ |
| v
Fast acceptance Retry / failure handling
Now the producer does not wait for the provider. Workers can process messages at a controlled pace, and a temporary provider outage becomes a backlog rather than an immediate failure of every calling application.
But the phrase “durable queue” hides a subtle consistency problem.
4. Durable Ingestion: The Database–Queue Dual-Write Problem
A typical notification API must record the notification's state in a database and publish work to a message broker.
Suppose we implement:
- Insert notification into the database.
- Publish the notification to Kafka, SQS, or RabbitMQ.
- Return success.
What if step 1 succeeds but step 2 fails?
The database says the notification exists, but no worker ever receives it.
Reversing the order doesn't solve it. The broker may accept the message while the database write fails, leaving workers processing a notification with no durable tracking record.
Transactional outbox
A practical solution is to write the notification record and an outbox event in the same database transaction:
Notification API
|
v
Single DB transaction
+------------------------+
| Notification record |
| Outbox event |
+------------------------+
|
v
Outbox Publisher ---> Message Broker ---> Workers
The API can acknowledge once that durable transaction commits, while the outbox publisher continues retrying publication if the broker is temporarily unavailable.
This closes the dual-write gap, but does not guarantee exactly-once processing. A publisher can send an event successfully and crash before marking it published, causing it to publish again after restart. Consumers still need idempotent handling.
The exact failure windows and guarantees are covered in Notification System Architecture Deep Dive — Durable Ingestion.
5. API Design: Accept Intent, Not Delivery Promises
A minimal API might expose:
POST /notifications
Content-Type: application/json
Idempotency-Key: order-8421-confirmation
{
"recipient_id": "user-123",
"channel": "email",
"priority": "normal",
"template_id": "order-confirmed",
"template_variables": {
"order_id": "8421"
}
}
A successful response returns a notification ID and ACCEPTED status — not DELIVERED.
{
"notification_id": "notif-9271",
"status": "ACCEPTED"
}
Other useful endpoints include:
-
GET /notifications/{id}/status— inspect the lifecycle state. -
POST /notifications/batch— accept a campaign definition and return a batch ID without synchronously expanding millions of recipients.
The producer-supplied idempotency key matters. If the producer times out after the API successfully commits, it may retry. The same key should return the existing logical notification instead of creating another one.
6. High-Level Architecture: Isolate Channels and Priorities
A single queue with one pool of workers looks attractive until the email provider fails while push notifications are healthy.
We want the channels to fail and scale independently.
Producer Services
|
v
Notification API
|
Durable Ingestion
|
v
Policy / Channel Router
|
+-----------+-----------+
| | |
Push Queue Email Queue SMS Queue
| | |
Workers Workers Workers
| | |
Rate Limiter Rate Limiter Rate Limiter
| | |
FCM/APNs Email Provider SMS Provider
\ | /
\ | /
Delivery State & Reconciliation
Each channel can have its own backlog, retry policy, worker concurrency, provider adapter, and rate-limit budget.
The architecture also needs priority-aware scheduling within or across channel queues. We'll address that next.
7. Priority Queues: Security Alerts Must Not Wait Behind Campaigns
Imagine a promotional campaign enqueues 10 million notifications. A security alert arrives one second later.
If everything shares one FIFO queue, that alert could wait behind enormous amounts of lower-value traffic.
Separate priority classes help:
| Priority | Typical traffic | Desired behavior |
|---|---|---|
| High | Security alerts, payment failures | Preferential low latency |
| Normal | Order confirmations, shipping updates | Predictable progress |
| Low | Promotions, newsletters | Can tolerate delay |
But a naive rule — always drain high priority before touching normal or low — creates starvation. If high-priority traffic never fully stops, lower-priority notifications may never be served.
Production alternatives include weighted polling, reserved capacity, and aging.
For example, a scheduler might use an illustrative ratio of eight high-priority messages, two normal, and one low per round. The exact weights are workload decisions, not universal defaults.
The objective is to protect urgent traffic while guaranteeing that every eligible class makes progress.
See Trade-offs & Decisions — Strict vs Fair Priority for the alternatives and their costs.
8. Recipient Preferences, Push Tokens, and Templates
Before a notification reaches a provider, the router must decide whether it is eligible to send.
That means checking the recipient's category preferences, channel permissions, opt-out status, quiet hours, and whether the destination is valid.
A promotional email should not be sent simply because it was queued before the user unsubscribed. For delayed messages, cached preferences may need to be checked again close to delivery time.
Push adds another complication: users have devices, and devices have tokens. One user may have several active tokens. A provider-reported invalid token should deactivate that specific destination, not trigger endless retries or disable the user's other devices.
Templates also need versioning. If an order confirmation sits in a queue for an hour and the template changes, should the message use the original or the new content? Referencing a specific template version at creation time avoids silently changing in-flight notifications.
These concerns are explored separately in the Architecture Deep Dive, including recipient eligibility, token lifecycle, localization, and render-before-queue versus render-at-delivery trade-offs.
9. Rate Limiting and Backpressure: More Workers Won't Fix a Provider Quota
Suppose workers are capable of making 5,000 provider requests per second, but the provider account allows only 1,000.
Adding workers does not raise that account limit. It may simply generate more 429 Too Many Requests responses.
The system needs provider-level rate limiting shared across distributed workers.
Two common designs are:
- Centralized token bucket: All workers coordinate through a shared rate limiter. This provides a global quota view but adds network calls and contention.
- Leased quota: Workers receive bounded slices of the total quota and spend them locally. This reduces coordination overhead but can leave unused quota stranded with idle workers.
When a provider returns 429 or Retry-After, the limiter should adjust rather than repeatedly sending at the same rejected rate.
The queue provides backpressure absorption while workers respect the provider's sustainable throughput. It does not remove the need for admission control when traffic stays above capacity.
For a detailed comparison, see Notification System Trade-offs & Decisions.
10. Retries: Why Exponential Backoff Needs Jitter and Deadlines
A provider outage often triggers a second incident caused by the notification system itself.
Picture thousands of workers receiving failures at the same moment. If they all retry one second later, then two seconds later, then four seconds later, the recovering provider gets hit by synchronized waves of traffic.
This is a retry storm.
Exponential backoff reduces the frequency of retries. Jitter randomizes retry timing so workers don't all return simultaneously.
Without jitter:
Workers fail together
| 1s
+----> retry spike
| 2s
+----> retry spike
With jitter:
Workers fail together
| randomized delays
+--> --> ---> ----> spread-out retries
Not every failure should retry:
| Failure | Action |
|---|---|
| Temporary timeout or provider 5xx | Retry with backoff and jitter |
| Provider 429 | Honor retry guidance; lower send rate |
| Invalid push token | Deactivate token; do not retry |
| Malformed destination | Permanent failure or appropriate terminal handling |
| Maximum retries exhausted | Move to DLQ when appropriate |
| Notification deadline passed | Expire instead of delivering stale content |
Retries should be scheduled using delayed queues or retry topics, not by keeping worker threads asleep.
And every time-sensitive notification needs a deadline. A one-time passcode delivered six hours late is not a successful user experience, even if the provider eventually accepts it.
11. Dead-Letter Queues: Failure Must Be Visible and Recoverable
After a bounded number of attempts, unresolved failures can be moved to a dead-letter queue (DLQ) for diagnosis.
The DLQ is not a trash can. It is an operational tool that preserves context such as the notification ID, error, attempt count, and relevant timestamps.
A growing DLQ should be monitored by channel, provider, error category, and message age, not only by one global threshold.
Replaying messages requires care:
- Fix the underlying failure first.
- Re-check expiry and recipient preferences.
- Replay at a controlled pace.
- Retain prior attempt history and observe whether the problem recurs.
Dumping the entire DLQ back into the main queue can recreate the original overload.
See Failure & Scale Scenarios — DLQ Growth and Backlog Recovery.
12. Idempotency: Can We Guarantee Exactly-Once Notifications?
No — not universally across an external notification provider.
There are three different boundaries to protect:
- Producer → API: A retried API call should not create a second logical notification. Use a producer idempotency key.
- Queue → Worker: Brokers may redeliver messages. Make worker state transitions idempotent.
- Worker → Provider: The external side effect can become ambiguous after a timeout or crash.
Consider this sequence:
Worker sends notification to provider
|
v
Provider accepts request
|
v
Worker crashes before recording success
|
v
Broker redelivers the message
|
v
New worker cannot tell whether it was sent
If the new worker retries, the user might receive the notification twice. If it doesn't, the notification might never arrive.
A pre-send deduplication check cannot resolve this ambiguity because the successful external action was never recorded.
Where supported, provider-side idempotency keys reduce the risk. Persisting provider message IDs and reconciling against callbacks or status APIs also helps.
The honest contract is durable at-least-once processing with strong duplicate reduction, not an unsupported promise of exactly-once end-user delivery.
The full crash-window analysis is in Notification System Architecture Deep Dive — Idempotency & Duplicate Suppression.
13. Provider Failover: A Timeout Is Not a Rejection
Suppose the primary SMS provider times out. Should we immediately send through a backup provider?
Not necessarily.
The primary may have already accepted the request but lost the response. Sending through a second provider can produce two real messages.
A known rejection is different from an ambiguous outcome. Failover policy should consider the signal, message urgency, duplicate tolerance, and the backup provider's capacity.
For a security alert, accepting some duplicate risk might be reasonable to improve the chance of timely delivery. For a promotional notification, conservative retry or reconciliation may be preferable.
Provider failover is a trade-off, not a universal on/off switch. The Trade-offs & Decisions page compares failover versus retrying the primary in detail.
14. Delivery Tracking: SENT Is Not DELIVERED
The system should maintain an observable notification lifecycle:
ACCEPTED → QUEUED → PROCESSING → PROVIDER_ACCEPTED
|
v
DELIVERED
Alternative outcomes:
RETRY_SCHEDULED | FAILED_PERMANENT | EXPIRED
SUPPRESSED | BOUNCED | UNDELIVERABLE | DLQ
PROVIDER_ACCEPTED means the provider took the request. It does not prove that the recipient received it.
Providers may send delivery-status webhooks later, but those callbacks can arrive duplicated, delayed, out of order, or not at all.
A robust reconciliation pipeline should:
- Authenticate provider callbacks.
- Deduplicate callback events.
- Apply idempotent, monotonic state transitions.
- Prevent a late
SENTcallback from overwritingDELIVERED. - Query provider status APIs when callbacks are missing and such APIs are available.
Not every channel provides a reliable end-device delivery receipt. The status model must describe what the system actually knows, rather than claim more than a provider can prove.
15. Production Failure Scenarios: Where the Architecture Gets Tested
A notification system's quality becomes clear when something breaks.
Email provider outage
Email workers back off and accumulate a durable backlog. Push and SMS continue independently. Recovery is paced according to provider capacity and message deadlines.
Broker outage
With a transactional outbox, the API can continue durably recording notification intent while the broker is down, subject to database capacity and the API's acceptance contract. Publishing resumes when the broker recovers. Without a durable fallback, the API should not falsely report successful acceptance.
Massive promotional campaign
A campaign targeting millions of recipients should be expanded into bounded fanout chunks, paced to sustainable capacity. Transactional and security traffic must retain protected capacity. The campaign should support cancellation and expiry.
Priority starvation
A strict high-first scheduler can leave normal and low-priority notifications waiting forever. Monitor oldest-message age by priority, not just queue depth, and use fair scheduling or reserved capacity.
Delivery tracker outage
A well-isolated tracker failure should degrade status visibility rather than stop all delivery. But provider IDs and state events must be preserved for later reconciliation; simply losing them creates a permanent audit gap.
Backlog recovery storm
After a long outage, releasing every queued notification at once can overwhelm a recovering provider. Drain gradually, respect priority, and discard or expire messages that are no longer useful.
These are part of the 14 concrete incidents covered in Notification System Failure & Scale Scenarios.
16. How Would We Evolve This Architecture?
A good system design grows in response to evidence, not because a diagram looks more impressive with more components.
| Stage | Architecture | Problem it solves |
|---|---|---|
| 1 | Simple notification service | Basic low-volume delivery |
| 2 | Durable ingestion and asynchronous workers | Provider latency and failure isolation |
| 3 | Per-channel queues, preferences, priority scheduling | Multi-channel correctness and urgent-message latency |
| 4 | Shared provider rate limits, retry scheduling, DLQ | Throttling, transient failures, and observability |
| 5 | Idempotency, reconciliation, campaign pacing | Duplicate reduction, trustworthy state, and large-scale fanout |
The precise order depends on product requirements. For example, recipient opt-outs and basic idempotency may be mandatory from the beginning rather than later enhancements.
17. Notification System Design Interview Questions
Why use a queue in a notification system?
To decouple the rate of notification creation from the rate of provider delivery, keep producer latency independent of provider health, and buffer temporary bursts durably.
Can Kafka or SQS guarantee exactly-once notification delivery?
No. Internal broker semantics cannot eliminate the ambiguous external provider-call window. Use at-least-once processing, idempotent internal transitions, provider idempotency where available, and reconciliation.
How do security alerts bypass promotional notifications?
Use priority classes with protected capacity or fair scheduling, so urgent traffic gets preferential latency without starving other classes indefinitely.
What happens if the provider allows only 1,000 messages per second?
Coordinate a rate-limit budget across workers. Additional workers cannot bypass the provider's account-level ceiling. Pace incoming campaigns and monitor backlog age.
Why do retries need jitter?
Exponential backoff alone can leave workers synchronized. Jitter spreads attempts across time and reduces retry spikes during recovery.
Should we fail over to a second provider after a timeout?
Only after considering whether the outcome is ambiguous, how urgent the notification is, and whether duplicate delivery is acceptable.
How do we send 100 million promotional notifications without delaying OTPs?
Accept the campaign quickly, enumerate recipients asynchronously in bounded chunks, pace fanout, and reserve downstream capacity for time-sensitive traffic.
What is the difference between SENT and DELIVERED?
SENT or PROVIDER_ACCEPTED means the external provider accepted the request; DELIVERED requires a stronger provider or downstream confirmation, where supported.
For more advanced follow-ups on outbox consistency, webhook ordering, push-token invalidation, and dedup-store failures, see the Notification System Interview Perspective.
Final Thoughts
The hardest part of a notification system isn't choosing Kafka over RabbitMQ or integrating an SMS SDK.
It's recognizing that notification creation, delivery, provider acceptance, and confirmed receipt are different events with different guarantees.
A durable queue protects producers from provider slowness. Channel isolation contains failures. Fair scheduling protects critical traffic. Rate limiting respects real provider ceilings. Backoff, jitter, deadlines, and DLQs make recovery controlled rather than chaotic. Idempotency and reconciliation reduce duplicate risk without pretending external side effects can be made universally exactly-once.
The architecture becomes easier to reason about when each component exists to solve a specific failure mode.
That's the deeper lesson of notification system design: reliability comes from making failure boundaries explicit, not from adding more infrastructure.
References & Further Reading
These SeeItFlow resources explore the same system in more detail:
- Notification System — Complete Design — Requirements, API contract, queue-based evolution, channel routing, priorities, retries, DLQ, and final architecture.
- Notification System — Architecture Deep Dive — Transactional outbox, three idempotency boundaries, fair scheduling, recipient policy, push tokens, templates, distributed rate limiting, provider adapters, webhooks, campaign fanout, and delivery-state storage.
- Notification System — Trade-offs & Decisions — Sync versus async, queue topology, retry policies, rate-limit coordination, provider failover, and template-rendering decisions.
- Notification System — Failure & Scale Scenarios — Fourteen operational scenarios covering provider outages, retry storms, queue failures, throttling, duplicate events, starvation, and backlog recovery.
- Notification System — Interview Perspective — Interview answer structure, common mistakes, and advanced follow-up questions.
Explore more free, visual explanations of distributed systems, backend engineering, and system design at SeeItFlow.




Top comments (2)
The part I'd underline is that retries are a rate-limit problem, not just a scheduling problem. EDI learned this the hard way: when a trading partner's endpoint flaps, every queued document retries on the same timer and you DDoS the exact partner you're trying to heal. The fix that worked for us was a per-partner retry budget — retries consume tokens from the same bucket as first attempts, so a backoff actually buys recovery instead of just deferring the pileup. Your queue design already has the shape for that; the token bucket just has to count retries too.
Thanks for sharing this real-world EDI example! It's a good illustration of how retries can amplify downstream instability.
Retry handling has several layers — classifying retryable failures, bounded retries with backoff and jitter, circuit breakers, and rate limiting. Each addresses a different aspect of the problem.
The per-partner retry budget you described was a practical solution for the behavior you encountered in production. It's always interesting to see how architectures evolve based on real operational experience.