PrepZone Logo
PrepZone

Notification System

Fan out push, email and SMS through queues, templates and user preference filters.

Read these first

Design StreamHub's notification service: when a followed streamer goes live, notify subscribers via push, email, or SMS based on their preferences — 10M events/day, 100M deliveries/day, under 30 seconds latency for push.

Requirements

Functional requirements

  • Send push, email, and SMS notifications triggered by platform events.
  • User preference filters: opt-in/out per channel and notification type.
  • Template rendering with dynamic content (streamer name, stream title).
  • Delivery status tracking (sent, delivered, failed, bounced).
  • Rate limiting: max N notifications per user per hour.

Non-functional requirements

  • Latency: push under 30 s from event; email under 5 min; SMS under 1 min.
  • Throughput: 100M deliveries/day ≈ 1.2K deliveries/sec average, ~5K peak.
  • Reliability: at-least-once delivery with deduplication.
  • Availability: 99.9% — degraded email acceptable; push is priority.

Estimation

  • 10M events/day × avg 10 recipients = 100M deliveries/day.
  • Push payload ~500 bytes; email ~5 KB with HTML template.
  • Storage for delivery logs: 100M × 200 bytes ≈ 20 GB/day — partition by date, TTL 90 days.

Notification platform (AWS)

INTEGRATION
MSK eventstream.started
COMPUTE
Notification ro…EKS
INTEGRATION
SNS → APNspush
INTEGRATION
Amazon SESemail
INTEGRATION
PinpointSMS
MSK event → router → SNS fan-out to push, SES email, Pinpoint SMS.

API design

Java
# Internal event ingestion (not public)
POST /internal/v1/notifications:
  body:
    event_type: "stream.started"
    actor_id: "sh_4420"
    payload:
      stream_id: "live_9912"
      title: "Friday Night Ranked"

# User preference management (public)
GET  /v1/users/me/notification-preferences
PUT  /v1/users/me/notification-preferences:
  body:
    push:  { stream_started: true,  new_follower: false }
    email: { stream_started: false, weekly_digest: true }
    sms:   { stream_started: false }
Java
CREATE TABLE notification_preferences (
    user_id       BIGINT NOT NULL,
    channel       VARCHAR(16) NOT NULL,   -- push, email, sms
    event_type    VARCHAR(64) NOT NULL,
    enabled       BOOLEAN DEFAULT true,
    PRIMARY KEY (user_id, channel, event_type)
);

CREATE TABLE delivery_log (
    id            BIGINT PRIMARY KEY,
    user_id       BIGINT NOT NULL,
    channel       VARCHAR(16) NOT NULL,
    event_type    VARCHAR(64) NOT NULL,
    status        VARCHAR(16) NOT NULL,   -- pending, sent, failed
    provider_id   VARCHAR(128),
    created_at    TIMESTAMPTZ DEFAULT now()
);

Architecture

StreamHub production architecture (AWS)

HTTPSstaticmissAPICLIENT
Mobile / WebStreamHub cli…
NETWORK
Route 53GeoDNS routing
NETWORK
CloudFrontCDN + WAF edge
NETWORK
AWS ALBTLS terminati…
NETWORK
API GatewayJWT · rate li…
STORAGE
Amazon S3media origin
COMPUTE
Amazon EKSAPI · auth · …
DATABASE
ElastiCachesessions · ho…
DATABASE
RDS Postgresprimary + rep…
INTEGRATION
Amazon MSKdomain events
ANALYTICS
OpenSearchstream discov…
OPS
CloudWatchmetrics · X-R…
End-to-end path from user to data — reference this when placing any new service.

Pipeline stages

  • Event ingestion: Kafka consumer receives stream.started events.
  • Recipient resolver: Query follower list; filter by preferences and rate limits.
  • Fan-out queue: Separate queues per channel (push, email, SMS) for independent scaling.
  • Template engine: Render channel-specific content from Jinja/Handlebars templates.
  • Provider adapters: FCM/APNs (push), SendGrid (email), Twilio (SMS) — swappable.
  • Delivery tracker: Update status; retry failures with exponential backoff; DLQ after max attempts.
Java
def process_event(event: StreamStarted):
    followers = db.get_followers(event.streamer_id, limit=100_000)
    for batch in chunks(followers, 1000):
        for user in batch:
            prefs = get_preferences(user.id, "stream.started")
            if prefs.push:
                push_queue.publish(Notification(user, event, channel="push"))
            if prefs.email:
                email_queue.publish(Notification(user, event, channel="email"))

Scaling fan-out for celebrities

AspectSynchronous fan-outQueue-based fan-out
Latency to producerMinutes for 10M followersMilliseconds — publish to queue and return
Failure handlingOne provider timeout blocks allPer-channel retry and DLQ
ScalingLimited by single threadScale consumers per channel independently
  • Latency to producer

    Synchronous fan-outMinutes for 10M followers
    Queue-based fan-outMilliseconds — publish to queue and return
  • Failure handling

    Synchronous fan-outOne provider timeout blocks all
    Queue-based fan-outPer-channel retry and DLQ
  • Scaling

    Synchronous fan-outLimited by single thread
    Queue-based fan-outScale consumers per channel independently

Quick recall

Everything you need if you only revisit this box.

  • Fan out via separate queues per channel — push, email, SMS scale independently.
  • Filter by user preferences before enqueueing; rate-limit per user per hour.
  • Celebrity fan-out (10M followers) requires batching and queue sharding — never synchronous.
  • Template engine renders channel-specific content; provider adapters are swappable.
  • At-least-once delivery with dedup by (user_id, event_id, channel).
  • Delivery log with TTL for debugging; DLQ for permanently failed notifications.

Test yourself

Answer these before moving on — recall is what makes it stick.