Skip to main content

Distributed Notification System

Sharpened prompt. Design a service that delivers 1B notifications/day across push (APNs/FCM), SMS, email, and in-app, respecting per-user preferences and quiet hours, never sending a duplicate, delivering an OTP in under 5 seconds even during a 100M-user marketing broadcast, and degrading gracefully when a vendor has a regional outage.

Everyone draws the same three boxes: queue, workers, providers. The interview is decided on deduplication, priority isolation, and what you do when APNs starts returning 429.

1. Problem framing​

Functional requirements​

  • Multi-channel: push, SMS, email, in-app, webhook — with per-notification channel preference and fallback chains.
  • Templating with localisation and variable substitution.
  • User preferences: category opt-outs, per-channel opt-outs, quiet hours in the user's local time zone.
  • Scheduled and immediate sends; campaign-scale broadcasts.
  • Delivery tracking: sent, delivered, opened, failed, and why.

Non-functional requirements​

PropertyTargetConsequence
Transactional latency (OTP)P99 < 5s end to endDedicated priority lane with reserved capacity
Broadcast throughput100M recipients in < 30 min≈ 55k sends/sec sustained
Duplicate rateEffectively zeroIdempotency at ingest and at the provider call
Availability99.95%Multi-provider per channel with automatic failover
Preference correctness100% for legal categoriesOpt-out is a compliance requirement, not a feature

Back-of-the-envelope​

Volume: 1B/day = 11,600/sec average; peak 5× = ~58,000/sec
Broadcast: 100M recipients / 30 min = 55,000/sec on top of baseline
Fan-out row: 1B × 300 B of delivery record = 300 GB/day (TTL 30 days -> 9 TB hot)
Dedup keys: 1B/day × 24h TTL × 50 B = 50 GB in Redis
Providers: APNs ~ 10k msg/s per HTTP/2 connection; ~20 connections for peak
Twilio: per-number limits (1 SMS/sec long code) -> a pool of short codes

The Twilio line is the one to notice: SMS throughput is bounded by regulation, not engineering. No amount of horizontal scaling raises a long-code limit; you buy short codes or you queue.

2. High-level architecture​

Note that eligibility is evaluated before the message enters a lane, not inside the delivery worker. Filtering 40% of a 100M-recipient broadcast at ingest saves 40M queue messages; filtering it at the worker saves nothing.

3. Component inventory​

ComponentConcrete choiceWhy this one
Ingest busKafka, one topic per priorityPartition count per topic is how you reserve capacity per lane
Dedup storeRedis with TTL (fast path) + DynamoDB conditional write (durable)Redis alone loses the guarantee on failover; DynamoDB alone costs too much per message
Preference storeAerospike or Redis, warmed from PostgresSub-millisecond lookups on the critical path for every recipient
Template engineHandlebars/Liquid with precompiled ASTs, ICU MessageFormat for pluralsRendering 100M messages means compile-once, execute-many
PushAPNs HTTP/2 multiplexed, FCM HTTP v1One long-lived connection carries thousands of concurrent streams
ResilienceResilience4j or Envoy outlier detectionPer-provider circuit breakers with half-open probes
SchedulingThe job schedulerQuiet-hour deferral is just a delayed job
Delivery recordsCassandra, partition by user_id, 30-day TTLWrite-heavy, read-by-user, naturally expiring

4. The toughest parts​

4.1 Deduplication across retries, replays, and channels​

Why it's hard. The same logical notification arrives multiple times through genuinely different routes: an upstream service's at-least-once Kafka producer, a consumer rebalance replaying uncommitted offsets, an operator replaying a topic after a bug fix, and a retry after a timeout where the send actually did succeed. Users notice duplicates immediately and lose trust; five identical "You have a new message" pushes is a product failure of the first order.

Solution — a deterministic idempotency key claimed atomically, at two layers.

def idempotency_key(n: Notification) -> str:
# Deterministic: same logical notification -> same key, from any producer,
# on any retry. A random UUID here defeats the entire mechanism.
return hashlib.sha256(
f"{n.user_id}|{n.event_type}|{n.entity_id}|{n.time_bucket(60)}".encode()
).hexdigest()[:32]

async def claim(key: str, ttl_s: int = 86400) -> bool:
# SET NX is a single atomic op: claim and check in one round trip.
return await redis.set(f"dedup:{key}", "1", nx=True, ex=ttl_s) is not None

The time_bucket(60) term is the design decision worth defending: it says "the same event for the same user within the same minute is one notification." That collapses genuine duplicates while still allowing a legitimate second notification about the same entity an hour later.

Redis alone is not sufficient — a failover can lose the last second of writes, precisely the window where duplicates cluster. For irreversible or high-value sends (payment receipts, OTPs), promote the claim to a DynamoDB conditional write. Two tiers, cost proportional to consequence.

A second, distinct problem is cross-channel duplication: push and email both firing for one event. Solve it at the eligibility stage with a channel-selection policy per category ("primary channel push, escalate to email only if undelivered after 10 minutes"), not with more dedup keys.

4.2 OTPs must never queue behind a marketing blast​

Why it's hard. A campaign enqueues 100M messages. An OTP for a user logging in right now enters the same queue at position 100,000,001. FIFO gives it a 30-minute wait; the login fails, the user is locked out, support lights up. A single shared priority field does not fix this — one queue means one head, and the head is occupied.

Solution — physical isolation, not logical priority. Separate Kafka topics and separate worker pools and separate provider connection pools per lane.

lanes:
p0_critical: # OTP, security alerts, fraud
topic: notif.p0
partitions: 64
workers: {min: 40, max: 200}
provider_connections: dedicated # its own APNs/Twilio pool, never shared
sla_p99_seconds: 5
p1_transactional: # receipts, order updates
topic: notif.p1
partitions: 128
workers: {min: 60, max: 400}
sla_p99_seconds: 60
p3_bulk: # marketing
topic: notif.p3
partitions: 256
workers: {min: 20, max: 100}
rate_limit_per_sec: 40000 # hard cap: bulk can never consume all provider quota
sla_p99_seconds: 1800

Three properties make this work. Dedicated provider connections mean a saturated bulk pool cannot starve p0 of sockets. A hard rate cap on bulk reserves headroom in the vendor's quota. And independent autoscaling means p0's 40 idle workers are a deliberately purchased insurance policy, not waste.

The line to say: "Priority queues within a single queue don't isolate — they only reorder. Under sustained load the low-priority backlog still consumes the shared resource. Isolation has to be physical."

4.3 Provider failure, throttling, and the cost of failing over​

Why it's hard. Vendors fail in three different ways, each needing a different response. Hard rejection (invalid token, unsubscribed) must never be retried — it will always fail and it damages your sender reputation. Throttling (429) means slow down, not switch. Regional brownout (5xx, timeouts) means switch. Treating all three as "retry" is what turns a vendor blip into your outage.

Solution — classify the failure, then act.

CircuitBreaker apns = CircuitBreaker.of("apns", CircuitBreakerConfig.custom()
.slidingWindowType(COUNT_BASED).slidingWindowSize(100)
.failureRateThreshold(50) // 50% of last 100 calls failed
.waitDurationInOpenState(Duration.ofSeconds(30))
.permittedNumberOfCallsInHalfOpenState(5) // probe before fully reopening
.recordException(e -> e instanceof ProviderUnavailable) // NOT InvalidToken
.build());

Result send(Push p) {
try {
return apns.executeSupplier(() -> apnsClient.send(p));
} catch (Throttled t) {
rateLimiter.reduceBy(0.5); // AIMD: halve, then ramp back linearly
return Result.retryAfter(t.retryAfter());
} catch (InvalidToken t) {
tokenStore.remove(p.token()); // terminal: clean up, never retry
return Result.dropped();
} catch (CallNotPermittedException open) {
return fallbackChannel(p); // circuit open: SMS instead of push
}
}

Two subtleties worth raising unprompted:

  • Failover is not free. Switching push→SMS costs real money (fractions of a cent versus ~$0.0075 per SMS) and is a worse user experience. Gate fallback by notification priority: p0 falls back immediately, p3 never falls back — it waits or is dropped.
  • AIMD on throttling. On a 429, multiplicatively decrease your send rate and additively recover. This is TCP congestion control applied to a vendor API, and naming it that way is a strong signal.

Also process the feedback channel religiously. APNs returns 410 Unregistered; FCM returns NOT_REGISTERED. Failing to purge those tokens means a growing fraction of your sends are wasted, and both Apple and Google will eventually throttle you for it.

4.4 Quiet hours across 400 time zones, at ingest speed​

Why it's hard. "Don't notify me between 22:00 and 08:00" must be evaluated in the user's local time, for 100M users, at 58k/sec, on the critical path. Time zones have DST transitions (a local time can be ambiguous or nonexistent), users travel, and a naive "store the UTC offset" model silently breaks twice a year for half the world.

Solution — store the IANA zone, not an offset; evaluate locally; defer rather than drop.

from zoneinfo import ZoneInfo # IANA database, DST-correct

def deliver_at(user: User, now_utc: datetime, category: str) -> datetime | None:
if category in CRITICAL: # OTP and security ignore quiet hours
return now_utc
tz = ZoneInfo(user.iana_tz) # e.g. "America/Sao_Paulo", never "-03:00"
local = now_utc.astimezone(tz)
start, end = user.quiet_start, user.quiet_end # 22:00, 08:00
in_quiet = (local.time() >= start or local.time() < end) if start > end \
else (start <= local.time() < end)
if not in_quiet:
return now_utc
if category in DROP_IF_QUIET: # a digest is worthless 10 hours late
return None
# Defer to the end of quiet hours, jittered so 10M users don't all fire at 08:00:00.
resume = datetime.combine(local.date(), end, tzinfo=tz)
if local.time() >= start: resume += timedelta(days=1)
return resume.astimezone(UTC) + timedelta(seconds=hash(user.id) % 900)

That final jitter line is the detail that separates a working system from an incident: without it, every deferred notification in a time zone fires in the same second and you have rebuilt the thundering herd from the job scheduler.

Preference evaluation also has to be fast: a Redis/Aerospike lookup per recipient at 58k/sec. Cache preferences in-process with a 60-second TTL and invalidate via Pub/Sub on change — a stale opt-out for 60 seconds is a real compliance risk, so for legally-binding categories (marketing under GDPR/CAN-SPAM) read through to the authoritative store and accept the latency.

4.5 Broadcasting to 100M users without melting the fan-out​

Why it's hard. "Notify every user about the new feature" is one API call that becomes 100M messages. Materialising 100M rows synchronously blocks the API for hours. Doing it in one worker takes days. Doing it in a thousand workers without coordination duplicates work when one restarts.

Solution — a two-phase, checkpointed, chunked expansion.

The campaign API returns immediately with a campaign ID. A segmenter splits the audience into 10,000-user chunks recorded in a chunk table. Expanders claim chunks with a conditional update (WHERE status='pending'), expand them into the bulk topic, and mark them complete. A crashed expander's chunk is reclaimed after its lease expires and re-expanded — at-least-once at the chunk level, made safe by the per-message idempotency key from 4.1.

This gives you three properties that matter operationally: progress visibility (chunks done / total), pause and resume (stop claiming chunks), and bounded blast radius on a bad campaign (kill it after 3% delivered, not 100%).

4.6 Frequency capping without a global lock​

Why it's hard. "At most 3 marketing notifications per user per week" is a cross-message constraint. Two campaigns running concurrently each check the cap, each see 2, each send — the user gets 4. And the check happens 100M times per campaign, so it must be cheap.

Solution — an atomic counter with a rolling window, checked and incremented in one operation.

-- freq_cap.lua: KEYS[1]=cap:{user}:{category}, ARGV={now_ms, window_ms, max}
redis.call('ZREMRANGEBYSCORE', KEYS[1], 0, ARGV[1] - ARGV[2]) -- drop old entries
local n = redis.call('ZCARD', KEYS[1])
if n >= tonumber(ARGV[3]) then return 0 end
redis.call('ZADD', KEYS[1], ARGV[1], ARGV[1] .. ':' .. math.random())
redis.call('PEXPIRE', KEYS[1], ARGV[2])
return 1

The check and the increment are one atomic script — the same lesson as the rate limiter, and it is worth pointing out that frequency capping is rate limiting with a business label on it.

One caveat to volunteer: the counter increments at eligibility time, but delivery can still fail. Either accept slight under-delivery (usually correct — caps are ceilings, not quotas) or compensate by decrementing on terminal failure.

4.7 Knowing whether anything actually arrived​

Why it's hard. "Sent" means you handed bytes to APNs. It does not mean the phone displayed anything. The device may be off for three days, the token may be stale, the user may have force-quit the app, the carrier may have silently dropped the SMS. Without real delivery data you cannot distinguish "our system is healthy" from "we have been shouting into a void for six hours."

Solution — model a state machine, and instrument the transitions.

accepted → suppressed (preference/cap/quiet)
→ queued → sent → delivered → opened
↘ soft_failed → retrying → sent
↘ hard_failed (invalid token, unsubscribed) → cleanup

Store every transition in Cassandra keyed by (user_id, notification_id) with a 30-day TTL, and stream them to your metrics pipeline. Then alarm on the ratios, not the counts: delivered/sent dropping below its 7-day baseline for one provider in one region is the earliest possible signal of a vendor brownout — usually 10–20 minutes before the vendor's own status page updates.

APNs delivery receipts require a feedback service; FCM provides them via BigQuery export; email needs webhook ingestion for bounces and complaints. Wire all three into the same feedback processor so token cleanup and suppression-list management are one code path.

5. What breaks first​

EventFirst failureMitigation
100M broadcast startsBulk lane saturates provider quotaHard rate cap on p3; dedicated p0 connections
APNs regional brownoutPush failures pile up, retries amplifyCircuit breaker + AIMD; fallback to SMS for p0 only
Redis dedup failoverDuplicate window of ~1sDynamoDB conditional write for high-value categories
Kafka consumer rebalanceReplayed offsets → duplicatesIdempotency key makes replay safe by construction
Preference DB slowEligibility becomes the bottleneck at 58k/secIn-process cache with 60s TTL + Pub/Sub invalidation
Quiet hours end at 08:0010M deferred notifications fire at onceDeterministic jitter across a 15-minute window
Bad template deployedMalformed messages to millionsCanary the first 0.1% of a campaign, auto-halt on error-rate spike

6. Cheat sheet​

  • Dedup: deterministic key = hash(user, event, entity, time_bucket); SET NX in Redis, conditional write in DynamoDB for high-value.
  • Priority: physical isolation — separate topics, workers, and provider connections. Logical priority does not isolate.
  • Providers: classify failures into terminal / throttle / unavailable, and respond differently to each. AIMD on 429.
  • Quiet hours: IANA zones (never offsets), defer with jitter, exempt critical categories.
  • Broadcast: chunked, checkpointed expansion with pause/resume and a kill switch.
  • Frequency caps: atomic Lua ZSET window — rate limiting with a product name.
  • Observability: alarm on delivered/sent ratio per provider per region, not on absolute counts.
  • The one-liner: "The queue is trivial. The system is a preference engine, a deduplicator, and a set of isolated lanes that keep a marketing campaign from delaying somebody's login code."