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)
| Question | Why 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.
| Metric | Calculation | Value |
|---|---|---|
| Concurrent Live Streams | Given (peak load assumption) | 50,000 |
| Peak Viewers on One Stream | Given (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 Size | Given (typical workload assumption) | 200 bytes |
| Raw Global Ingress Payload | Derived | 100 MB/s |
| WebSocket Fan-Out Messages / sec | 10M viewers x 50 displayed comments/sec | 500,000,000 |
| Approx Raw WebSocket Egress | 500M messages/sec x 200 bytes | 100 GB/s |
| Approx Raw Egress / 100K-Connection Gateway | 5M messages/sec x 200 bytes | 1 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.
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:
- 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.
- 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). - 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.
- 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. - Edge Gateway Local Fan-Out: Across the cluster, approximately 100 WebSocket gateway servers subscribed to
stream:{stream_id}:commentsreceive 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. - 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:
stream:{stream_id}:comments). Redis handles channel fan-out to registered backend subscribers with sub-millisecond latency.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:
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:
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:
// 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:
// 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:
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:
-- 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 horizonClickHouse: 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:
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:
# 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 Scenario | Mitigation Strategy |
|---|---|
| WebSocket Gateway Crash | Clients 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 Latency | Fall back gracefully to synchronous in memory pre-filters, publishing clean comments instantly while dispatching deep ML evaluation asynchronously via Kafka. |
| Redis Pub/Sub Connection Drop | Because 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:
- The gateway registers the client socket with the stream's Redis Pub/Sub channel first and temporarily buffers live frames.
- The gateway fetches the last 100 comments from the Redis serving list and sends the available backfill.
- 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-commentsKafka 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 adelete_commentcontrol 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_idin 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 thanlast_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
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.