DEV Community

Avaneesh Yadav
Avaneesh Yadav

Posted on Originally published at buildingai.in

System Design: Building a Notification Service for 50 Million Users

Every non-trivial application sends notifications. Order confirmations. Password resets. Fraud alerts. Weekly digests. Yet the notification system is one of the most commonly underdesigned pieces of infrastructure — a single service that tries to do everything, with a shared database table that becomes a bottleneck at scale, and no retry logic until the first production outage.

This post designs a notification service that handles 50 million users, multiple channels (email, SMS, push, in-app), and 500,000 notifications per minute at peak. It's structured the way I'd walk through it in a system design interview — requirements first, capacity estimate, high-level design, component deep dive, trade-off decisions — with architecture diagrams at each stage and Spring Boot code for the non-obvious parts.

Requirements

Functional requirements:

  • Send notifications via email, SMS, push notification, and in-app feed
  • User preference management: per-channel, per-notification-type opt-in/opt-out
  • Template management: variable substitution, localization, A/B variants
  • Scheduling: send now, send at a specific time, send in user's local timezone
  • Deduplication: guarantee each notification is delivered at most once
  • Delivery status tracking: sent, delivered, failed, opened

Non-functional requirements:

  • 50 million users, up to 500,000 notifications/minute at peak (campaign sends)
  • P99 delivery latency < 5 seconds for transactional notifications (password reset, fraud alert)
  • P99 delivery latency < 5 minutes for bulk campaign sends
  • At-least-once delivery (retry on failure), with idempotency to prevent duplicates
  • Multi-region deployment for geo-based regulations (GDPR, SMS routing)

Out of scope: Real-time bidirectional chat (that's a separate WebSocket service), notification analytics beyond delivery status.

Capacity Estimate

Before touching architecture, numbers:

Metric Value Reasoning
Daily notifications 200M 4 per user per day average
Peak notifications/min 500K Campaign sends: 10M users × 3 notifications over 1 hour
Notification record size 2 KB Template rendered + metadata
Daily storage (hot) 400 GB 200M × 2 KB
Storage retention 90 days For audit + re-delivery
Kafka throughput needed ~8K messages/sec sustained, 8.5K peak 500K/min = 8,333/sec

This tells us we need a message queue (not direct DB writes) to absorb campaign spikes, and columnar or time-series storage for 90-day retention at 36 TB scale.

High-Level Architecture

flowchart TD
    A[API Clients\nMobile · Web · Backend Services] --> B[Notification API\nSpring Boot · REST]

    B --> C{Preference\nService}
    C -->|User opted in| D[Notification Router]
    C -->|User opted out| Z[/Drop/]

    D --> E[(Template\nService)]
    E --> F[Kafka\nnotification-events topic]

    F --> G[Email Worker\nSES / SendGrid]
    F --> H[SMS Worker\nTwilio / Vonage]
    F --> I[Push Worker\nFCM / APNs]
    F --> J[In-App Worker\nPostgreSQL + WebSocket]

    G & H & I & J --> K[(Delivery Log\nTimescaleDB)]
    G & H & I --> L[Dead Letter Queue\nnotification-dlq topic]

    L --> M[DLQ Processor\nRetry with backoff]
    M --> F

    style A fill:#1e293b,stroke:#334155,color:#e2e8f0
    style B fill:#0e7490,stroke:#0891b2,color:#f0f9ff
    style C fill:#7c3aed,stroke:#8b5cf6,color:#faf5ff
    style D fill:#0e7490,stroke:#0891b2,color:#f0f9ff
    style E fill:#0e7490,stroke:#0891b2,color:#f0f9ff
    style F fill:#b45309,stroke:#d97706,color:#fffbeb
    style G fill:#065f46,stroke:#059669,color:#f0fdf4
    style H fill:#065f46,stroke:#059669,color:#f0fdf4
    style I fill:#065f46,stroke:#059669,color:#f0fdf4
    style J fill:#065f46,stroke:#059669,color:#f0fdf4
    style K fill:#1e3a5f,stroke:#2563eb,color:#eff6ff
    style L fill:#7f1d1d,stroke:#dc2626,color:#fef2f2
    style M fill:#7f1d1d,stroke:#dc2626,color:#fef2f2
    style Z fill:#374151,stroke:#6b7280,color:#f9fafb

Seven distinct responsibilities, seven components. Let me walk through each.

Component Deep Dive

Notification API

The entry point. Two types of callers:

Transactional triggers — backend services (Order Service, Auth Service) calling directly with a structured event:

POST /api/v1/notifications
{
  "type": "ORDER_CONFIRMED",
  "recipientId": "usr_abc123",
  "templateVariables": {
    "orderId": "ord_xyz789",
    "totalAmount": "$149.99",
    "estimatedDelivery": "Sep 12"
  },
  "channels": ["EMAIL", "PUSH"],     // override preference? No — preference service decides
  "deduplicationKey": "order-confirm-ord_xyz789"  // idempotency
}
Enter fullscreen mode Exit fullscreen mode

Campaign sends — marketing platform sending in bulk. The API accepts a campaign job (not individual notifications) and fans it out internally:

POST /api/v1/campaigns
{
  "templateId": "tmpl_monthly_digest",
  "audienceSegmentId": "seg_premium_users",
  "scheduledAt": "2026-09-10T09:00:00",
  "timezoneStrategy": "RECIPIENT_LOCAL"  // send at 9 AM in each user's timezone
}
Enter fullscreen mode Exit fullscreen mode

The API layer does nothing except validate, deduplicate, and enqueue. No channel logic, no template rendering, no external calls. It must return in under 50ms.

Idempotency by deduplicationKey:

@PostMapping("/notifications")
public ResponseEntity<NotificationAck> send(@RequestBody @Valid SendRequest request) {
    String dedupKey = request.deduplicationKey();

    // Check Redis for recent duplicate (TTL = 24 hours)
    if (deduplicationCache.exists(dedupKey)) {
        return ResponseEntity.ok(NotificationAck.duplicate(dedupKey));
    }

    NotificationEvent event = eventFactory.create(request);
    kafkaProducer.send("notification-events", event.recipientId().toString(), event);
    deduplicationCache.set(dedupKey, Duration.ofHours(24));

    return ResponseEntity.accepted().body(NotificationAck.accepted(event.id()));
}
Enter fullscreen mode Exit fullscreen mode

Using recipientId as the Kafka partition key ensures all notifications for one user land in the same partition — guaranteeing ordering for that user's in-app feed without requiring global ordering.

Preference Service

The preference check runs before any message hits the queue. A notification not worth sending is not worth routing, rendering, or failing on.

flowchart LR
    A[Incoming Event\ntype · recipientId] --> B{Global\nopt-out?}
    B -->|Yes| Z[/Drop silently/]
    B -->|No| C{Channel\npreference?}
    C -->|EMAIL disabled| D[Remove EMAIL\nfrom channels]
    C -->|SMS disabled| E[Remove SMS\nfrom channels]
    C -->|No channels left| Z
    C -->|Channels remain| F{Quiet hours?}
    F -->|In quiet hours| G[Schedule for\nend of quiet period]
    F -->|Not in quiet hours| H[Route to queue]

    style Z fill:#374151,stroke:#6b7280,color:#f9fafb
    style G fill:#7c3aed,stroke:#8b5cf6,color:#faf5ff
    style H fill:#065f46,stroke:#059669,color:#f0fdf4

Preference data lives in a Redis hash per user — fast enough to check before every notification without being a bottleneck:

HGETALL user:pref:usr_abc123
 {
    "global_optout": "false",
    "email_enabled": "true",
    "sms_enabled": "false",
    "push_enabled": "true",
    "quiet_start": "22:00",
    "quiet_end": "08:00",
    "timezone": "America/New_York"
  }
Enter fullscreen mode Exit fullscreen mode

Preference data is also written to PostgreSQL as the source of truth. Redis is the read cache with a 5-minute TTL. Preference updates invalidate the Redis key immediately.

Template Service

Templates are versioned, localized, and support A/B testing:

tmpl_order_confirmed/
├── v3/
│   ├── en/
│   │   ├── email.html
│   │   ├── email.subject
│   │   ├── sms.txt          (160 char limit enforced)
│   │   └── push.json        (title + body + data payload)
│   └── es/
│       ├── email.html
│       └── ...
└── v2/                      ← previous version, still served for in-flight notifications
Enter fullscreen mode Exit fullscreen mode

Template rendering is synchronous within the worker — not a separate network call — because the worker already has the template variables from the event payload:

@Service
class TemplateRenderer {

    private final TemplateEngine templateEngine;  // Thymeleaf or Pebble

    public RenderedNotification render(NotificationEvent event, Channel channel) {
        Template template = templateRepository.getTemplate(
            event.type(), channel, event.locale(), event.templateVariant()
        );
        return new RenderedNotification(
            templateEngine.process(template.content(), event.variables()),
            channel
        );
    }
}
Enter fullscreen mode Exit fullscreen mode

Templates are cached in-process (Caffeine, 10-minute TTL) in each worker. Template cache invalidation sends a Kafka event to all worker instances — they evict and reload on next render.

Kafka: The Central Backbone

One topic, multiple consumer groups — one per channel. This is the fan-out mechanism:

notification-events (topic)
├── email-notification-group    → Email Workers (12 instances)
├── sms-notification-group      → SMS Workers (8 instances)
├── push-notification-group     → Push Workers (10 instances)
└── inapp-notification-group    → In-App Workers (6 instances)
Enter fullscreen mode Exit fullscreen mode

Each consumer group processes at its own pace. Email might be slower (due to ESP rate limits) without backing up SMS delivery. The channels are completely decoupled by the queue.

Partition count is the key scaling decision. At 8,333 messages/second peak and aiming for 1,000 messages/second per partition (leaving headroom), we need at minimum 9 partitions. Use 24 partitions — round number divisible by 1, 2, 3, 4, 6, 8, 12 consumers, giving flexibility to scale workers without rebalancing.

Channel Workers

Each worker follows the same structure regardless of channel:

@KafkaListener(
    topics = "notification-events",
    groupId = "email-notification-group",
    containerFactory = "emailKafkaListenerContainerFactory"
)
public class EmailNotificationWorker {

    private final TemplateRenderer templateRenderer;
    private final EmailProvider emailProvider;
    private final DeliveryLogService deliveryLog;

    @KafkaHandler
    public void handle(NotificationEvent event, Acknowledgment ack) {
        String notificationId = event.id();
        try {
            RenderedNotification rendered = templateRenderer.render(event, Channel.EMAIL);
            EmailResult result = emailProvider.send(rendered, event.recipientEmail());

            deliveryLog.record(notificationId, Channel.EMAIL, DeliveryStatus.SENT,
                Map.of("messageId", result.providerId(), "provider", "SES"));
            ack.acknowledge();

        } catch (RetryableException e) {
            // Don't ack — Kafka will redeliver to this consumer group
            log.warn("[{}] Retryable failure on email send, will redeliver: {}", notificationId, e.getMessage());
            // Don't ack — message will be redelivered after max.poll.interval.ms

        } catch (PermanentException e) {
            // Ack the message (remove from main topic) and send to DLQ
            deliveryLog.record(notificationId, Channel.EMAIL, DeliveryStatus.FAILED,
                Map.of("reason", e.getMessage()));
            dlqProducer.send("notification-dlq", event.recipientId(), new DlqEvent(event, e));
            ack.acknowledge();
        }
    }
}
Enter fullscreen mode Exit fullscreen mode

Retryable vs. permanent exceptions:

Exception Class Behavior
Rate limit (429) Retryable No ack, redeliver after backoff
Network timeout Retryable No ack, redeliver
Invalid email address Permanent Ack + DLQ
Unsubscribed (ESP bounce) Permanent Ack + update preference + skip DLQ
Template rendering error Permanent Ack + DLQ + alert

Retry and DLQ Flow

sequenceDiagram
    participant W as Channel Worker
    participant K as Kafka Main Topic
    participant D as DLQ Topic
    participant P as DLQ Processor
    participant L as Delivery Log

    W->>K: Consume notification event
    W->>W: Attempt delivery

    alt Delivery succeeds
        W->>L: Record SENT status
        W->>K: Acknowledge (offset committed)

    else Retryable failure (429, timeout)
        Note over W,K: Worker does NOT acknowledge
        K-->>W: Redeliver after poll interval
        W->>W: Attempt 2 (exponential backoff: 1s, 2s, 4s)
        W->>W: Attempt 3

        alt Succeeds on retry
            W->>L: Record SENT status
            W->>K: Acknowledge
        else Still failing after 3 attempts
            W->>D: Publish to DLQ with attempt count + error
            W->>L: Record FAILED status
            W->>K: Acknowledge (remove from main topic)
        end

    else Permanent failure
        W->>D: Publish to DLQ immediately
        W->>L: Record PERMANENTLY_FAILED status
        W->>K: Acknowledge
    end

    P->>D: Consume DLQ event
    P->>P: Check if recoverable (age, error type, attempt count)

    alt Recoverable and attempt < 5
        P->>K: Re-publish to main topic with attempt_count+1
    else Max attempts exceeded
        P->>L: Record ABANDONED status
        P->>P: Alert on-call (PagerDuty)
    end

The DLQ processor runs on a 15-minute schedule for most failures. For transactional notifications (password reset, fraud alert), it runs every 60 seconds — delivery SLA matters more for these.

Delivery Log (TimescaleDB)

Delivery tracking generates time-series data. One row per notification per channel, written at every state transition:

CREATE TABLE notification_delivery_log (
    notification_id  UUID        NOT NULL,
    channel          VARCHAR(20) NOT NULL,
    recipient_id     UUID        NOT NULL,
    status           VARCHAR(30) NOT NULL,   -- SENT, DELIVERED, OPENED, FAILED, ABANDONED
    provider_msg_id  VARCHAR(100),           -- ESP/SMS provider's message ID for callbacks
    error_code       VARCHAR(50),
    attempt_count    INTEGER DEFAULT 1,
    metadata         JSONB,
    occurred_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    PRIMARY KEY (notification_id, channel, occurred_at)
);

SELECT create_hypertable('notification_delivery_log', 'occurred_at');

-- Retention: auto-drop partitions older than 90 days
SELECT add_retention_policy('notification_delivery_log', INTERVAL '90 days');
Enter fullscreen mode Exit fullscreen mode

TimescaleDB's hypertable partitioning handles the write volume (10,000+ inserts/second at peak) and makes time-range queries efficient: "show me all failed email notifications in the last 24 hours" hits one or two partitions instead of scanning 36 TB.

Database Schema

erDiagram
    NOTIFICATIONS {
        uuid id PK
        varchar type
        uuid recipient_id FK
        varchar dedup_key UK
        jsonb template_variables
        varchar[] channels
        varchar status
        timestamptz scheduled_at
        timestamptz created_at
    }

    USER_PREFERENCES {
        uuid user_id PK
        boolean global_optout
        jsonb channel_preferences
        varchar quiet_start
        varchar quiet_end
        varchar timezone
        timestamptz updated_at
    }

    TEMPLATES {
        uuid id PK
        varchar notification_type
        varchar channel
        varchar locale
        integer version
        text content
        text subject
        boolean active
        timestamptz created_at
    }

    CAMPAIGNS {
        uuid id PK
        uuid template_id FK
        varchar segment_id
        varchar timezone_strategy
        varchar status
        timestamptz scheduled_at
        integer total_recipients
        integer sent_count
        timestamptz created_at
    }

    NOTIFICATIONS ||--o{ USER_PREFERENCES : "recipient"
    NOTIFICATIONS ||--o{ TEMPLATES : "uses"
    CAMPAIGNS ||--|| TEMPLATES : "uses"

The Trade-offs That Matter

1. Fan-out on write vs. fan-out on read

For campaign sends (10 million recipients), you have two strategies:

  • Fan-out on write: When the campaign is submitted, create one notification record per recipient immediately, then process from the queue. Pros: consistent delivery tracking, simple worker logic. Cons: 10M database writes before a single notification is sent.
  • Fan-out on read: Store one campaign record. Workers resolve the recipient list at delivery time. Pros: minimal writes upfront. Cons: recipient list must be re-resolved at send time, harder to track per-recipient status.

Recommendation: Fan-out on write, but asynchronously. The campaign scheduler writes notification records in batches of 1,000 to Kafka directly, not to PostgreSQL first. Workers consume from Kafka and persist delivery status. The notification record exists in the delivery log, not as a pre-created row in the notifications table.

2. At-least-once vs. exactly-once delivery

Exactly-once delivery with Kafka requires Kafka transactions and idempotent producers — significant operational complexity. At-least-once with deduplication in the API layer (the deduplicationKey) gives 99.9% of the benefit with 10% of the complexity.

The remaining 0.1% case: network partition causes the Kafka producer to retry after a successful write, sending the event twice. The deduplication cache (Redis with 24-hour TTL) catches this if the deduplicationKey is set. Callers that don't set a deduplicationKey accept at-least-once semantics explicitly.

3. Notification ordering

In-app notifications need to appear in chronological order for a given user. Email and SMS do not — the user sees them in their email client's sort order anyway.

The partition-by-recipientId strategy guarantees ordering within a partition for a given consumer group. For in-app, we additionally attach a sequence number from a per-user Redis counter when writing to the in-app feed. This handles the edge case where two events for the same user land in the same partition but the in-app worker processes them out of order (rare, but possible with concurrent workers).

4. Template rendering location

Render at the API layer (before enqueue) or at the worker layer (after dequeue)?

Rendering at the worker layer is better because:

  • Template data can be stale at render time rather than at enqueue time (template variables fetched fresh)
  • Failed renders become visible at delivery time with full context for debugging
  • Template updates between enqueue and delivery pick up the new template (important for campaigns scheduled days ahead)

The cost: the template variables must be carried in the Kafka event payload, not fetched again at render time. Include all needed data in the event at enqueue time.

Scaling to 500,000 Notifications Per Minute

flowchart LR
    subgraph "API Tier (stateless, auto-scale)"
        A1[API Instance 1]
        A2[API Instance 2]
        A3[API Instance N]
    end

    subgraph "Kafka (24 partitions)"
        K[notification-events]
    end

    subgraph "Email Workers (12 instances = 12 partitions)"
        E1[Email Worker 1]
        E2[Email Worker 2]
        E3[Email Worker 3..12]
    end

    subgraph "SMS Workers (8 instances)"
        S1[SMS Worker 1..8]
    end

    subgraph "Push Workers (10 instances)"
        P1[Push Worker 1..10]
    end

    subgraph "In-App Workers (6 instances)"
        I1[InApp Worker 1..6]
    end

    A1 & A2 & A3 --> K
    K --> E1 & E2 & E3
    K --> S1
    K --> P1
    K --> I1

At 500,000/min (8,333/sec), with 24 partitions, each partition handles ~347 events/sec. Email workers (slowest channel due to ESP rate limits) use 12 partitions. SMS, push, and in-app share the remaining 12 partitions across their consumer groups — each consumer group reads all 24 partitions independently.

Bottleneck analysis at 500K/min:

Channel Max throughput per worker Workers needed Kafka partitions
Email 50 sends/sec (SES rate limit) 167 → use 200 24 (shared)
SMS 200 sends/sec (Twilio) 42 → use 50 24 (shared)
Push 1,000 sends/sec (FCM batch API) 8.3 → use 10 24 (shared)
In-App 2,000 writes/sec (Postgres) 4.2 → use 6 24 (shared)

Email is the bottleneck. The solution isn't more partitions — it's multiple ESP accounts (SendGrid primary, SES fallback), each with their own rate limit, and distributing sends across accounts in a round-robin pattern within the email worker.

What Goes Wrong in Production (And How to Handle It)

1. Preference service is down

Fail open (send the notification, skip preference check) for transactional notifications (password reset, fraud alert). Fail closed (drop the notification) for marketing content. The distinction is in the notification type's config.

2. ESP (email service provider) outage

The circuit breaker opens on the primary ESP. Workers automatically route to the secondary ESP. The DLQ processor drains backed-up notifications when the primary recovers.

3. Template rendering fails for one user

The worker sends the notification to the DLQ with error_type: TEMPLATE_RENDER_ERROR. The DLQ processor alerts engineering. The rest of the campaign continues — one bad template variable doesn't stop 9,999,999 other sends.

4. Kafka consumer lag spikes

Set an alert when consumer group lag exceeds 100,000 messages (about 12 seconds of peak throughput). Root cause is usually: a slow downstream (ESP rate-limited), workers processing too slowly (increase parallelism), or a Kafka partition imbalance (rebalance the consumer group).

Spring Boot Wiring

@Configuration
class NotificationKafkaConfig {

    @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, NotificationEvent>>
    emailKafkaListenerContainerFactory(ConsumerFactory<String, NotificationEvent> cf) {
        var factory = new ConcurrentKafkaListenerContainerFactory<String, NotificationEvent>();
        factory.setConsumerFactory(cf);
        factory.setConcurrency(12);  // 12 threads = 12 partitions for email group
        factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
        factory.getContainerProperties().setPollTimeout(3_000);

        // Dead letter publishing: after 3 failed attempts, route to DLQ
        var recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate,
            (rec, ex) -> new TopicPartition("notification-dlq", rec.partition()));
        var errorHandler = new DefaultErrorHandler(recoverer,
            new FixedBackOff(1_000L, 3L));
        factory.setCommonErrorHandler(errorHandler);

        return factory;
    }
}
Enter fullscreen mode Exit fullscreen mode

Interview Cheat Sheet

If you're preparing for system design interviews, here's what interviewers are listening for on this question:

Topic Key point
Fan-out mechanism Consumer groups on Kafka, not fan-out in the API
Idempotency Dedup key in Redis, not in the database (too slow)
Preference check Before the queue, not inside the worker
Failure handling Three categories: success, retryable, permanent
DLQ Separate topic, separate processor, not just retry in-place
At-least-once Acceptable with dedup key; exactly-once is complexity with marginal gain
Partition key recipientId for ordering guarantee per user
Scale number Know your throughput math (events/min → events/sec → partitions needed)
Database choice TimescaleDB for delivery log time-series, Redis for preferences cache

The question interviewers are waiting for you to ask: "What's the SLA difference between transactional and marketing notifications?" Answering this before they prompt it shows you understand the business constraints that drive architecture, not just the technology.

Notification systems are where distributed systems complexity meets business requirements that change constantly (new channels, new preference types, new regulatory requirements). The architecture above is designed to absorb those changes: add a new channel by adding a consumer group, add a new preference type by extending the Redis hash, add a new region by adding a Kafka cluster in that region.

The design is the easy part. The hard part is the retries.

Avaneesh Yadav is Engineering Manager at HashedIn by Deloitte, designing enterprise systems at scale. He writes about system design and AI architecture at buildingai.in.

Top comments (0)