System Design Problem

Design a Live Comments System (like Facebook Live / YouTube Live)

Commonly Asked By:MetaGoogleTwitchByteDance

Interview Setup

Interview Prompt

Design a live comment system for streaming video. Viewers post comments during a live stream, see new comments in near real time, and can scroll comments posted before they joined.

Clarifying Questions (ask before designing)

QuestionWhy it matters
How many concurrent viewers on a peak stream?10M viewers changes fan-out from a simple broadcast to a hierarchical WebSocket cluster.
Do viewers need every comment, or is sampling OK?100K comments/sec is unreadable. Sampling ~50/sec to viewers while persisting all comments is the standard trade-off.
How strict does ordering need to be?Snowflake IDs give approximate chronological order. Perfect ordering across regions is expensive and usually unnecessary for live chat.
Sync or async moderation?Blocking on ML adds latency. Fast dictionary filter + async ML with delete_comment is the usual hybrid.

Scope

In scope

  • Post and broadcast comments in near real time
  • Hierarchical fan-out to 10M viewers
  • Server-side sampling with priority bypass
  • Late joiner history (last 100 comments)
  • Profanity filtering and rate limiting
  • Full persistence for replay after stream ends

Out of scope (state explicitly)

  • Reply threads and reactions on comments
  • Video streaming infrastructure
  • Training moderation ML models from scratch

Functional Requirements

Start by asking about moderation and ordering requirements. Confirm that the system needs comment posting, real time delivery to viewers, pagination for late joiners, spam filtering, and a decision on whether threaded replies are in scope.

In the room: ask whether deleted comments must disappear from all clients instantly or can lag briefly.

  • Post comments: Allow users to submit real time comments on active live videos or streaming events.
  • Real-time delivery: Broadcast sampled comments to active viewers within 200 ms in their serving region (human perception threshold for "instant").
  • Pinned comments: Enable hosts and moderators to pin selected comments to the top of the chat panel.
  • Flat comment stream: Comments appear in a single chronological panel (nested reply threads are out of scope for this real time design as detailed in the additional considerations section).
  • Automated moderation: Filter profanity, detect spam bots, and enforce block lists dynamically.
  • Comment rate throttling: Implement adaptive sampling to handle viral streams without degrading UI performance.
  • Historical comments: Persist comments long-term, keeping them scrollable after live events conclude.

Non-Functional Requirements

Broadcast latency under a few seconds and high write throughput during viral streams. Your interviewer will ask about ordering guarantees and toxic content at scale.

  • Low Latency: Deliver comments to millions of concurrent viewers within 200 ms in their serving region.
  • High Throughput: Design for 100K+ incoming comments per second during peak viral streams.
  • Viewer Scalability: Support up to 10 Million concurrent active connections on a single massive stream.
  • High Availability: Deliver 99.99% availability on the comment submission and viewing pipelines.
  • Approximate Ordering: Ensure comments scroll in roughly chronological order.
  • Rapid Filtering: Execute profanity checks and moderation lookups in under 50 ms.

Capacity Estimations

Comment rate during a viral live stream drives Kafka partitions, regional Redis fan-out, and WebSocket socket writes, so compute both ingress and downstream egress before sizing cluster nodes.

MetricCalculationValue
Concurrent Live StreamsGiven (peak load assumption)50,000
Peak Viewers on One StreamGiven (peak load assumption)10,000,000
Comments / sec (One Hot Stream)From Comments / day ÷ 86400 (+ peak factor in value)100,000
Comments / sec (Global)From Comments / day ÷ 86400 (+ peak factor in value)500,000
Avg Comment SizeGiven (typical workload assumption)200 bytes
Raw Global Ingress PayloadDerived100 MB/s
WebSocket Fan-Out Messages / sec10M viewers x 50 displayed comments/sec500,000,000
Approx Raw WebSocket Egress500M messages/sec x 200 bytes100 GB/s
Approx Raw Egress / 100K-Connection Gateway5M messages/sec x 200 bytes1 GB/s (~8 Gbps)
I/O and Bandwidth Calculations:
- Ingestion Data Rate: 500,000 comments/sec (global) x 200 bytes = 100 MB/s write load.
- Storage Footprint (Raw):
  - 100 MB/s x 3,600 seconds = 360 GB per hour.
  - Cassandra compression reduces this by ~3x, requiring ~120 GB per hour of persistent disk.
- WebSocket fan-out:
  - 10M viewers x 50 displayed comments/sec = 500M outbound WebSocket messages/sec.
  - 500M messages/sec x 200 bytes = 100 GB/s raw payload egress before protocol, JSON, and TLS overhead.
  - With the 100-node connection baseline, each gateway handles about 5M socket writes/sec and ~1 GB/s (~8 Gbps) raw payload at 50 comments/sec.

Architecture Diagram

We ingest comments asynchronously, moderate in a fast path, sample the live feed, and fan out the approved sampled stream through regional Redis hubs to viewers over WebSocket while persisting every comment for replay.

Loading...

Component Deep Dives

1. End-to-End Comment Lifecycle

Processing comments from the viewer's keystroke to global broadcast requires coordinated execution across low latency ingestion, filtering, fan-out, and durable persistence tiers:

  1. Comment Ingestion via WebSocket: The viewer authenticates the WebSocket session and establishes stream authorization before submitting a comment frame over the persistent connection to an ingress gateway node.
  2. Validation, Identity, and Snowflake Sequencing: The Comment Service validates text length (up to 200 bytes), assigns a 64-bit time ordered Snowflake ID to provide approximate creation time ordering, checks the author against the banned:{stream_id} Redis set in O(1) time, and enforces per-user token-bucket rate limiting (for example, maximum 5 comments per second per user).
  3. Fast Synchronous Inline Filtering: An in memory dictionary profanity check evaluates the text in under 1 millisecond. Clean messages enter the real time sampler and the Kafka publish path in parallel. The client acknowledgement waits for Kafka confirmation, while deeper machine learning toxicity scoring runs asynchronously from the Kafka stream.
  4. Windowed Dynamic Sampling: The Comment Sampler aggregates eligible comments across a 500 millisecond sliding micro window. It selects approximately 50 comments per second for broad broadcast and pushes the batch to the Redis Pub/Sub channel stream:{stream_id}:comments. Host, moderator, and pinned comments bypass sampling drops unconditionally.
  5. Edge Gateway Local Fan-Out: Across the cluster, approximately 100 WebSocket gateway servers subscribed to stream:{stream_id}:comments receive the 50 comments per second from Redis and push them to their local client sockets, with each node serving up to 100,000 active viewers.
  6. Durable Ingestion and Archival: Concurrently, the Comment Service writes every incoming comment (up to the full 100,000 comments/sec peak) to an Apache Kafka topic. The Cassandra consumer writes the complete un-sampled log for post stream replay and compliance audits.

2. Hierarchical Fan-Out Architecture (Scaling to 10M Viewers)

Directly fanning out 100 comments per second to 10,000,000 concurrent viewers requires broadcasting 1,000,000,000 messages every second from the backend. Single-node broadcast topologies immediately collapse under socket buffer exhaustion and kernel context switching. We resolve this bottleneck using a three-tier decoupled fan-out hierarchy:

Tier 1: Centralized Pub/Sub Hub: The Comment Sampler writes the filtered stream of 50 to 100 messages per second into a Redis Pub/Sub channel keyed by stream identifier (stream:{stream_id}:comments). Redis handles channel fan-out to registered backend subscribers with sub-millisecond latency.
Tier 2: Edge Gateway Subscriber Tier: A horizontally scaled pool of gateway servers subscribes to the local regional channel. The 100-server figure is the connection-based baseline for one 10M-viewer stream. Instead of 10M viewers connecting to a single point of failure, each gateway server maintains a dedicated subscription to the Redis channel, consuming only 50 to 100 messages per second over a single network socket. Total Redis cluster egress across all 100 gateway nodes is merely 5,000 to 10,000 messages per second.
Tier 3: Local Socket Broadcast: Each WebSocket gateway node manages an isolated pool of up to 100,000 persistent client TCP/TLS connections. When a message arrives from Redis, the local gateway server loops through its in memory connection registry and pushes the frame out to local sockets without incurring cross-server coordination or distributed locks.

Regional Fan-Out Note:

For globally distributed viewers, the Comment Sampler replicates only the sampled and priority stream, roughly 50 to 100 messages per second, to regional Redis hubs. Each region then fans out locally to its WebSocket gateways. The full 100,000 comments per second never crosses every regional fan-out link.

Fan-Out Capacity Summary:

  • Total concurrent viewers: 10,000,000
  • Client capacity per WebSocket gateway: 100,000 connections
  • Required gateway cluster size: 100 servers (10M / 100K)
  • Redis channel consumption per gateway: ~50 to 100 messages/sec
  • Total Redis subscriber delivery work: 100 gateway subscriptions x 100 messages/sec = 10,000 deliveries/sec, with payload size, connection count, and Redis egress monitored as practical limits

3. Dynamic Server-Side Sampling

During viral live events where incoming comments reach 100,000 per second, displaying every message in the viewer interface is both visually unreadable for humans and technically prohibitive for client network cards and mobile rendering loops. The Comment Sampler enforces adaptive server side rate control to maintain chat readability while preserving every submitted message for long-term record:

  • Uniform Random Sampling: The sampler collects incoming comments in 500ms micro-batches and extracts a uniform random sample of approximately 50 comments per second to broadcast to ordinary viewers.
  • Priority Bypass Lane: Certain classes of comments must never be dropped by the random sampler. High-priority comments bypass sampling unconditionally: host and co-host comments, channel moderators, pinned announcements, and verified creators.
  • Full Persistence Decoupling: While live broadcast pushes only 50 sampled comments per second over Redis Pub/Sub, the ingest pipeline streams all 100,000 comments per second directly into Apache Kafka. As a result, the complete unedited conversation is persisted to Apache Cassandra, ensuring full fidelity for post stream replay and moderation audits.

4. Event Bus Design and Ingestion Pipeline

Apache Kafka decouples high volume comment ingestion from downstream storage and offline intelligence pipelines. Partitioning by stream identifier preserves event order for normal streams while enabling parallel consumption across streams. A mega stream that exceeds one partition's sustainable throughput can add a shard identifier and accept approximate cross-shard ordering:

YAML
topics:
  live-comments:
    partitions: 128 # Partitioned by stream_id for normal streams, while hot streams can shard by stream_id + shard_id
    partition_key: stream_id (or stream_id + shard_id for a hot stream)
    retention: 90d # Matches Cassandra default TTL of 90 days
    replication_factor: 3
    min_insync_replicas: 2
    producer_idempotence: true
    producers:
      - comment-service # Emits event following 64-bit Snowflake ID assignment
    event_schema:
      event_id: "UUID"
      client_message_id: "UUID"
      stream_id: "UUID"
      shard_id: "int32"
      comment_id: "int64 (Snowflake ID)"
      user_id: "UUID"
      text: "string (max 200 bytes UTF-8)"
      is_pinned: "boolean"
      moderation_state: "visible | shadow_banned | deleted"
      timestamp: "int64 (epoch milliseconds)"
    consumer_groups:
      cassandra-writer:
        delivery: at_least_once
        purpose: Persists all incoming comments for durable storage and post stream replay
        deduplication: Uses the Cassandra primary key and event identity to make retries idempotent
      moderation-ml:
        purpose: Asynchronous toxicity and spam scoring, emitting delete_comment if flagged
      analytics-pipeline:
        purpose: Ingests into ClickHouse for real time engagement and viewer sentiment metrics
    dead_letter_queue: live-comments-dlq
    alert_threshold: "consumer_lag > 30s"

execution_paths:
  real_time_broadcast: "WebSocket -> Comment Service -> Comment Sampler -> regional Redis Pub/Sub -> WebSocket Gateways (< 200ms target in the serving region)"
  durable_persistence: "Comment Service -> Kafka -> Cassandra Writer Consumer -> Cassandra (all 100K comments/sec)"

API Design

Domain Model and Frame Signatures

The API consists of persistent WebSocket frames for bidirectional real time comment ingestion and batched broadcast delivery, complemented by a REST cursor pagination endpoint for late joining viewers and post stream replay. Pin and unpin control frames are accepted only after the gateway verifies host or moderator permissions for the target stream.

Strongly typed contracts define client submissions, server acknowledgements, batched push events, and historical cursor queries:

TYPESCRIPT
type StreamId = string;
type UserId = string;
type CommentId = string; // 64-bit Snowflake ID encoded as a string
type ClientMessageId = string; // UUID generated by the client and stable across retries
type UnixMillis = number;
type ModerationState = "visible" | "shadow_banned" | "deleted";
type HistoryCursor = string; // Opaque cursor encoding the Cassandra bucket position

interface PostCommentFrame {
  type: "comment";
  stream_id: StreamId;
  client_message_id: ClientMessageId;
  text: string; // UTF-8 body, max 200 bytes
}

interface CommentAckFrame {
  type: "comment_ack";
  comment_id: CommentId;
  client_message_id: ClientMessageId;
  status: "accepted" | "duplicate" | "rate_limited" | "blocked";
}

interface BroadcastCommentPayload {
  id: CommentId;
  moderation_state: ModerationState;
  user: {
    id: UserId;
    name: string;
    avatar_url?: string;
    is_host?: boolean;
    is_moderator?: boolean;
    is_verified?: boolean;
  };
  text: string;
  is_pinned: boolean;
  timestamp: UnixMillis;
}

interface BroadcastBatchFrame {
  type: "new_comments";
  stream_id: StreamId;
  comments: BroadcastCommentPayload[];
}

interface PinCommentFrame {
  type: "pin_comment";
  stream_id: StreamId;
  comment_id: CommentId;
  action: "pin" | "unpin";
}

interface ModerationControlFrame {
  type: "comment_removed";
  stream_id: StreamId;
  comment_id: CommentId;
}

interface HistoricalComment {
  comment_id: CommentId;
  user_id: UserId;
  text: string;
  is_pinned: boolean;
  created_at: string;
}

interface HistoricalCommentsResponse {
  comments: HistoricalComment[];
  next_cursor: HistoryCursor | null;
}

1. Post Comment (WebSocket Frame)

Clients submit comment frames over the persistent WebSocket connection. The Comment Service publishes the event to Kafka and returns an acknowledgement after Kafka confirms the write. The acknowledgement confirms durable admission to the event pipeline, while Cassandra materialization and downstream moderation or analytics continue asynchronously:

JSON
// Client submits comment over persistent WebSocket:
{
  "type": "comment",
  "stream_id": "live-123",
  "client_message_id": "cmsg-550e8400-e29b-41d4-a716-446655440000",
  "text": "Amazing performance! 🔥"
}

// Server returns acknowledgement after Kafka confirms the event:
{
  "type": "comment_ack",
  "comment_id": "881239589218311",
  "client_message_id": "cmsg-550e8400-e29b-41d4-a716-446655440000",
  "status": "accepted"
}

2. Receive Comments (Server Push Batch)

WebSocket gateways push batched comments to connected clients to amortize TCP frame and JSON serialization overhead:

JSON
// WebSocket gateway pushes batched comments to connected clients:
{
  "type": "new_comments",
  "stream_id": "live-123",
  "comments": [
    {
      "id": "881239589218311",
      "user": {
        "id": "u-42",
        "name": "Alice",
        "avatar_url": "https://cdn.example.com/a.png",
        "is_verified": true
      },
      "text": "Incredible sound quality!",
      "is_pinned": false,
      "timestamp": 1710320000000
    },
    {
      "id": "881239589218312",
      "user": {
        "id": "u-88",
        "name": "Bob",
        "is_moderator": true
      },
      "text": "Please keep the chat respectful everyone.",
      "is_pinned": true,
      "timestamp": 1710320001000
    }
  ]
}

3. Get Historical Comments (REST Endpoint)

Late joiners and replay viewers query historical comments using cursor-based pagination backed by Snowflake IDs:

HTTP
GET /api/v1/streams/{stream_id}/comments?before={cursor}&limit=50
Authorization: Bearer <token>

Response: 200 OK
{
  "comments": [
    {
      "comment_id": "881239589218311",
      "user_id": "u-42",
      "text": "Incredible sound quality!",
      "is_pinned": false,
      "created_at": "2026-03-13T10:00:00Z"
    }
  ],
  "next_cursor": "opaque-cursor-token"
}

Common Error Responses

Standardized error envelopes returned when client requests fail validation or exceed rate limits:

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

Cassandra: Long-Term Comments Database

Apache Cassandra serves as the durable comment store. Partitioning by stream identifier, a five minute UTC bucket, and a deterministic hash shard spreads hot streams across bounded partitions. Clustering descending on the 64-bit Snowflake ID enables reverse chronological range scans within each shard, while historical reads merge the relevant shards under an automatic 90 day time to live:

SQL
-- Durable comments storage in Apache Cassandra
-- bucket_start: UTC 5-minute window
-- bucket_shard: hash(comment_id) % 16
-- At 100K comments/sec, one shard contains about 1.875M rows per 5-minute bucket before TTL expiry
CREATE TABLE comments (
    stream_id         UUID,
    bucket_start      TIMESTAMP,
    bucket_shard      SMALLINT,
    comment_id        BIGINT,
    user_id           UUID,
    text              TEXT,
    is_pinned         BOOLEAN,
    moderation_state  TEXT,       -- visible, shadow_banned, deleted
    client_message_id UUID,
    created_at        TIMESTAMP,
    PRIMARY KEY ((stream_id, bucket_start, bucket_shard), comment_id)
) WITH CLUSTERING ORDER BY (comment_id DESC)
  AND default_time_to_live = 7776000;  -- 90 days

-- Durable retry identity backstop
CREATE TABLE comment_idempotency (
    stream_id         UUID,
    user_id           UUID,
    client_message_id UUID,
    comment_id        BIGINT,
    event_id          UUID,
    created_at        TIMESTAMP,
    PRIMARY KEY ((stream_id, user_id), client_message_id)
) WITH default_time_to_live = 86400;  -- 24-hour retry horizon

ClickHouse: Analytics Store

ClickHouse receives asynchronous analytics events for engagement and viewer sentiment. It supports analytical queries and dashboards, but it is not the source of truth for comment durability and it is not used on the live WebSocket delivery path.

Redis Topology and In-Memory Keyspaces

Redis acts as the low latency pub/sub message bus and ephemeral state cache. Channels distribute sampled comments to WebSocket gateways, while dedicated keyspaces store banned user rosters and recent comment ring buffers:

YAML
channels:
  live_broadcast:
    pattern: "stream:{stream_id}:comments"
    type: "Pub/Sub Channel"
    purpose: Broadcasts sampled and priority comments to subscribed WebSocket gateways

keyspaces:
  banned_users:
    key: "banned:{stream_id}"
    type: "SET"
    purpose: Provides O(1) membership checks to filter hard-blocked users before broadcast
    example_members: ["user-921", "user-331"]

  shadow_banned_users:
    key: "shadow_banned:{stream_id}"
    type: "SET"
    purpose: Marks users whose comments are persisted and acknowledged to the author but never broadcast to other viewers
    example_members: ["user-777"]

  recent_comments_buffer:
    key: "recent_comments:{stream_id}"
    type: "LIST"
    purpose: Capped serving buffer retaining the last 100 broadcast comments for late joining viewers
    lifecycle: "Capped via LPUSH + LTRIM with 24-hour expiration after stream ends"

  comment_dedupe:
    key: "comment_dedupe:{stream_id}:{user_id}:{client_message_id}"
    type: "STRING"
    purpose: Maps a client retry token to the originally assigned comment_id and event_id
    lifecycle: "Short TTL fast path, backed by Cassandra comment_idempotency as the durable backstop for the 24-hour retry horizon"

  moderation_tombstones:
    key: "moderation_tombstones:{stream_id}"
    type: "SET"
    purpose: Prevents deleted comments from reappearing from the recent buffer or reconnect replay
    lifecycle: "Retained through the active stream plus a short post stream grace period"

Gateways and sampling workers execute pipeline-friendly O(1) Redis commands to push updates, evaluate moderation block lists, and maintain the late-joiner buffer:

REDIS
# 1. Publish sampled comment batch to stream channel
PUBLISH stream:live-123:comments '{"type":"new_comments","comments":[{"id":"881239589218311","text":"Great show!"}]}'

# 2. Check if a commenting user is banned from the stream (O(1))
SISMEMBER banned:live-123 "user-921"

# 3. Check if a user is shadow-banned from the stream (O(1))
SISMEMBER shadow_banned:live-123 "user-777"

# 4. Buffer recent comments for late joiners (capped at 100 entries)
LPUSH recent_comments:live-123 '{"id":"881239589218311","text":"Great show!"}'
LTRIM recent_comments:live-123 0 99

# 5. Resolve a client retry to the original comment ID
GET comment_dedupe:live-123:user-921:cmsg-550e8400-e29b-41d4-a716-446655440000

# 6. Record a moderation tombstone for a deleted comment
SADD moderation_tombstones:live-123 "881239589218311"

Fault Tolerance

Concern ScenarioMitigation Strategy
WebSocket Gateway CrashClients reconnect with jittered backoff and supply their last seen Snowflake ID. The gateway subscribes first, buffers live frames briefly, replays the retained Redis gap, and then drains the live buffer with client side deduplication. If the gap exceeds the retained window, the gateway uses Cassandra history.
Moderation Pool LatencyFall back gracefully to synchronous in memory pre-filters, publishing clean comments instantly while dispatching deep ML evaluation asynchronously via Kafka.
Redis Pub/Sub Connection DropBecause Pub/Sub is ephemeral, surviving gateways serve recent history buffers while missing comments are fetched from Cassandra if necessary.

1. The Late Joiner Catch-Up Problem

When a user opens a live video stream that has been broadcasting for an hour, presenting an empty comment panel feels jarring and disconnected. However, querying the persistent Cassandra database on every WebSocket connection would generate an unsustainable read surge during viral viewer influxes.

We solve this bootstrap challenge using a dedicated Redis circular buffer cache. For every active broadcast, Redis maintains a serving list keyed as recent_comments:{stream_id} capped at 100 items using atomic LPUSH and LTRIM operations.

When a new client establishes a WebSocket connection:

  1. The gateway registers the client socket with the stream's Redis Pub/Sub channel first and temporarily buffers live frames.
  2. The gateway fetches the last 100 comments from the Redis serving list and sends the available backfill.
  3. The gateway drains the temporarily buffered live frames. The client application deduplicates by the unique 64-bit comment identifier so the handoff does not create duplicate visible comments.

2. Moderation Pipeline Latency and Shadow-Banning

Filtering toxic content, hate speech, and spam bots cannot introduce perceptible delays into legitimate user conversations. A purely synchronous machine-learning evaluation pipeline adds 50 to 100 milliseconds of latency per comment and causes entire chat feeds to stall if the moderation inference cluster experiences a spike in queue depth.

We employ a three-tier hybrid moderation architecture:

  • Fast Synchronous Pre-Filter: An in memory dictionary and regex hash set evaluates incoming comments within 1 millisecond on the Comment Service. Obvious profanity and known abusive patterns are rejected instantly before the comment ever reaches the broadcast or persistence layers.
  • Asynchronous Deep Scoring: Clean comments publish immediately to the live broadcast sampler, while simultaneously being written to the live-comments Kafka topic. A pool of background machine learning moderation consumers ingests from Kafka to score contextual toxicity, spam similarity, and user reputation. If a comment is flagged after broadcast, the moderation service records the durable deleted state and a Redis tombstone before emitting a delete_comment control frame over Redis Pub/Sub. Gateways remove the offending comment ID from viewer viewports within approximately 2 seconds after receiving the control frame.
  • User Shadow-Banning: When abusive users or spam bots are detected, they are placed on a shadow-ban list stored in Redis (shadow_banned:{stream_id}). Comments submitted by a shadow-banned account return an immediate success acknowledgement to the author, and their client interface displays the comment normally. The comment is still persisted with moderation state set to shadow_banned, but the Comment Service does not forward it to Redis Pub/Sub or other viewers. The author still sees the local result, which makes the restriction less obvious to the abusive account.

3. Ephemerality vs Strict Comment Ordering

In a globally distributed streaming platform, network propagation delays between diverse geographical edge locations mean packets inevitably arrive out of physical order at backend ingress gateways. Attempting to enforce globally synchronized total order across millions of distributed clients would require expensive consensus protocols that ruin real time performance.

Because live chat is ephemeral and fast-scrolling, human viewers cannot distinguish minor timing discrepancies on the order of 50 to 100 milliseconds. We balance ordering and performance using Snowflake identifiers:

  • Snowflake ID Generation: When the Comment Service receives a comment, it stamps the record with a 64-bit Snowflake ID generator containing a 41-bit millisecond timestamp, 10-bit machine identifier, and 12-bit sequence number. This provides approximate creation time ordering without requiring distributed synchronization.
  • Client-Side Micro-Buffering: During critical scenarios such as live host Q&A sessions where sequential causality matters, client applications buffer incoming frames for 100 to 200 milliseconds, sorting items by Snowflake identifier in memory before appending them to the rendered chat container. This eliminates out of order jitter caused by variable network transit times without stalling the visual presentation.

4. WebSocket Gateway Failure and Thundering Herd Mitigation

A single WebSocket gateway node accommodates up to 100,000 persistent viewer connections. If that gateway crashes or experiences an abrupt hardware fault, 100,000 clients disconnect simultaneously. If all clients instantly attempt to reconnect to the surviving cluster, the resulting reconnect storm can cascade across remaining gateway instances.

We mitigate reconnect storms and ensure uninterrupted chat continuity through three synchronized mechanisms:

  • Exponential Backoff with Jitter: Disconnected client applications apply an exponential backoff retry strategy with randomized jitter between 0 and 500 milliseconds. Spreading connection attempts across this window flattens the reconnection spike on edge load balancers.
  • Resumption Cursor Handoff: When reconnecting, the client includes its last_seen_comment_id in the WebSocket connection query parameters.
  • Server-Side Gap Replay: The new WebSocket gateway first subscribes the socket to the live Redis Pub/Sub channel and temporarily buffers those live frames. It then reads the Redis recent comments buffer (recent_comments:{stream_id}) for entries newer than last_seen_comment_id, sends the available gap, and drains the temporarily buffered live frames with client side deduplication. If the requested cursor is older than the retained Redis buffer, the gateway falls back to historical Cassandra pagination instead of assuming Redis can replay an unbounded gap. This ordering avoids the race between backfill and live subscription.

Additional Considerations

Reply Threads: Deferred Staff Extension

Nested reply threads are listed in many product specifications but conflict with mega-stream sampling, because a thread with 50,000 replies still cannot be pushed to every viewer without saturating client interfaces. When scoping this in an interview, treat threads as a post stream video-on-demand feature backed by Cassandra with parent_comment_id clustering and paginated REST reads, rather than placing them on the live WebSocket broadcast path. Live chat remains flat, keeping threaded conversations as a separate, lower-traffic surface after the event concludes.

Verified User and Host Prioritization

To maintain high-quality streaming discussions, comments originating from verified accounts, moderators, or the stream hosts themselves are tagged with priority flags. The sampling engine routes these comments into a dedicated high-priority queue, guaranteeing they bypass random sampling drops entirely.

Interview Walkthrough

  • 25-minute cut

    Skip arch50/arch75 depth unless staff.

    • Frame the throughput mismatch: 100,000 comments/sec incoming, viewers need batched delivery rather than per-comment pushes (5 minutes).
    • Propose a tiered pipeline: ingest, Kafka queuing, stream sampling, and WebSocket fan-out to edge gateways (6 minutes).
    • Compare Redis Pub/Sub (low latency, ephemeral) vs Apache Kafka (durable, replayable) for the broadcast and persistence layers (5 minutes).
    • WebSocket fan-out through edge servers rather than pushing from a single origin to millions of client sockets (5 minutes).
    • Apply per-user rate limits and asynchronous moderation to prevent bot floods during viral streams (4 minutes).
  • Frame the problem as a throughput mismatch: 100,000 comments/sec incoming traffic while viewers can only comfortably consume approximately 10 to 50 comments/sec, making server side sampling mandatory rather than optional.
  • Propose a tiered delivery pipeline consisting of ingest, filtering and sampling, and edge broadcast, with a dedicated high priority bypass lane for hosts and verified accounts.
  • Compare Redis Pub/Sub (low latency, fire and forget) vs Apache Kafka (durable, replayable) for the broadcast layer and justify your architectural choice.
  • Use WebSocket fan-out through edge servers, avoiding direct pushes from a single origin to millions of clients (similar to patterns examined in Real-Time Chat System).
  • Apply per-user rate limiting to prevent bot floods from starving legitimate viewers in the sampling pool.
  • Persist all comments via Event Sourcing and CQRS even if only a sample is broadcast, because replay and moderation require the full log.
  • Quantify with Back-of-the-Envelope Estimation: 10M viewers x 50 comments/sec displayed = 500M WebSocket messages/sec at peak across the edge fleet.
  • Common pitfall: attempting to deliver every comment to every viewer on a viral stream, which causes cluster collapse at the fan-out layer.

Engineering Trade-offs

1. Redis Pub/Sub vs Apache Kafka for Real Time Broadcasting

Balancing delivery speed against durability, strict ordering against throughput, and safety against latency dictates the core architectural trade offs of live streaming chat.

Redis Pub/Sub provides sub-millisecond in memory dispatch and lightweight channel subscriptions without disk persistence overhead. However, it is fire and forget. Disconnected subscribers immediately miss dropped messages and cannot rewind history. Conversely, Apache Kafka offers durable, multi-subscriber append-only logs with configurable retention, replayability, and consumer group offset management, but introduces 5 to 15 milliseconds of disk sync and broker batching latency, which adds unnecessary overhead on real time interactive chat feeds.

Architecture Decision: We adopt a hybrid integration. Redis Pub/Sub powers the ephemeral real time regional fan-out tier to WebSocket gateways to achieve sub-200ms broadcast latency in each serving region. Simultaneously, all incoming comments route to an Apache Kafka topic for durable ingestion into Cassandra, offline moderation, and long-term analytics.

2. Synchronous Inline Moderation vs Asynchronous Post-Broadcast Auditing

Synchronous moderation inspects every comment before broadcast, ensuring that abusive or violating content never reaches viewer viewports. However, deep neural language models require 50 to 150 milliseconds of inference latency. Under viral stream peaks of 100,000 comments per second, GPU inference pools saturate, creating backpressure that blocks clean comments and degrades chat responsiveness. Asynchronous moderation publishes comments immediately to viewers and dispatches evaluation jobs in the background. While this guarantees sub-50ms ingestion latency and protects system availability, viewers may be exposed to toxic content for 1 to 2 seconds until an asynchronous delete event arrives.

Architecture Decision: We implement a tiered compromise. An in memory dictionary and regex rule engine executes synchronous pre-checks in under 1 millisecond. Clean comments are broadcast immediately and enqueued into Kafka for deep asynchronous machine learning evaluation. If an asynchronous model flags a comment, a deletion frame retracts it from all connected clients, combining instantaneous delivery with comprehensive safety.

3. Strict Global Total Ordering vs Approximate Snowflake Clustering

Strict global total ordering requires all comments across all regions to pass through a single coordinator or consensus group (such as Raft or a centralized database sequencer). This introduces cross-region round-trip network latency and forms a severe throughput bottleneck that cannot scale to 100,000 writes per second. Approximate ordering generates 64-bit time-ordered Snowflake ID generator tokens at the ingest tier. Comments retain millisecond timestamp ordering, but network jitter can cause packets from different regions to arrive with minor millisecond deviations.

Architecture Decision: Live streaming chat is ephemeral and scrolls rapidly, meaning viewers cannot perceive 50 to 100 millisecond ordering variations between disparate commenters. We adopt approximate ordering via Snowflake identifiers for live broadcasts, coupled with Cassandra clustering ordering by comment_id DESC for deterministic, chronological post stream replay.

4. 100% Broadcast Delivery vs Dynamic Server-Side Sampling

Delivering 100% of comments during a 100,000 comments per second mega-stream to 10,000,000 viewers would require transmitting 1,000,000,000,000 messages every second across the gateway fleet. This instantly saturates edge network bandwidth, exhausts client device battery and CPU, and causes the chat interface to blur into an illegible cascade. Server-side sampling caps client delivery at approximately 50 comments per second, maintaining an engaging and readable chat velocity while dramatically slashing egress bandwidth. However, users whose comments are sampled out may feel ignored if not handled transparently.

Architecture Decision: We enforce server side sampling with a dedicated priority lane. The Comment Sampler publishes a representative sample of 50 comments per second to Redis Pub/Sub, but unconditionally prioritizes hosts, moderators, pinned comments, and verified users. Meanwhile, all 100,000 comments per second are written durably to Cassandra, guaranteeing that complete chat transcripts remain available for creator reviews, search indexes, and on-demand VOD replays.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...