Notification System - HLD
Difficulty: Intermediate Prerequisites:Message Queues, Fan-Out, and Caching
TL;DR
A multi-channel notification platform that delivers push, SMS, email, and in-app messages. It decouples βsomething happenedβ from βtell the userβ using an event bus.
flowchart LR
SERVICES["Internal Services<br/>order placed and payment etc"]:::client
BUS["Kafka"]:::async
NS["Notification Service<br/>template and routing"]:::service
PUSH["Push<br/>FCM APNs"]:::external
SMS["SMS<br/>Twilio"]:::external
EMAIL["Email<br/>SES"]:::external
USER["User"]:::client
SERVICES -->|"1. Publish notification event"| BUS
BUS -->|"2. Route by channel"| NS
NS -->|"3. Deliver via push"| PUSH
NS -->|"4. Deliver via SMS"| SMS
NS -->|"5. Deliver via email"| EMAIL
PUSH -->|"6. Push to device"| USER
SMS -->|"7. Deliver"| USER
EMAIL -->|"8. Deliver"| USER
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
classDef external fill:#4a1942,stroke:#f472b6,color:#e2e8f0
| Color | Role |
|---|---|
| Blue | Edge / API gateway |
| Green | Service |
| Purple | Async / message broker |
| Pink | External dependency |
| Yellow | Data store |
| Orange | Client |
In 3 sentences: Backend services emit events (βorder confirmedβ) to Kafka. The notification service consumes these, renders a template, picks the channel (push/SMS/email based on user preferences), and dispatches. Delivery tracking, retries, and send-time optimization ensure messages reach users when theyβre most likely to engage.
1. Understanding the Problem
A notification system lets product surfaces across a company send messages to users across multiple channels - push (mobile), email, SMS, in-app - without every team re-implementing delivery, preferences, retries, and rate limiting. The system must handle billions of events per day, respect user preferences, dedupe noisy senders, and prove delivery.
2. Naive First Cut
The whiteboard sketch before any real thought:
flowchart LR
APP["Product Service"]:::service
NS["Notification Sender"]:::service
APNS["APNs / FCM / SES / Twilio"]:::external
DB[("Users DB")]:::data
APP --> NS
NS --> DB
NS --> APNS
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef external fill:#fbcfe8,stroke:#be185d,color:#500724
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
Why it breaks under real load:
- One synchronous hop - if APNs is slow, the product service blocks. A single downstream hiccup freezes every checkout / signup / message-send.
- No retries - APNs drops a packet, the user never hears back. No receipt, no replay.
- No user preferences - customer unsubscribed from marketing but still gets marketing. Legal problem (GDPR, CAN-SPAM), trust problem.
- No rate limiting - batch job fires a million sends in a minute, exceeds APNs quota, APNs throttles us for everyone.
- No deduplication - two product teams both decide to welcome the user; they get two welcomes.
- No quiet hours - user in Sydney gets a 3am push.
- No observability - something failed for 10% of users today. Which 10%? We donβt know.
The rest of the doc evolves this into a queue-based, multi-channel, preference-aware notification platform.
3. Prior Art Weβre Drawing From
- Uber Consumer Communication Gateway (CCG) - central intelligence layer that manages quality, ranking, timing, and frequency of push notifications at the per-user level. Introduced after Uber hit 15+ hours a week of manual coordination trying to keep internal teams from stepping on each other. (blog)
- Airbnb notification platform - channel abstraction + per-user preference service, feeding multiple providers. Treats the sender API as a single pipe regardless of channel.
- LinkedIn Air Traffic Control - deduping + frequency capping layer that rides on top of all outbound member communications to prevent over-notification.
- Stripe Webhooks - the canonical βreliable delivery via outbox + retryβ pattern that applies directly to transactional notifications: persist before send, deliver asynchronously, expose attempt history.
- Courier, Braze, OneSignal, SendGrid - commercial templates for how this kind of platform is exposed to product teams: a single
send(user, template, data)API on top of a multi-channel router.
4. Functional Requirements
Core (top 3)
- Send a notification to a user through one or more channels (push, email, SMS, in-app) given a template ID and template variables.
- Respect user preferences - channel opt-ins, category opt-ins (marketing vs transactional), quiet hours, locale.
- Guaranteed at-least-once delivery with retries for transient failures, and exposed per-attempt status for debugging.
Below the line (out of scope)
- Rich content creation (images, carousels, deeplink generation).
- Campaign management UI / marketer console.
- ML-driven send-time optimization (Uber CCGβs specialty - weβll mention it but not build it).
- Reply handling (SMS 2-way conversations).
- Delivery receipts beyond what the provider returns.
5. Non-Functional Requirements
Core
- Scale - 1B notifications/day peak, ~50k/sec sustained, bursts to 500k/sec (campaign).
- Latency - transactional (OTP, security): P95 end-to-end < 5s. Marketing: P95 < 5 min.
- Reliability - at-least-once delivery, no silent loss. Acceptable duplicate rate < 0.1%.
- Multi-region - failover; user in EU hits EU stack.
Below the line
- Cost optimization (provider routing for cheapest path).
- Per-tenant isolation (if multi-tenant SaaS).
- Compliance reports (TCPA/GDPR opt-in audit trails).
6. Scale Estimation (Back-of-Envelope)
- Users: 500M DAU, generating events across all product surfaces
- Write QPS: 100K notifications/sec peak (5B notifications/day across 3 channels)
- Read QPS: 50K preference lookups/sec, 10K status queries/sec
- Storage: ~2TB notification metadata/year (attempts + audit trail)
- Bandwidth: ~10 Gbps outbound to providers at peak campaign burst
7. Core Entities
- Notification - one logical message intended for a user, with a template ID, variables, channels, and priority.
- User - the recipient. Has channel identifiers (device tokens, email, phone), preferences, locale, time zone.
- Template - a channel-specific content template with variable interpolation (``).
- Preference - per-user, per-category, per-channel opt-in/out. Quiet hours.
- Delivery Attempt - one try to hand off a notification to a channel provider. Status: SENT / DELIVERED / FAILED / THROTTLED.
- Campaign (optional) - a batch of notifications sharing a template, targeting a segment.
- Channel Provider - APNs, FCM, SES, Twilio, etc. Adapter per provider.
8. API / System Interface
One primary send API. Callers are internal product services, authenticated via service-to-service JWT.
Send a notification
POST /v1/notifications -> NotificationReceipt
Header: Idempotency-Key: <uuid>
Header: Authorization: Bearer <service-jwt>
Request body:
{
"userId": "u_1293847",
"templateId": "order_shipped_v3",
"variables": {
"orderId": "A-88273",
"trackingUrl": "https://tr.ck/X7Y8Z",
"carrier": "BlueDart"
},
"channels": ["PUSH", "EMAIL"],
"category": "TRANSACTIONAL",
"priority": "HIGH",
"dedupKey": "order-A-88273-shipped",
"deeplink": "myapp://orders/A-88273",
"overrides": {
"sendAt": null,
"expireAt": "2026-05-05T10:00:00Z"
}
}
Response (202 Accepted):
{
"notificationId": "n_7a3f2e91",
"state": "ACCEPTED",
"createdAt": "2026-05-04T12:01:33.412Z"
}
Get status and attempts (for support / debugging)
GET /v1/notifications/:id -> Notification + attempts
Response:
{
"id": "n_7a3f2e91",
"userId": "u_1293847",
"templateId": "order_shipped_v3",
"state": "DELIVERED",
"category": "TRANSACTIONAL",
"attempts": [
{ "channel": "PUSH", "provider": "APNs", "status": "SENT", "at": "2026-05-04T12:01:33.900Z", "providerResp": "200 APNs accepted" },
{ "channel": "PUSH", "provider": "APNs", "status": "DELIVERED", "at": "2026-05-04T12:01:34.102Z", "receipt": "apns-receipt-abc" },
{ "channel": "EMAIL","provider": "SES", "status": "SENT", "at": "2026-05-04T12:01:34.200Z", "providerResp": "250 OK" },
{ "channel": "EMAIL","provider": "SES", "status": "OPENED", "at": "2026-05-04T12:04:11.512Z", "userAgent": "Mail.app iOS 17" }
]
}
User preferences
GET /v1/users/:id/preferences -> Preference
PUT /v1/users/:id/preferences -> Preference
Preference shape:
{
"userId": "u_1293847",
"channels": { "PUSH": true, "EMAIL": true, "SMS": false, "IN_APP": true },
"categories": {
"TRANSACTIONAL": { "enabled": true, "channels": ["PUSH","EMAIL","SMS"] },
"SECURITY": { "enabled": true, "channels": ["PUSH","EMAIL","SMS"] },
"MARKETING": { "enabled": false, "channels": ["EMAIL"] }
},
"quietHours": { "start": "22:00", "end": "07:00", "timezone": "Asia/Kolkata" },
"frequencyCaps": { "MARKETING": 5, "SOCIAL": 20 },
"locale": "en-IN"
}
Device registration (push)
POST /v1/users/:id/devices -> register token
DELETE /v1/users/:id/devices/:token -> revoke
{
"deviceToken": "apns-token-base64=",
"platform": "IOS",
"appVersion": "12.3.0",
"timezone": "Asia/Kolkata"
}
Campaigns
POST /v1/campaigns -> Campaign
{
"name": "diwali_offers_2026",
"templateId": "promo_diwali_v1",
"segment": { "query": "country=IN AND active_last_7d=true AND age BETWEEN 18 AND 34" },
"scheduledAt": "2026-10-28T09:00:00+05:30",
"rateLimit": { "perSec": 20000, "perMin": 600000 },
"useSendTimeOptimization": true
}
Real-time subscription (in-app)
WSS /v1/users/:id/stream -> WebSocket
Header: Authorization: Bearer <user-jwt>
Server pushes JSON frames as notifications fire:
{ "type": "notification", "id": "n_7a3f2e91", "title": "Your order shipped", "body": "...", "at": "2026-05-04T12:01:33.412Z" }
Security notes:
- All HTTP APIs require service JWT; product services pass the end-userβs ID as data, not identity.
- The
/streamWebSocket uses the end-userβs JWT - authenticates the subscriber against theuserIdin the path. - Templates are pre-approved and versioned; raw text body isnβt accepted from callers to prevent content injection and compliance bypass.
- Device tokens stored encrypted at rest; tokens rotate when app reinstalls.
- Idempotency key is scoped per
(service, key)and kept 24h in Redis.
9. High-Level Design
Weβll grow the architecture in three passes - one per core functional requirement.
6.1 FR-1: Send a notification through multiple channels
Start with the minimum viable pipeline: accept, enqueue, fan out per channel, dispatch to the provider.
New components we need:
- Notification API - the single entry point for all product services. Receives βsend a notification to user Xβ requests, validates them, and enqueues for processing.
- Message Broker (Kafka) - decouples notification intake from delivery.
π‘ Kafka here acts as a buffer - if push notifications are slow today, the queue absorbs the backlog instead of slowing down the checkout flow that triggered the notification. Learn more β - Router - reads each notification event, decides which channels to use (push? email? SMS?), and fans out one message per channel to channel-specific topics.
- Channel Workers (Push, Email, SMS) - each specialized worker renders the template and calls the external provider. Isolated so a Twilio outage doesnβt affect push delivery.
- External Providers (APNs, SES, Twilio) - the actual delivery services. We donβt send emails ourselves - we hand them to SES/Mailgun, which handles the SMTP complexity.
π‘ FCM (Firebase Cloud Messaging) and APNs (Apple Push Notification Service) are the only way to send push notifications to Android and iOS devices respectively. Your server canβt push directly to phones - it must go through these gateways.
flowchart TD
APP["Product Service"]:::service
API["Notification API"]:::edge
Q["Message broker<br/>Kafka"]:::async
ROUTE["Router"]:::service
PUSH["Push Worker"]:::service
EMAIL["Email Worker"]:::service
SMS["SMS Worker"]:::service
APNS["APNs / FCM"]:::external
SES["SES / Mailgun"]:::external
TWIL["Twilio"]:::external
APP -->|"1. POST notification request"| API
API -->|"2. Enqueue for routing"| Q
Q -->|"3. Route by channel"| ROUTE
ROUTE -->|"4. Deliver via push"| PUSH
ROUTE -->|"5. Deliver via email"| EMAIL
ROUTE -->|"6. Deliver via SMS"| SMS
PUSH -->|"7. Send via APNs FCM"| APNS
EMAIL -->|"8. Render template"| SES
SMS -->|"9. Render template"| TWIL
classDef edge fill:#bfdbfe,stroke:#1d4ed8,color:#0c1f4a
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef external fill:#fbcfe8,stroke:#be185d,color:#500724
Legend
Step-by-step flow:
- Order Service calls
POST /v1/notificationsβ βHey, order A-88273 shipped, tell the user via push + emailβ - Notification API validates the payload, looks up the userβs preferences and locale, and persists the intent to a notifications table (audit + dedup)
- API publishes an event to Kafka and returns
202 Acceptedin ~20ms - the product service is free to move on - Router consumes the event, checks template rules + user preferences, and fans out: one message to
pushtopic, one toemailtopic - Push Worker picks up its message, renders the template in the userβs locale (βYour order has shipped! π¦β), and POSTs to APNs/FCM
- Provider response (accepted/rejected) is recorded as a delivery attempt for debugging (βwhy didnβt my user get their notification?β)
Why async? The product service call path has a strict latency budget (checkout is running). Hitting APNs synchronously is a timebomb - provider slowness becomes product slowness. Enqueueing gives us 10-50ms end-to-end on the hot path; the actual send happens on the workerβs clock.
Why persist before publish? Safety net. If the broker is down, we still have the row. A reconciler (covered later) sweeps notifications stuck in PENDING_PUBLISH and re-publishes.
6.2 FR-2: Respect user preferences
Preferences live in a dedicated service. Both the Notification API (at intake) and the Router (before fan-out) consult it.
New components we need (in addition to the ones above):
- Preference Service - owns all user notification settings: which channels are on/off, which categories theyβve opted out of, quiet hours, locale, and frequency caps.
π‘ This is the βdo not disturbβ brain - it prevents us from waking someone at 3am with a marketing push. - Preference Cache (Redis) - since every single notification triggers a preference lookup (50K/sec!), we cache preferences in Redis for microsecond reads instead of hammering Postgres.
- Preference DB (Postgres) - the durable source of truth for preferences. Updated when users change settings.
flowchart LR
API["Notification API"]:::edge
ROUTE["Router"]:::service
PREFS["Preference Service"]:::service
PREFDB[("Postgres<br/>user_preferences")]:::data
PREFCACHE[("Redis<br/>pref cache")]:::data
API -->|"1. Load user preferences"| PREFS
ROUTE -->|"2. Check user preferences"| PREFS
PREFS -->|"3. Lookup cached prefs"| PREFCACHE
PREFCACHE -. miss .-> PREFDB
classDef edge fill:#bfdbfe,stroke:#1d4ed8,color:#0c1f4a
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
What preferences capture:
- Channel opt-ins - SMS: off. Email: on. Push: on.
- Category opt-ins - Marketing: off. Transactional: always on (legally required in many jurisdictions). Security: always on.
- Quiet hours - per-user time window in userβs local time. Non-urgent notifications during quiet hours get deferred to a delayed queue; urgent (security, fraud, OTP) bypass.
- Locale + time zone - for template localization and quiet-hour computation.
- Frequency cap - max N marketing notifications per day.
Step-by-step flow (Router consulting preferences):
- Router receives the event: βsend MARKETING notification to user U via PUSH and EMAILβ
- Router asks Preference Service: βWhat are Uβs notification preferences?β
- Preference Service checks Redis cache (hit 99% of the time) β returns preferences
- Router filters: user has Marketing=on, Push=on, Email=on - both channels stay. If theyβd opted out of Marketing, the notification would be silently dropped here
- Router checks quiet hours: userβs timezone says itβs 2:30am β enqueue to a delayed queue that will fire at 7am when quiet hours end. (Security and transactional notifications bypass quiet hours - your OTP still arrives at 3am)
- Router checks frequency cap: βuser got 4/5 marketing notifications todayβ - still under the limit, proceed
- Fans out to the push and email channel topics
Why cache preferences in Redis? Every notification triggers a preference lookup. 50k/sec sustained = 50k/sec reads minimum. Postgres can do it, but Redis drops the latency from millis to microseconds and takes load off the DB for campaigns.
Write path: user updates prefs β API writes Postgres β invalidates Redis entry (write-through not worth the complexity; read-through handles the miss).
6.3 FR-3: Guaranteed at-least-once delivery with retries
Channel workers own the retry logic. The key mechanism is the outbox pattern between provider state and our own DB.
New components we need (in addition to the ones above):
- Retry Queue (delayed Kafka topic) - when a provider returns a transient error (timeout, 5xx, rate-limit 429), the failed message goes here with exponential backoff timing.
π‘ Exponential backoff means: wait 2s, then 10s, then 60s, then 5min before each retry. This prevents hammering a struggling provider. Learn more β - Dead Letter Queue (DLQ) - where permanently-failed messages go after exhausting all retries. These get reviewed by a human or an automated reconciler.
- Delivery Attempts table (Postgres) - records every single attempt to deliver a notification, including the providerβs response. Essential for debugging βwhy didnβt user X get their OTP?β
flowchart LR
PUSH["Push Worker"]:::service
APNS["APNs"]:::external
ATTEMPT[("Postgres<br/>delivery_attempts")]:::data
RETRY["Retry Queue<br/>delayed topic"]:::async
DLQ["Dead Letter Queue"]:::async
PUSH -->|"1. Send via APNs FCM"| APNS
PUSH -->|"2. Record delivery attempt"| ATTEMPT
PUSH -.timeout or 5xx.-> RETRY
RETRY -->|"3. Retry delivery"| PUSH
PUSH -.permanent failure.-> DLQ
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef external fill:#fbcfe8,stroke:#be185d,color:#500724
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
How a channel worker handles each notification (with retries):
What a worker does for each message:
- Read from channel topic.
- Render template with user variables + locale.
- POST to provider (APNs / SES / Twilio) with a timeout (2s for push, 5s for email/SMS).
- Write a
delivery_attemptsrow: status (SENT / FAILED / THROTTLED), provider response code, timestamp. - Classify the outcome:
- Success β commit Kafka offset, move on.
- Transient failure (5xx, timeout, 429 throttle) β publish to a delayed retry topic with exponential backoff (2s β 10s β 60s β 5min β 30min; max 5 retries).
- Permanent failure (400 bad token, invalid phone) β revoke the token in our user device table, send to DLQ.
Why not just retry forever? Permanent failures (invalid device token, phone number doesnβt exist) will never succeed no matter how many times we retry. Sending to DLQ and revoking the bad token prevents infinite loops and keeps the queue healthy. Transient failures (provider overloaded) usually resolve within minutes, so retries with backoff are the right call.
Why at-least-once, not exactly-once? Distributed systems canβt do exactly-once delivery across a network boundary - only at-least-once + idempotency on the receiving side. The Idempotency-Key header upstream and the dedupKey on the notification row let us detect dups on retry. Providers also dedupe on apns-collapse-id / FCM collapse_key - we pass our notification ID as collapse key so a retry doesnβt produce two banners on the device.
Why persist every attempt? Debugging (βwhy didnβt my user get the OTP?β) requires per-attempt receipts. When someone files a ticket, ops need to see the exact provider response.
DLQ ownership: a sweeper runs every 5 minutes, reviews DLQ messages, and either retries with longer backoff, or surfaces to a human dashboard after N tries.
10. Technology Choices
Vendor-agnostic with alternatives. Swap to a specific cloudβs services if youβre targeting one.
| Tier / purpose | What it stores | Access pattern | Primary pick | Alternatives |
|---|---|---|---|---|
| Notification primary | notifications - intent, template ID, state, dedupKey |
high write on send, index by (user, createdAt), point-read by id | PostgreSQL partitioned by day | MySQL, CockroachDB, Aurora |
| Delivery attempts | one row per provider call, with response | insert-heavy, query by notification_id for support | PostgreSQL monthly-partitioned, or Cassandra at very high rates | DynamoDB with TTL |
| User preferences | channel + category opt-ins, quiet hours, frequency state | low write, very high read (every send) | PostgreSQL + Redis cache | DynamoDB + DAX |
| Device registry | push tokens per user per device | medium write (registrations), high read | PostgreSQL sharded by user_id | DynamoDB, Cassandra |
| Event backbone | notification.requested, notification.dispatched, etc. |
ordered per user_id, replayable, at-least-once | Kafka | Kinesis, Google Pub/Sub, Pulsar |
| Delayed queue (quiet hours, retries) | messages keyed by wake-up time | producer writes with a readyAt; consumer only pulls ready ones |
Redis sorted sets keyed by timestamp, or Kafka with timer-topic wheel | SQS with visibility delay, RabbitMQ delayed exchange |
| Template store | rendered templates per channel per locale | read-heavy, versioned | S3 / object storage + Postgres metadata | Git-backed templates (GitOps) |
| Rate-limit counters | per-user / per-tenant counters | very high read+write, TTLβd | Redis token bucket / sliding window | DynamoDB with atomic counters |
| Analytics / reporting | daily sends, delivery rates, opt-out trends | OLAP scans, dashboards | Snowflake / BigQuery / ClickHouse via CDC | Redshift, Druid |
| Secrets | provider API keys, APNs certs | very low read, high sensitivity | Vault / AWS Secrets Manager | 1Password Connect, Parameter Store |
Why Postgres for notifications + attempts, not Cassandra or DynamoDB?
You get ACID within a send: persist the notification + initial state atomically. Support queries hit an indexed point lookup, not a full table scan. Partition by day to keep hot data in tiny tables. For companies sending 10B+/day, Cassandra becomes the right call because append-only writes beat Postgres WAL. Below that, Postgres is simpler and has full SQL.
Why Redis for the delayed queue?
Redis sorted sets (ZSET) give O(log N) insert and O(log N) range query by score. βGive me everything with score β€ now()β β pop them. This is the cleanest delayed queue pattern. Kafka timer-topic wheels work too but are harder to get right; use them only when volumes exceed what a Redis cluster handles.
Why Kafka for the event backbone?
- Ordering per partition - critical for βthis userβs notificationsβ to stay in order.
- Replayable - if the Push Worker had a bug yesterday, we reprocess the topic.
- Durable - tolerates consumer downtime; messages persist until we ACK.
Kinesis and Pub/Sub are equivalent on managed clouds.
11. Data Modeling
Postgres (Notification Primary β intent and state):
CREATE TABLE notifications (
notification_id UUID PRIMARY KEY,
user_id UUID NOT NULL,
channel VARCHAR(10) NOT NULL, -- push, email, sms, in_app
template_id VARCHAR(64) NOT NULL,
category VARCHAR(30), -- marketing, transactional, alert
status VARCHAR(15) NOT NULL, -- PENDING, DISPATCHED, DELIVERED, FAILED, SUPPRESSED
dedup_key VARCHAR(128),
payload JSONB NOT NULL,
created_at TIMESTAMP NOT NULL,
dispatched_at TIMESTAMP
) PARTITION BY RANGE (created_at);
CREATE INDEX idx_notifs_user ON notifications(user_id, created_at DESC);
CREATE INDEX idx_notifs_status ON notifications(status) WHERE status IN ('PENDING', 'DISPATCHED');
Redis (Rate-Limit Counters + Delayed Queue):
-- Rate limiting
Key: "rl:notif:{userId}:{category}:{window}" β Integer (count in current window)
TTL: window_size (e.g., 3600s for "max 5 marketing per hour")
-- Delayed queue (quiet hours, scheduled sends)
Key: "notif:delayed" β Sorted Set (score = send_at_timestamp, member = notification_id)
Postgres (User Preferences + Device Registry):
CREATE TABLE user_preferences (
user_id UUID NOT NULL,
channel VARCHAR(10) NOT NULL,
category VARCHAR(30) NOT NULL,
enabled BOOLEAN DEFAULT true,
quiet_hours_start TIME,
quiet_hours_end TIME,
PRIMARY KEY (user_id, channel, category)
);
CREATE TABLE device_tokens (
user_id UUID NOT NULL,
device_id VARCHAR(128) NOT NULL,
platform VARCHAR(10) NOT NULL, -- ios, android, web
push_token TEXT NOT NULL,
is_active BOOLEAN DEFAULT true,
updated_at TIMESTAMP NOT NULL,
PRIMARY KEY (user_id, device_id)
);
Access Patterns:
| Query | Data Source | How |
|---|---|---|
| Send notification (dispatch) | Kafka β Workers | Consume from notification.requested topic, check prefs + rate limit, call provider |
| Check user preferences | Redis (cached) β Postgres | Cache prefs in Redis with 5-min TTL; fallback to DB |
| Rate limit check | Redis | INCR rl:notif:{userId}:{category}:{window} β if > limit, suppress |
| Quiet hours handling | Redis ZSET | If in quiet hours: ZADD notif:delayed <quiet_end_timestamp> notifId β pick up later |
| Get notification history | Postgres | SELECT * FROM notifications WHERE user_id = ? ORDER BY created_at DESC LIMIT 20 |
How a Notification Is Rate-Limited and Delayed for Quiet Hours:
- Marketing notification triggered for user X β API inserts to Postgres with status=PENDING, publishes to Kafka
- Worker consumes event β checks Redis:
INCR rl:notif:{X}:marketing:1hβ returns 6 - Rule says max 5/hour β SUPPRESSED. Update Postgres status, stop here.
- If count was 3 (under limit): check user preferences. User has quiet hours 10pm-8am and itβs 11pm.
- Worker adds to delayed queue:
ZADD notif:delayed <8am_timestamp> notifId - A scheduler worker runs every minute:
ZRANGEBYSCORE notif:delayed 0 <now>β pops ready notifications and re-publishes to Kafka for delivery - On delivery: call FCM/APNs/SMTP, update status to DELIVERED, publish
notification.deliveredevent
12. Deep Dives
Running the flows above against the checklist surfaces eleven worth doing. Deep Dives 1-6 cover the core delivery mechanics; 7-11 address real-time delivery, template management, send-time optimization, engagement tracking, and broadcast.
Deep Dive 1 - Hot write path: notification intake at scale
Problem: Product teams send 500K notification requests per second during peak (flash sales, morning digests). The system that receives these must not become the bottleneck.
In simple terms: Imagine 500K βsend this notificationβ requests arriving every second. If each one requires a database write before responding, the database melts. We need a way to accept requests instantly and process them asynchronously.
Bad: Product service inserts directly into notifications + publishes to Kafka + writes an audit log. Three writes on the critical path. At 500k/sec burst, the DB is the bottleneck.
Good: Notification API does one insert with INSERT ... RETURNING id, then publishes. Two writes, still DB-bound. Campaigns that fire 10M sends in a minute still melt Postgres.
Great - outbox + CDC:
- Notification API writes one row in Postgres in a transaction that includes the event payload in an
outboxtable. - Debezium / logical replication tails the WAL and emits to Kafka - exactly once from WAL to Kafka.
- The product-facing API returns immediately after the DB commit (~5ms).
- Campaign batch writer uses
COPY FROMto bulk-insert millions of rows in seconds.
flowchart LR
API["Notification API"]:::edge
DB[("Postgres<br/>notifications + outbox")]:::data
CDC["Debezium"]:::async
KAFKA["Kafka"]:::async
ROUTE["Router"]:::service
API -->|"1. Persist notification"| DB
DB -->|"2. CDC stream"| CDC
CDC -->|"3. Stream changes"| KAFKA
KAFKA -->|"4. Route by channel"| ROUTE
classDef edge fill:#bfdbfe,stroke:#1d4ed8,color:#0c1f4a
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
Deep Dive 2 - Fan-out amplification: one event to N devices
Bad: Router looks up a userβs device tokens inline. User has 4 devices (phone, tablet, desktop, kiosk). Push worker fans out 4x. Fine for one user - not for a campaign that hits 50M users = 200M push sends.
Good: Router fans out once per device into the push topic. Each device is an independent delivery attempt. Works until we need to dedupe across devices for the same notification (e.g., user opens on phone, donβt ring the tablet 30s later).
Great - logical notification + per-device attempts + device-collapse:
- One
notificationrow, Ndelivery_attemptsrows. apns-collapse-id= notification_id: if the same notification retries, APNs replaces rather than stacking.- A βread receiptβ from the client marks the notification read and any still-pending retries are cancelled.
- For campaigns: partitioned bulk insert, each partition processes in parallel by a pool of Router workers.
Deep Dive 3 - Provider throttling and backpressure
Bad: Blast sends at APNsβ max rate. APNs rate-limits the whole tenant, legitimate transactional sends also get 429βd.
Good: Token bucket per provider, per channel. Workers pull from the channel topic only when a token is available.
Great - weighted bucket + priority lanes:
- Two topics per channel:
push.transactional(high priority),push.marketing(bulk). - Transactional workers have higher bucket capacity and refill rate.
- Marketing workers share a smaller bucket, self-throttle.
- When provider returns 429, workers exponential-backoff the bucket refill rate for that provider - system-wide TTLβd override in Redis.
- Noisy-neighbor isolation: per-tenant token buckets on top of the per-provider limit, so one product team canβt starve others.
flowchart LR
KAFKA1["push.transactional<br/>(high priority)"]:::async
KAFKA2["push.marketing<br/>(bulk)"]:::async
BUCKET1["Bucket: 10k/s"]:::service
BUCKET2["Bucket: 2k/s"]:::service
WORKERS1["Txn Workers"]:::service
WORKERS2["Mkt Workers"]:::service
APNS["APNs"]:::external
KAFKA1 -->|"1. Deliver"| BUCKET1
BUCKET1 -->|"2. Read file"| WORKERS1
KAFKA2 -->|"3. Deliver"| BUCKET2
BUCKET2 -->|"4. Read file"| WORKERS2
WORKERS1 -->|"5. Deliver via APNs"| APNS
WORKERS2 -->|"6. Deliver via APNs"| APNS
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef external fill:#fbcfe8,stroke:#be185d,color:#500724
Deep Dive 4 - Deduplication and frequency capping
Bad: Two product teams both emit welcome_user for a new signup. User gets 2 welcomes.
Good: dedupKey on the notifications table with a unique constraint. Second insert fails β second teamβs send is dropped. Works for exact dups.
Great - Air Traffic Control layer (from LinkedInβs playbook):
- A βpolicy checkβ step between Router and channel workers.
- Input: user ID, category, timestamp, content hash.
- Policies:
- Dedup - exact dedupKey match within last 24h β drop.
- Frequency cap - β€5 marketing per day per user. Counter in Redis (
INCR notif:mkt:{userId}:2026-05-04+ EXPIRE). - Category quota - no more than 3 βnew reviewβ notifications per hour.
- Global mute - user churned to quiet mode for N days β all marketing dropped.
- Tokens live in Redis, sharded by user_id. Per-user consistency holds; cross-user global quotas use a separate counter.
This is the layer Uber calls the CCG. Itβs where the ML logic for send-time optimization would eventually slot in.
Deep Dive 5 - Quiet hours and scheduled delivery
Bad: Check quiet hours synchronously, if in-quiet-hours, Thread.sleep(untilSomeTime). Worker threads pile up.
Good: If in quiet hours, compute readyAt = end_of_quiet_hours_in_user_tz and insert into a delayed queue. A scheduler wakes up and re-injects when ready.
Great - Redis ZSET with a dequeue poller:
ZADD delayed:notifications <readyAtEpoch> <notificationId>
A scheduler service polls every second:
ZRANGEBYSCORE delayed:notifications 0 <now> LIMIT 0 1000
For each returned ID, re-publish to the channel topic and ZREM. O(log N) inserts, O(log N + k) poll where k = batch size. Scales to hundreds of millions of scheduled notifications.
Alternatives: Kafka timer-topic tumbling wheel, DynamoDB TTL streams, SQS visibility timeout tricks. Pick Redis for simplicity and latency, others if volumes push past a single Redis cluster.
Edge case: user changes time zone mid-wait. Two reasonable answers:
- Let the scheduled time fire as originally computed (simplicity wins).
- Re-enqueue on preference change with the new tz (more accurate, more complex).
Most teams pick option 1 and accept occasional mis-timing.
Deep Dive 6 - Observability and delivery proof
Bad: βDid user X get their OTP?β - grep logs across 50 hosts. Hope someone logged what we need.
Good: Structured logs in ELK. Search by notification_id.
Great - first-class attempt history + delivery webhooks + dashboards:
delivery_attemptstable holds per-try status. Indexed by notification_id for O(1) support lookups.- Provider delivery callbacks (APNs feedback, SES SNS, Twilio webhook) write to an
inbound_receiptstable; a reconciler joins attempts β receipts to compute true delivery rate. - Real-time dashboard: delivery rate per channel, per provider, per category, per locale. Alert on sudden drops.
- Backstop for lost webhooks: periodic
/statuspoll for providers that support it (APNs Feedback Service); if we have a SENT attempt with no receipt after 30 min, we poll.
Deep Dive 7 - In-app notifications in real time
Bad: In-app notifications rely on polling. The mobile app hits GET /notifications?since=... every 30 seconds. Users see a 30-second lag; 1M DAUs = 33k req/sec of wasted polling.
Good: Server-sent events (SSE) over HTTP/2. The server keeps a unidirectional stream open; when a notification arrives for this user, the server writes a frame. SSE is one-way, text-only, and works through most proxies without configuration.
Great - WebSocket gateway with a presence layer + a fallback poll:
flowchart LR
APP["Mobile or Web"]:::client
LB["L4 Load Balancer"]:::edge
WS["WebSocket Gateway<br/>(sticky sessions)"]:::service
PRESENCE[("Redis<br/>presence: userId -> nodeId")]:::data
KAFKA["Kafka<br/>in-app topic"]:::async
FANOUT["In-App Fan-out"]:::service
APP -->|"1. Connect WebSocket"| LB
LB -->|"2. Route"| WS
WS -->|"3. Register presence"| PRESENCE
FANOUT -->|"4. Read from Kafka"| KAFKA
KAFKA -->|"5. Push to WS gateway"| WS
WS -->|"6. Push to client"| APP
classDef client fill:#fed7aa,stroke:#c2410c,color:#431407
classDef edge fill:#bfdbfe,stroke:#1d4ed8,color:#0c1f4a
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
How it works:
- Client opens a WebSocket to
/v1/users/:id/stream. Load balancer uses consistent hashing onuserIdto pin the connection to a specific WS gateway node (sticky sessions). This means a given user is always on the same node, simplifying routing. - On connect, the WS gateway writes
userId -> nodeIdto Redis (TTL 30s, refreshed by heartbeat every 10s). - When the In-App Fan-out worker consumes a notification, it looks up the target userβs
nodeIdin Redis. If present, it publishes to a Kafka topic keyed by node; each WS gateway node consumes its own keyed partition and delivers to the live socket. - If
nodeIdis missing (user offline), the worker writes the notification to an βinboxβ - a Redis listinbox:{userId}capped at 100 items + a Postgres backup. Next time the app connects, it drains the inbox as the first thing. - Fallback poll: even with WS, the app periodically (every 60s) calls
GET /notifications?since=<lastSeenId>as a safety net. This catches any notification lost to a transient WS hiccup. Itβs rare, but belt-and-suspenders.
Why WS over SSE: bidirectional frames give us acknowledgments (client says βgot it, showed badgeβ), and modern load balancers + browsers handle WebSocket fine. Also: we can multiplex multiple event types over one WS (notifications, typing indicators, presence).
Scale numbers: a single modern WS gateway node handles 50k-100k open sockets. For 10M concurrent users, 100-200 gateway nodes behind a consistent-hash LB.
Trade-off: sticky sessions complicate rolling deploys. Mitigation: graceful drain - new deploy tells existing sockets to reconnect, they get routed to the new node via the LB.
Deep Dive 8 - Template service: versioning, localization, rendering
Bad: Templates as hardcoded strings in worker code. Every copy change requires a deploy. Marketing canβt iterate. Translators need a developer.
Good: Templates in a database, fetched by ID at render time. Versioned. Marketing uses a console.
Great - immutable template versions + pre-compiled renderer + per-locale cache:
flowchart LR
CONSOLE["Marketing Console"]:::client
TMPL["Template Service"]:::service
S3[("Object Storage<br/>template artifacts")]:::data
DB[("Postgres<br/>template metadata")]:::data
CACHE[("Redis<br/>compiled templates")]:::data
WORKER["Channel Worker"]:::service
CONSOLE -->|"1. Create campaign"| TMPL
TMPL -->|"2. Load template from DB"| DB
TMPL -->|"3. Store rendered template"| S3
WORKER -->|"4. Render notification"| TMPL
TMPL -->|"5. Lookup cached template"| CACHE
CACHE -. miss .-> S3
classDef client fill:#fed7aa,stroke:#c2410c,color:#431407
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
Model:
- A template has an ID, a version, a channel (PUSH / EMAIL / SMS / IN_APP), and a locale (en, hi-IN, ja-JP, etc.).
- Each version is immutable. You donβt edit v3; you publish v4.
- The
templateIdon a send request resolves to the current active version per locale. - Old versions remain queryable - needed to render historical notifications in the support UI correctly.
Publishing flow:
- Marketer drafts a template in the console.
- Console uploads the template file (Mustache/Handlebars/MJML for email) to object storage at
templates/order_shipped_v4/en.mustache. - Pre-flight validation: variables referenced in the template must all exist in a registered
variableSchema. Catches typos before a real send fails. - Compliance reviewer approves (required for MARKETING category).
- Activate atomically: Postgres row updates the
active_versionpointer for(templateId, locale).
Render flow (from a channel worker):
- Look up
(templateId, userLocale)via Template Service. - Fetch compiled template from Redis. Miss β pull from S3 β compile (parse Mustache to AST) β cache.
- Render with the userβs variables. Run output sanitization (HTML escape for email, length-cap for SMS, JSON-safe for push payload).
- Return the rendered content to the worker.
Caching: compiled template objects stay in Redis with a long TTL (24h) because versions are immutable - thereβs no staleness risk. On activation, the service bumps a global version number which workers check cheaply to detect new active versions.
Locale fallback: hi-IN not found β try hi β try templateβs declared default locale β fail send. All falls through Template Service so workers donβt reinvent fallback logic.
Why pre-compile: a template is parsed once per node per version; subsequent renders are ~10-20 ΞΌs instead of parsing the template string each time. At 500k/sec we canβt afford the parser on every send.
Integration with channels:
- Push: renders to title + body + data payload. Max 4KB (APNs) / 4KB (FCM).
- Email: renders MJML β responsive HTML. Separate plain-text fallback.
- SMS: renders to plain text, segmented if > 160 GSM characters.
- In-app: renders to a structured JSON the client knows how to display.
Deep Dive 9 - Send-time optimization (Uber CCG-style)
Bad: All marketing fires immediately when the campaign is scheduled. Users get a 9am blast, half ignore it. Open rates tank.
Good: Default quiet hours + frequency cap. Better, but still one-size-fits-all. Your βmarketing hits at 9am localβ misses the user who opens the app at 7pm every day.
Great - per-user send-time prediction + constrained ranking:
flowchart LR
EVENTS["User Engagement Events<br/>Kafka"]:::async
FEATURE["Feature Store"]:::data
MODEL["Send-Time Model<br/>(trained offline)"]:::service
SERVING["Model Serving<br/>online inference"]:::service
ATC["Policy ATC"]:::service
RANKER["Ranker"]:::service
DELAYQ[("Redis ZSET<br/>per-user delayed queue")]:::data
EVENTS -->|"1. Publish user event"| FEATURE
FEATURE -->|"2. Compute features"| MODEL
MODEL -->|"3. Return prediction"| SERVING
ATC -->|"4. Score urgency"| SERVING
SERVING -->|"5. Return prediction"| RANKER
RANKER -->|"6. Return prediction"| DELAYQ
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
Three layers:
-
Feature store - per-user features: open-rate-by-hour histogram, click-rate-by-hour, last-active-hour, timezone, days since last notification. Updated in near real time from the engagement Kafka topic.
-
Send-time model - offline-trained (weekly) gradient-boosted model that predicts P(open user, hour-of-day, category). Lightweight: ~KBs per user, serves in <1ms. For each incoming notification with useSendTimeOptimization=true, inference returns the best hour in the next 24h. - Ranker - multiple notifications compete for attention. Per-user ranker runs a linear program (Uberβs actual approach - linear programming) that picks at most N notifications per day subject to:
- category priority (transactional > security > social > marketing),
- frequency caps,
- minimum spacing between notifications (15 min),
- predicted open probability.
The output is a scheduled list:
(notificationId, sendAt)pairs. Each is enqueued in the delayed Redis ZSET from Deep Dive 5.
Why linear programming rather than greedy: the βmax 5 marketing per day + min 15 min spacing + top-N by predicted openβ problem has conflicting constraints. Greedy picks the highest-scoring notification first and loses optimal coverage. LP gets the globally best schedule in milliseconds for per-user problems of this size (~dozens of candidates per user per day).
Only marketing and social categories go through this layer. Transactional and security bypass - they fire immediately.
Cost control: inference serving is the hot spot. At 100M users Γ 10 marketing candidates per day = 1B inferences/day. A small Redis-cached βbest hour per user per categoryβ result valid for 24h absorbs 95% of those.
Trade-off: adds latency to marketing sends. Thatβs fine - theyβre not time-sensitive. For βthis product just restockedβ youβd still fire with a shorter urgency window.
Deep Dive 10 - Engagement tracking (opens, clicks)
Bad: βDid the user see it?β - no clue. Only the provider knows they accepted it.
Good: Client-side reporting. App SDK pings POST /v1/notifications/:id/opened when user taps. Email has a tracking pixel. But: no reliable way to know for push without client SDK, and clients can lie or double-report.
Great - multi-source engagement ingestion + deduped event store:
flowchart LR
APP["Mobile SDK"]:::client
EMAIL["Email Pixel<br/>& tracking links"]:::client
PROVIDER["Provider Webhooks<br/>APNs Feedback / SES SNS"]:::external
COLLECT["Engagement Collector"]:::service
KAFKA["Kafka<br/>engagement topic"]:::async
DEDUP["Deduper"]:::service
DWH[("Analytics Store<br/>events - 90d hot")]:::data
COLD[("Parquet on S3<br/>cold - 2y")]:::data
APP -->|"1. Receive notification"| COLLECT
EMAIL -->|"2. Notify"| COLLECT
PROVIDER -->|"3. Delivery receipt"| COLLECT
COLLECT -->|"4. Publish engagement event"| KAFKA
KAFKA -->|"5. Dedup events"| DEDUP
DEDUP -->|"6. Write to warehouse"| DWH
DWH -->|"7. Sink data"| COLD
classDef client fill:#fed7aa,stroke:#c2410c,color:#431407
classDef service fill:#bbf7d0,stroke:#16a34a,color:#052e16
classDef async fill:#e9d5ff,stroke:#7c3aed,color:#3b0764
classDef data fill:#fde68a,stroke:#b45309,color:#451a03
classDef external fill:#fbcfe8,stroke:#be185d,color:#500724
Event types tracked:
SENT- worker handed off to provider.DELIVERED- provider confirmed or SDK ackβd receipt.OPENED- user tapped push / opened email.CLICKED- user tapped a tracked link.DISMISSED- user swiped away without opening.UNSUBSCRIBED- user hit unsubscribe link.BOUNCED- email bounce (hard / soft).COMPLAINED- spam report.
Pipeline:
- Sources normalize to a common envelope:
{notificationId, userId, event, at, source}. - Collector publishes to the engagement Kafka topic, partitioned by
notificationId. - Deduper drops duplicates per
(notificationId, event)within a 7-day window using a Bloom filter in Redis. Fixes double-reports from SDK + provider webhook for the same open. - Events land in ClickHouse for real-time dashboards (campaign open rate, per-locale performance). Daily roll-up to S3 Parquet for long-term analysis and ML training.
- Events also feed back into the feature store (Deep Dive 9) within minutes.
Why ClickHouse: column-store gives fast aggregation for SELECT category, hour, COUNT(*) FROM events WHERE date=today GROUP BY ... dashboards. Postgres would be too slow at the event volume (10B events/day).
Why a Bloom filter for dedup: exact dedup would require storing 10B IDs per week. Bloom filter accepts ~0.1% false positives (we occasionally drop a real second open) but uses ~100x less memory.
Engagement data flows back into:
- The send-time optimization model (better predictions).
- The ATC layer (noisy unsubscribe β suppress future marketing).
- Product dashboards (which templates work, which donβt).
Deep Dive 11 - Broadcast to large segments (optional)
Bad: βSend this to all 500M usersβ - the Router iterates the user list one-by-one. Takes hours.
Good: Parallelize the iteration. Shard the user list into batches of 10k, each processed by a worker pool. Still bounded by provider rate limits.
Great - pre-computed segment materialized view + push-time content personalization:
- For any large segment (country = IN, engaged_users, power_users_top_10pct), Segment Service materializes the user list nightly into an S3 object or a dedicated table.
- Broadcast Worker reads the segment in parallel partitions (e.g., 1000 partitions for 500M users) - each partition processed by a separate consumer group.
- Content is static across the segment (same template ID, same variables) but rendered per-user-locale at worker time.
- Provider rate limit is the ceiling: even fully parallelized, APNs caps total throughput at ~1M/sec. 500M users β ~8 minutes wall-clock minimum.
For truly instant large-audience cases (emergency civic alerts, security advisories), broadcast via channels that support topic-based fan-out natively:
- APNs Topic, FCM Topics - subscribers receive by topic subscription. Kicks fan-out to the provider.
- SMS carrier broadcast features for regional alerts (governmental use only).
For normal business broadcast, the segmented-queue approach is what you want.
13. Design Self-Audit
Weak spots checked:
- Text search - yes, if support team needs to search notifications by content. Push content through a search index (Elasticsearch) with a 14-day retention. Not core but worth mentioning.
- Stale prefs after write - user opts out, gets one more marketing send because the cache hasnβt invalidated. Fix: 1-second TTL on the cache entry plus event-driven invalidation; acceptable window.
- Single-region failure - primary Postgres region goes down. Active-passive: async replica in DR region; on failover, promote and re-route traffic. In-flight notifications in Kafka β consumer group repositions to the DR cluster thatβs mirrored via MirrorMaker.
- DLQ reconciliation - ops dashboard lists DLQ entries by reason. Auto-retry once after 1h, then require human decision.
- Cost at scale - egress to providers is free-ish for APNs/FCM, paid per send for SES/Twilio. Cost dashboard per category so Marketing knows their send cost.
- Hot user / broadcast - celebrityβs account triggers 100k notifications to followers. Covered in Deep Dive 11.
- In-app real-time delivery - user opens the app and expects to see the unread badge instantly. Covered in Deep Dive 7 via WebSocket gateway + presence + inbox.
- Template staleness and localization - marketing canβt edit strings without a deploy. Covered in Deep Dive 8.
- Over-notification and smart scheduling - users ignore poorly-timed marketing. Covered in Deep Dive 9 with Uber CCG-style send-time optimization.
- Engagement blindness - we send and hope. Covered in Deep Dive 10 with a unified engagement pipeline.
14. Core Flows
Flow 1 - Transactional send (OTP)
sequenceDiagram
actor User
participant AuthSvc as Auth Service
participant NotifAPI as Notification API
participant DB as Postgres
participant Kafka
participant Router
participant Prefs as Preference Svc
participant Worker as SMS Worker
participant Twilio
participant Attempts as attempts table
AuthSvc->>NotifAPI: POST /notifications (OTP, category=SECURITY)
NotifAPI->>DB: INSERT notification + outbox
DB-->>NotifAPI: id
NotifAPI-->>AuthSvc: 202 accepted (id)
DB-->>Kafka: CDC β notification.requested
Kafka->>Router: consume
Router->>Prefs: get prefs(userId)
Prefs-->>Router: prefs (security always-on, ignore quiet hours)
Router->>Kafka: publish sms.transactional (renderedPayload)
Kafka->>Worker: consume
Worker->>Twilio: POST /Messages
alt success
Twilio-->>Worker: 201
Worker->>Attempts: INSERT status=SENT
else transient 5xx
Twilio-->>Worker: 503
Worker->>Attempts: INSERT status=FAILED
Worker->>Kafka: publish retry with backoff
end
Walkthrough:
- Auth service calls the Notification API with the OTP template and category=SECURITY.
- API writes the notification and its outbox entry in one transaction.
- It returns 202 in under 20ms - the user sees βcode sentβ immediately.
- CDC picks up the commit and publishes to Kafka.
- Router consults Preference Service. Security overrides quiet hours and marketing opt-out.
- Router fans out to
sms.transactional(high-priority topic). - SMS worker renders and hits Twilio with a 5s timeout.
- On success, we record the attempt; on 5xx, we retry with exponential backoff; on 4xx we mark permanent failure and alert.
Failure case: Twilio webhooks tell us 15s later the SMS was actually undelivered (number disconnected). The reconciler joins the webhook to our delivery_attempts, flips the status to UNDELIVERED, and notifies Auth Service via its own outbound webhook so it can offer the user an alternate channel.
Flow 2 - Marketing campaign send (10M users)
sequenceDiagram
participant Marketer
participant CampaignAPI
participant SegBuilder as Segment Builder
participant DB as Postgres
participant Kafka
participant Router
participant ATC as Policy (ATC)
participant Prefs
participant DelayQ as Redis ZSET
participant Worker as Push Worker
participant APNs
Marketer->>CampaignAPI: POST /campaigns (segment, templateId)
CampaignAPI->>SegBuilder: resolve segment β user IDs
SegBuilder-->>CampaignAPI: 10M user IDs (stream)
CampaignAPI->>DB: bulk COPY notifications + outbox
DB-->>Kafka: CDC (batched)
loop per user
Kafka->>Router: consume
Router->>Prefs: get prefs
Prefs-->>Router: marketing=on, quiet=22-07 in Sydney
alt in quiet hours
Router->>DelayQ: ZADD readyAt=7am-Sydney
else not in quiet hours
Router->>ATC: check dedup + freq cap
ATC-->>Router: allowed
Router->>Kafka: push.marketing
Kafka->>Worker: consume
Worker->>APNs: POST /push (rate-limited bucket)
APNs-->>Worker: 200
end
end
Note over DelayQ,Router: scheduler polls every 1s<br/>ZRANGEBYSCORE 0 now
Walkthrough:
- Marketer calls CampaignAPI with a segment definition (e.g., βIndian users, 18-34, active in last 7 daysβ).
- Segment Builder streams the user IDs out of the user data warehouse.
- CampaignAPI bulk-writes notifications using
COPY- seconds, not minutes. - CDC streams events to Kafka in order.
- Router consults prefs + ATC for each. Users in quiet hours get deferred via Redis ZSET.
- Rate-limited workers drain the topic, respecting APNs throughput caps.
- Failures β retry topic with backoff.
Non-obvious failure: campaign writes succeed, CDC is behind by 10 minutes. We donβt block. Marketer sees βcampaign queuedβ because the row is committed; the delay is at most CDC lag, which alerting monitors. Acceptable for marketing.
Flow 3 - User updates preferences
sequenceDiagram
actor User
participant App as Mobile App
participant PrefAPI as Preference API
participant DB as Postgres
participant Cache as Redis
participant Kafka
User->>App: toggle marketing off
App->>PrefAPI: PUT /users/:id/preferences
PrefAPI->>DB: UPDATE preferences
PrefAPI->>Cache: DEL user:prefs:{id}
PrefAPI->>Kafka: publish preference.changed
PrefAPI-->>App: 200
Note over Kafka: downstream ATC listens<br/>to reset per-user counters
- App calls PUT.
- Preference API updates Postgres.
- Invalidates Redis cache (next read rebuilds).
- Publishes a
preference.changedevent - downstream consumers (ATC, counters) can react. - Returns 200.
Within milliseconds of the update, the next notification fan-out sees the new preference on cache miss β DB hit β re-cache.
State machine - a notificationβs lifecycle
stateDiagram-v2
[*] --> ACCEPTED
ACCEPTED --> QUEUED: published to Kafka
QUEUED --> DEFERRED: in quiet hours
DEFERRED --> QUEUED: wake up
QUEUED --> SUPPRESSED: policy drop
QUEUED --> DISPATCHING: worker picks up
DISPATCHING --> SENT: provider 200
DISPATCHING --> RETRYING: transient failure
RETRYING --> DISPATCHING: backoff elapsed
RETRYING --> FAILED: retries exhausted
SENT --> DELIVERED: provider receipt
SENT --> UNDELIVERED: provider receipt (failure)
FAILED --> [*]
SUPPRESSED --> [*]
DELIVERED --> [*]
UNDELIVERED --> [*]
15. Final Architecture
flowchart TB
SVC(["Product Services"]):::client
API["Notification API"]:::edge
KF[["Kafka"]]:::async
ROUTE["Router<br>policy and ranking"]:::service
PREFS["Preference Service"]:::service
TMPL["Template Service"]:::service
WORK["Channel Workers<br>push email sms"]:::service
INAPP["In-App Gateway<br>websocket"]:::edge
PG[("Postgres<br>notifications and prefs")]:::data
RD[("Redis<br>presence and delay queue")]:::data
DLQ[["Dead Letter Queue"]]:::async
ANALYTICS[("Analytics Store<br>engagement")]:::data
VENDORS[/"APNs SES Twilio"/]:::external
USER(["End Users"]):::client
SVC -->|"send notification"| API
API --> PG
API -->|"enqueue"| KF
KF --> ROUTE
ROUTE -->|"is this allowed"| PREFS
PREFS --> PG
ROUTE -->|"quiet hours and presence"| RD
ROUTE -->|"render body"| TMPL
ROUTE --> WORK
ROUTE --> INAPP
WORK --> VENDORS
WORK -->|"gave up"| DLQ
VENDORS --> USER
INAPP --> USER
USER -->|"opened or clicked"| ANALYTICS
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
classDef external fill:#4a1942,stroke:#f472b6,color:#e2e8f0
How it works end-to-end:
- Product Service sends notification β calls the Notification API with event type, user IDs, and payload
- Outbox pattern captures event β written to Postgres; Debezium CDC streams the change to Kafka
- Router consumes and evaluates β checks user preferences (Redis pref cache), applies frequency caps (ATC/Policy)
- Channel-specific topics populated β Router publishes to per-channel Kafka topics (push, email, SMS, in-app)
- Workers render and deliver β Push/Email/SMS Workers call Template Service for localized content, then send via APNs/FCM, SES, or Twilio
- In-App fan-out via WebSocket β In-App worker pushes through the WebSocket Gateway to connected users
- Failures retry with backoff β exponential backoff on transient errors; permanent failures go to DLQ
- Engagement tracked β opens, clicks, and provider webhooks flow to the Engagement Collector, feeding the analytics store and ML feature store for send-time optimization
Key Technologies
| Term | What it is |
|---|---|
| Kafka | Distributed event log decoupling notification intake from delivery - absorbs burst traffic so product services never block on slow providers. |
| APNs (Apple Push) | Apple Push Notification Service - the only gateway for delivering push notifications to iOS devices. |
| FCM (Firebase Cloud Messaging) | Googleβs push notification gateway for Android (and web) - your server canβt push directly to phones without going through FCM. |
| SES / Twilio | Amazon SES for email delivery, Twilio for SMS - external provider adapters wrapped behind channel workers. |
| Template Engine | Renders per-channel, per-locale message content from pre-approved templates with variable interpolation (e.g., ``). |
| Dead Letter Queue | Parking spot for permanently-failed messages after max retries - reviewed by ops or an automated reconciler. |
| Exponential Backoff | Retry strategy that waits progressively longer between attempts (2s β 10s β 60s β 5min) to avoid hammering a struggling provider. |
| Rate Limiting | Token bucket per provider and per tenant preventing any single sender from exhausting push/email/SMS quotas for everyone. |
Whatβs Expected at Each Level
This section helps you calibrate your depth. You donβt need to cover everything - just know whatβs expected for your level.
Mid-level
Design a basic system that receives events and sends notifications via push/email. Propose a queue between event generation and delivery. Understand why async processing matters - synchronous dispatch blocks the product service and canβt handle provider slowdowns gracefully.
Senior
Propose Kafka for event ingestion with consumer groups per channel. Discuss template service for message rendering, user preference management (opt-in/out per channel), and rate limiting per user. Explain retry strategies with exponential backoff and DLQ for permanently failed deliveries. Articulate the outbox pattern for guaranteed event capture.
Staff+
Address notification deduplication across channels (Air Traffic Control pattern from LinkedIn), priority queuing (critical alerts skip the queue), and A/B testing delivery times for engagement optimization. Discuss cost analysis across channels (SMS costs $0.01/msg vs push at $0) and provider routing for cost optimization. Cover regulatory compliance (CAN-SPAM, GDPR consent) and the operational cost of maintaining per-user frequency caps at scale.
π― Key Takeaways
- Multi-channel (push + email + SMS) with per-user preference routing
- Kafka decouples event producers from notification delivery
- Template engine separates content from channel logic
- At-least-once delivery with dedup on the client side
Related Designs
- Chat System - WebSocket real-time delivery
- Job Scheduler - scheduled and delayed notifications
- Twitter Feed - fan-out patterns
Related Concepts
Understand the building blocks used in this design:
- Message Queues β β decouple event producers from the per-channel senders
- Fan-Out Patterns β β deliver one event across many users and channels (push, email, SMS)
- Idempotency β β dedupe so a user isnβt notified twice for the same event
- Dead Letter Queue β β parks notifications that permanently fail delivery for inspection
- Retry & Backoff β β retries transient provider failures without hammering them
Discussion
Newest first