Every product eventually needs to tell users things — an order shipped, a login from a new device, a friend request. The naive version is a function call from application code straight to a provider’s API, and it works fine until you have five event sources, three channels, and a provider outage that starts silently dropping password reset emails. This is the design for the system that replaces that function call.

Problem

We’re building a notification service that other internal services call when something happens that a user should know about: OrderShipped, PasswordChanged, NewMessage, PromotionalOffer. The service decides how to deliver that event — email, SMS, push, or some combination — and guarantees the message reaches the user’s device or inbox, or fails in a way that’s visible to whoever needs to know.

The system sits between dozens of upstream producers (checkout, auth, messaging, marketing) and a handful of downstream providers (SendGrid for email, Twilio for SMS, FCM/APNs for push). It owns templating, user preferences, rate limiting per user and per channel, retries, and delivery tracking. It explicitly does not own business logic — a producer decides that a notification should fire, we decide how it gets there.

Requirements

Functional

  • Accept a notification request via API or internal event bus, addressed to a user id and an event type.
  • Resolve the user’s channel preferences (email/SMS/push, quiet hours, opt-outs) before sending.
  • Render a template per channel with event-specific variables.
  • Deliver through the right provider, retrying transient failures.
  • Deduplicate: the same logical event should never be delivered twice to the same user on the same channel.
  • Expose delivery status (queued, sent, delivered, failed, bounced) to producers and to an internal dashboard.

Non-functional

  • At-least-once delivery to the provider, with idempotent dispatch so retries don’t double-send.
  • p99 end-to-end latency (event received to provider accepted) under 5 seconds for transactional notifications; best-effort for bulk.
  • Durable: an event accepted by the API must not be lost even if downstream providers are down for an hour.
  • Provider-agnostic: swapping SendGrid for SES shouldn’t touch producer code.
  • Backpressure-safe: a slow provider must not block ingestion of new events.

Capacity Estimation

Assume 50M monthly active users, averaging 6-7 notifications each per day across all channels, giving roughly 10M notifications/day as a baseline, with marketing pushes spiking traffic 3x during campaign windows.

  • Average QPS: 10,000,000 / 86,400 ≈ 116/s
  • Peak QPS (3x for campaigns, plus daily peak-hour concentration of ~2x): 116 × 3 × 2 ≈ 700/s at the ingestion layer; dispatch workers need to sustain roughly 350/s sustained against providers after batching and dedup collapse some of that volume.
  • Storage per notification record: event id, user id, channel, template id, rendered payload reference, status, timestamps ≈ 400 bytes average.
  • Daily storage: 10M × 400 bytes ≈ 4 GB/day.
  • Annual storage: 4 GB × 365 ≈ 1.46 TB/year, before replication. With 3x replication and a year of retention for audit/compliance, budget roughly 4.4 TB.
  • Provider bandwidth: email payloads average 15 KB (HTML template), push/SMS payloads average 1 KB. At a rough 60/30/10 split (email/push/SMS), daily egress ≈ (6M × 15KB) + (3M × 1KB) + (1M × 1KB) ≈ 94 GB/day, or about 1.1 MB/s sustained, bursting higher during campaigns.
Back-of-the-envelope
Daily notifications
10M
peak 3x during campaigns
Ingest QPS (peak)
~700/s
10M / 86400 x 3 x 2
Dispatch QPS (sustained)
~350/s
after batching + dedup
Storage / year
~4.4 TB
400B/record, 3x replication
Provider egress
~1.1 MB/s avg
60% email, 30% push, 10% SMS

API Design

Producers talk to a thin HTTP API (also mirrored as an async event on Kafka for high-volume producers who don’t want a synchronous call in their critical path).

POST /v1/notifications HTTP/1.1
Content-Type: application/json
Idempotency-Key: order-48213-shipped

{
  "eventType": "order.shipped",
  "userId": "usr_9f21ac",
  "channels": ["email", "push"],
  "priority": "transactional",
  "templateData": {
    "orderId": "48213",
    "trackingUrl": "https://ship.example.com/t/48213",
    "carrier": "UPS"
  }
}

HTTP/1.1 202 Accepted
Content-Type: application/json

{
  "notificationId": "ntf_7c19e0",
  "status": "queued",
  "acceptedChannels": ["email", "push"]
}
GET /v1/notifications/ntf_7c19e0 HTTP/1.1

HTTP/1.1 200 OK
Content-Type: application/json

{
  "notificationId": "ntf_7c19e0",
  "eventType": "order.shipped",
  "status": "delivered",
  "channels": {
    "email": { "status": "delivered", "providerId": "sg_a91c", "deliveredAt": "2026-08-02T14:03:21Z" },
    "push":  { "status": "sent", "providerId": "fcm_28d1", "sentAt": "2026-08-02T14:03:19Z" }
  }
}
PUT /v1/users/usr_9f21ac/preferences HTTP/1.1
Content-Type: application/json

{
  "channels": { "email": true, "sms": false, "push": true },
  "quietHours": { "start": "22:00", "end": "07:00", "tz": "America/New_York" },
  "categories": { "marketing": false, "transactional": true }
}

The Idempotency-Key header is mandatory for producers — it’s how we collapse retried producer calls into a single logical notification. priority drives which queue the request lands in, which we cover in Deep Dive.

Data Model

CREATE TABLE notifications (
  id              BIGINT PRIMARY KEY,
  idempotency_key VARCHAR(128) NOT NULL,
  event_type      VARCHAR(64)  NOT NULL,
  user_id         BIGINT       NOT NULL,
  priority        VARCHAR(16)  NOT NULL DEFAULT 'transactional',
  template_data   JSONB        NOT NULL,
  status          VARCHAR(16)  NOT NULL DEFAULT 'received',
  created_at      TIMESTAMPTZ  NOT NULL DEFAULT now(),
  UNIQUE (idempotency_key)
);

CREATE TABLE notification_deliveries (
  id               BIGINT PRIMARY KEY,
  notification_id  BIGINT      NOT NULL REFERENCES notifications(id),
  channel          VARCHAR(16) NOT NULL,
  provider         VARCHAR(32) NOT NULL,
  provider_msg_id  VARCHAR(128),
  status           VARCHAR(16) NOT NULL DEFAULT 'queued',
  attempt_count    INT         NOT NULL DEFAULT 0,
  last_error       TEXT,
  sent_at          TIMESTAMPTZ,
  delivered_at     TIMESTAMPTZ,
  UNIQUE (notification_id, channel)
);

CREATE TABLE user_preferences (
  user_id     BIGINT PRIMARY KEY,
  channels    JSONB NOT NULL,
  categories  JSONB NOT NULL,
  quiet_hours JSONB,
  updated_at  TIMESTAMPTZ NOT NULL DEFAULT now()
);

CREATE INDEX idx_deliveries_status ON notification_deliveries (status, sent_at);
CREATE INDEX idx_notifications_user ON notifications (user_id, created_at DESC);

notifications is the logical event; notification_deliveries is one row per channel per notification, which is where retries, provider ids, and per-channel status live. Splitting them lets a single event fan out to N channels without duplicating the event payload, and lets us retry one channel’s failure without touching the others.

Field Purpose
idempotency_key unique constraint Collapses duplicate producer calls at write time, not read time
notification_deliveries.attempt_count Drives exponential backoff and dead-letter threshold
user_preferences.quiet_hours Checked at render time, not enqueue time — see Deep Dive

High-Level Architecture

High-level architecture

The ingestion API writes the notification synchronously to the database (this is our durability boundary — once that write commits, we own the delivery) and publishes to a priority queue. Dispatch workers pull from the queue, check preferences, render the template, and call the provider. Providers call back through webhooks (delivered, bounced, clicked) which update notification_deliveries.

Deep Dive

Fan-out and the preference check

A single event can target multiple channels, and each channel has independent delivery semantics. We fan out at enqueue time: the API writes one notifications row and N notification_deliveries rows (one per requested channel), then pushes N messages onto the queue — one per channel. This means a slow SMS provider never blocks the email send for the same event, and a channel-level retry doesn’t need to re-derive which channels were requested.

Preference checks happen at dispatch time, not enqueue time, because preferences can change between enqueue and send — a user might toggle off SMS thirty seconds after an event fires but before the queue drains during a backlog. Checking late means we never send to an opted-out channel, at the cost of one extra read (mitigated with a preference cache, discussed below).

Idempotency and exactly-once-feeling delivery

We can’t get true exactly-once delivery to an external provider over an unreliable network — the provider might accept our request and then the acknowledgment gets lost, and we’d retry into a duplicate send. We get as close as practical with two layers:

  1. Producer-level dedup: the Idempotency-Key unique constraint on notifications means a retried producer call is a no-op — same key returns the existing record instead of creating a second one.
  2. Provider-level dedup: before calling a provider, the dispatch worker acquires a Redis key SETNX dispatch:{delivery_id} with a 24-hour TTL. If the key already exists, another worker (or a retried job) is already handling this delivery, so we skip. This closes the “two workers pick up the same queue message during a rebalance” race that at-least-once queues create.

Decision

Deduplicate at the dispatcher

Redis SETNX on delivery_id with a 24h TTL, checked before every provider call

This is cheap (single round trip), works across worker restarts, and the TTL bounds memory growth. The cost is a hard dependency on Redis being available on the hot path — we treat a Redis outage as “fail closed and let the queue redeliver later” rather than sending unprotected, because a duplicate SMS is worse than a delayed one.

Template rendering

Templates are versioned and stored outside the hot path (S3 + a small in-memory LRU per worker), keyed by (templateId, channel, locale). Rendering is pure — no I/O beyond the initial fetch — so it’s cheap to retry. We compile templates to a simple AST once and cache the compiled form, because re-parsing Handlebars-style templates for every one of 700 requests/second at peak is a measurable amount of avoidable CPU.

Priority queues

We run three logical queues — critical (2FA, security alerts), transactional (order updates, receipts), and bulk (digests, marketing) — backed by separate SQS queues (or separate Kafka topics with dedicated consumer groups). Workers are provisioned with weighted concurrency: critical gets its own dedicated worker pool sized for its SLA regardless of bulk traffic, so a marketing blast never starves a password reset email. This is the single most important operational decision in this system — conflating queues is the most common way these systems fail during a big campaign send.

Rate limiting per user

Users can be over-notified — a chatty product sending 40 pushes in an hour is a churn risk, not an engagement win. We enforce a per-user, per-category rate limit (e.g., max 5 marketing notifications/day) using a sliding window counter in Redis, checked at dispatch time alongside preferences. Transactional notifications bypass this limit; the product decision is that a security alert always gets through.

Failure Handling

Failure Detection Response
Provider API down HTTP 5xx / timeout on send call Exponential backoff retry (1s, 5s, 30s, 5m, 30m), then dead-letter after 5 attempts
Provider accepts but never confirms delivery No webhook within SLA window Mark sent but flag as unconfirmed after 15 min; surfaced in dashboard, not retried (avoids dupes)
Worker crashes mid-dispatch Queue visibility timeout expires Message becomes visible again; Redis dedup key prevents double-send if the original call actually landed
Template render error Exception during render Delivery marked failed, alert fired, does NOT retry (a broken template won’t fix itself)
Preference service unavailable Timeout on preference lookup Fail closed for marketing (skip send), fail open for critical security notifications, both logged
Poison message (malformed payload) Repeated dispatch failure, same error After 3 identical failures, dead-letter immediately instead of exhausting full retry budget

The dead-letter queue isn’t a graveyard — a daily job replays DLQ entries after confirming the root cause is fixed, and anything older than 7 days gets escalated to a human for manual review, since silently dropping an event a user was owed is a support ticket waiting to happen.

Scaling

Horizontal scaling of dispatch workers

Dispatch workers are stateless consumers, so scaling is adding pods behind the queue’s consumer group — the queue depth (messages waiting, oldest message age) is the autoscaling signal, not CPU, because a slow provider can leave workers idle-waiting on network I/O while the queue backs up.

Provider throughput limits

Every provider imposes its own rate limit (SendGrid: ~10K/s on enterprise tiers with per-IP warmup requirements; Twilio: varies by number type; FCM: effectively unbounded per app but batched). We maintain a token bucket per provider, shared across all dispatch workers via Redis, so we never exceed the contracted rate regardless of how many workers are running — a burst of dispatch capacity should never turn into a provider throttling our entire account.

Database write scaling

At 10M notifications/day with a 3-channel average fan-out, we’re writing roughly 20-30M delivery rows/day. We shard notification_deliveries by user_id hash across a set of Postgres instances (or move to a wide-column store like Cassandra if we outgrow single-shard write throughput), since almost every read pattern — “show this user’s notification history,” “check this delivery’s status” — is naturally scoped to a user or a single delivery id, not a cross-shard scan.

Multi-region

For a global user base, we run dispatch workers regionally close to the provider endpoints they call (SES has regional endpoints; Twilio doesn’t care) and keep the notifications database as a single source of truth with read replicas in each region, since notification writes aren’t latency-critical enough to justify multi-region write conflict resolution — a 200ms cross-region write is invisible against a 5-second delivery SLA.

Trade-off

Kafka for the ingestion/dispatch queue

  • Replayable log — reprocess a bad deploy's worth of events
  • Natural fit if other services already consume notification events
  • You operate partitioning, rebalancing, and consumer lag

SQS (or equivalent managed queue) per priority tier

  • Fully managed, near-zero ops
  • Built-in dead-letter queues and visibility timeouts
  • No replay after a message is acked — a bug that mis-processes a batch is unrecoverable

Recommendation — Start on managed SQS-style queues split by priority; adopt Kafka only when you need other services to consume notification events as a stream, or hit a proven need for replay.

Observability

Three signals matter more than the rest:

  • Queue age (oldest unprocessed message per priority tier) — the earliest indicator that dispatch is falling behind, well before end-to-end latency dashboards catch it.
  • Delivery funnel per channel per provider: queued → sent → delivered → bounced, as a real-time counter, so a provider issue (e.g., SendGrid IP reputation drop) shows up as a bounce-rate spike within minutes, not as a support ticket a day later.
  • Idempotency collision rate: how often the Idempotency-Key unique constraint actually catches a duplicate. A rate near zero is fine (producers are well-behaved); a rising rate flags a retrying producer with a bug.

Structured logs carry notification_id and delivery_id on every line so a support engineer can trace one user’s one notification end-to-end without grepping. Alerts fire on DLQ depth, provider error rate over 5% in a 5-minute window, and critical-tier queue age over 30 seconds.

Trade-offs

The biggest trade-off in this design is synchronous durability write versus throughput. Writing every notification to Postgres before acknowledging the producer bounds our ingestion throughput to what a single primary can commit, but it’s what lets us promise a producer that an accepted event will eventually be delivered even through an hour-long provider outage. An alternative — acknowledge on queue publish alone — is faster but means a queue failure between publish and persistence silently loses the event. Given that notifications frequently carry security and compliance weight (password resets, fraud alerts), we accept the throughput ceiling.

The second trade-off is template rendering at dispatch time versus enqueue time. Rendering late means preference and locale changes are always reflected, at the cost of re-fetching template data per attempt on retry. We accept this because a stale rendered payload (wrong language, stale unsubscribe link) is a worse failure mode than a slightly more expensive retry path.

Final Architecture

Final architecture with dedup, priority tiers, and delivery tracking

The finished system separates concerns cleanly: the ingestion API owns durability and dedup at the door, three isolated priority queues protect the critical path from bulk traffic, dispatch workers own preference checks and rendering, and a shared Redis layer enforces per-delivery idempotency and per-provider rate limits across every worker. Every delivery is traceable from producer request to provider confirmation, and every failure mode has a defined, tested response rather than a silent drop.