Real-Time and Streaming System Design Questions
Designing client-facing, always-on live delivery systems: real-time communication transports (WebSockets, server-sent events, long-polling, WebRTC, MQTT), connection lifecycle and scaling for millions of persistent connections, presence, pub/sub fan-out to online users, chat and notifications, live feeds and tickers, real-time collaboration (CRDT vs OT, offline sync and reconciliation), and live and on-demand video delivery (ingest, transcoding, adaptive bitrate, CDN delivery, low-latency protocols, playback entitlement and content protection). Covers latency budgets, per-client ordering, delivery guarantees across disconnect and reconnect, per-client backpressure, capacity estimation for connection fleets, authentication, authorization, and revocation for real-time channels and premium content, and operational readiness (SLOs and error budgets, observability and incident response, rate limiting, safe rollout, and load and chaos testing) for live-delivery platforms. Data-pipeline stream processing (Kafka or Flink jobs, windowed aggregation, exactly-once pipelines, real-time analytics ingestion) is out of scope.
Design an architecture to support 100k user-generated live channels, most with small audiences. Discuss cost-efficient transcoding approaches (per-channel vs shared encoder pools), segment retention policies, cold/warm origin caching, and autoscaling strategies for sudden popularity spikes.
Sample Answer
Direct answer
The defining fact is the audience shape: of 100,000 live channels, most have a handful of viewers, and a few suddenly have thousands. So I would not transcode by default. Small channels are delivered as the broadcaster's own stream, only repackaged into segments (transmuxing, which changes the container without re-encoding and is cheap). Channels get a reduced or full quality ladder from a shared encoder pool only when their viewer count crosses thresholds. Retention is a short live window in fast storage, with archive only for channels that opt in. Delivery uses a tiered CDN (content delivery network: a geographically distributed set of edge servers that cache content close to viewers, instead of every viewer fetching from your origin) with a shield layer that collapses duplicate requests, because for a 3-viewer channel the edge cache barely helps. Popularity spikes are absorbed by the CDN automatically and by a warm encoder pool that can upgrade a channel in seconds.
Requirements and assumptions
Stated as assumptions so the arithmetic can be redone with real numbers:
- 100,000 concurrent live channels at peak, each sending about 6 Mbit/s (1080p) of source.
- Viewer distribution: 70% of channels have 0 to 2 viewers, 25% have 3 to 49, 5% have 50 or more; about 1 million concurrent viewers in total, averaging 3 Mbit/s each.
- Target latency of a few seconds (2-second segments), viewer-count-driven quality.
Capacity arithmetic
ingestviewer egressone channel-hour of sourceall channels, one hour=100,000×6 Mbit/s=600 Gbit/s=1,000,000×3 Mbit/s=3 Tbit/s=86×106×3600=2.7×109 bytes=2.7 GB=100,000×2.7 GB=270 TBAt 2-second segments, the origin produces 100,000 / 2 = 50,000 new segments per second (per rendition). Viewer egress (3 Tbit/s) is the largest number here, but it is served from the CDN's edge fleet, which scales with total viewer traffic more or less automatically regardless of how any single channel is encoded or retained; nothing in this section changes how many bytes an edge sends to viewers. What the per-channel choices in this section (transcoding tier, retention window, caching) actually change is how much gets written and stored at the origin, and how often the origin gets asked for something new. That is why storage and origin write rate, not viewer egress, is what the per-channel design choices control.
Transcoding: per-channel encoders vs shared pools
Per-channel encoders (a dedicated encoder process or machine per stream) are simple and isolate failures, but allocate capacity for a full ladder to channels nobody watches. Shared encoder pools schedule many channels' encode jobs onto a pool of GPUs or hardware encoders, bin-packing by rung (grouping jobs of the same quality level onto the fewest machines that fit them, the way you would pack boxes onto trucks to minimize how many trucks you need), so capacity follows demand.
Commit: shared pools, with tiered transcoding by viewer count. Measure cost in "encode units", where one unit is one low rung (for example 480p) and a full 4-rung ladder is about 4 units (an assumption; measure it on your hardware):
| Channel tier | Viewers | Treatment | Units per channel |
|---|---|---|---|
| Passthrough | 0 to 2 | Transmux source only | 0 |
| Light | 3 to 49 | Source plus one 480p rung for weak connections | 1 |
| Full | 50+ | Full ABR ladder (adaptive bitrate: several qualities the player switches between) | 4 |
The cost of this choice: a viewer on a weak connection watching a passthrough channel gets only the source bitrate and may buffer. That is the right trade for a 2-viewer channel; the light tier fixes it once anyone else is watching.
Hysteresis: upgrade at 50 viewers, downgrade only after the count stays below 30 for 10 minutes, so a channel hovering near the threshold does not flip encoders every minute.
Mid-stream upgrade caveat: in HLS (HTTP Live Streaming, the most common segment-based delivery protocol) the player reads the list of available qualities (the multivariant playlist) when it starts. New viewers see the new ladder immediately; existing viewers stay on the source until their player reloads that playlist. Either accept that (the spike is mostly new viewers anyway) or have the player re-fetch the multivariant playlist periodically.
Segment retention
- Live window: keep the last few minutes per channel (enough for a short rewind, say 5 minutes) on origin SSD or memory. 5 minutes of 6 Mbit/s source is 225 MB per channel, 22.5 TB across all channels: fine in a fleet's memory and SSD.
- Archive (VOD of past broadcasts): only for channels that opt in or pass a popularity threshold, written to object storage and deleted after a retention period (for example 14 days) unless promoted. Archiving everything for 14 days would be 270 TB x 24 x 14, about 90 PB; archiving 10% of channels is about 9 PB, which is why archiving is opt-in.
- Store the source rendition only for archives and re-transcode on demand if an old broadcast becomes popular.
Cold and warm origin caching
flowchart LR
B[Broadcasters] --> ING[Ingest and transmux]
ING --> POOL[Shared encoder pool]
ING --> ORG[Origin: live window]
POOL --> ORG
ORG --> SH[CDN shield tier]
SH --> E1[Edge PoP A]
SH --> E2[Edge PoP B]
E1 --> V1[Viewers]
E2 --> V2[Viewers]
- The small-channel problem: a channel with 3 viewers in 3 cities hits 3 different edge locations (PoPs, points of presence), each missing cache and each asking the origin. Edge caching gives almost nothing for the long tail.
- Shield tier (a mid-layer cache between edges and origin) plus request coalescing (many simultaneous requests for the same segment become one origin fetch) caps origin load at about one fetch per segment per shield region, regardless of the edge count.
- Warm vs cold: popular channels' segments are hot in edge caches; long-tail channels are cold everywhere and served by the shield. Short cache TTLs (time-to-live) on playlists (a second or two, they change every segment) and long TTLs on segments (they never change once written).
- Origin placement: route each channel to an origin in the region its broadcaster ingests into, and consistent-hash channels across origin nodes (assign each channel to a node using a hash function built so that removing or adding one node only reshuffles the small slice of channels it owned, instead of remapping every channel) so a node failure moves only its channels.
Autoscaling for sudden popularity spikes
A small channel gets "raided" (a large streamer sends their audience over) and goes from 3 to 20,000 viewers in a minute.
- Delivery scales by itself: the CDN absorbs it, and the shield's request coalescing means 20,000 viewers still produce about one origin fetch per segment per shield.
- Transcode upgrade is triggered by viewer-count events (not a periodic job), pulling encoders from a warm pool: pre-started encoder capacity held idle for this purpose, sized to the peak number of upgrades per minute. A cold start of a new machine is too slow for a raid.
- Pool autoscaling follows pool utilisation with headroom (for example keep 20% idle), plus scheduled scale-up before known peak hours.
- Degradation order when the pool is exhausted: new upgrades get the light tier instead of the full ladder; no running full ladder is ever torn down to make room.
Trade-offs and pitfalls
- Transcoding everything is the default design and the most expensive mistake at this scale (the 88.75% in the example is the size of that mistake under these assumptions).
- Passthrough exposes broadcaster settings directly to viewers: an odd keyframe interval from broadcasting software breaks segmenting. Enforce ingest rules (keyframe every 2 seconds) and reject or fix non-conforming streams at ingest.
- Per-channel cost floors (ingest connection, origin memory, monitoring series) matter more than encode cost for the long tail; avoid per-channel metrics labels (a metrics system stores one time series per distinct label value, so a label like
channel_idon every metric multiplies your stored series by 100,000, one per channel, which is expensive and mostly useless to query). - Viewer counts can be gamed to trigger free transcodes; count authenticated, playing sessions, not page loads.
Design an architecture to handle 10 million concurrent persistent connections globally across multiple regions, supporting failover, low latency, and broadcasting messages to any subset of users, with regional latency SLOs under 500ms. Discuss regional affinity, global routing (GeoDNS/GSLB), sticky-session implications, edge proxies, how connection state is stored, consistency of presence state, cross-region replication of pub/sub metadata, and failover design.
Sample Answer
Direct answer
I would run the same stateless connection tier in four regions, route each client to its nearest healthy region with latency-based DNS or anycast, and treat every connection as disposable: any gateway in any region can accept any user, and a reconnect resumes from a durable per-user sequence number rather than from server memory. Within a region, a registry maps each user to the gateway holding their socket, and a regional pub/sub bus moves messages between gateways. Across regions, only coarse metadata replicates (which regions have subscribers for which channels, and presence deltas), so a regional publish crosses the ocean only when someone over there is listening. Each region is sized so the other three can absorb it if it fails.
Terms: GeoDNS answers a DNS query with an address chosen by the client's location; GSLB (global server load balancing) is the same idea plus health checks, so a dead region stops being handed out; anycast advertises one IP address from many locations and lets internet routing deliver packets to the nearest. Presence is the "who is online" state. Pub/sub (publish/subscribe) is messaging where senders publish to a named channel and every subscriber receives a copy.
Requirements and numbers
- 10 million concurrent persistent connections (WebSocket), global.
- 1 billion messages per day, delivered to arbitrary subsets of users.
- Regional latency SLO (service-level objective): sender to recipient in the same region under 500 ms.
- Survive the loss of a whole region.
With four regions, each carries 2.5 million connections normally. If one region fails, the other three must hold 10 million / 3, about 3.33 million each. At an assumed 100,000 connections per gateway (a sizing to validate by load test), that is 34 gateways per region, 136 in total, against 100 needed with no failure headroom.
Architecture
flowchart TB
U[Clients] --> D[GeoDNS or anycast GSLB]
D --> EA[Region A edge TCP load balancer]
D --> EB[Region B edge TCP load balancer]
EA --> GA[Gateways A]
EB --> GB[Gateways B]
GA --> RA[(Registry and presence A)]
GA --> BA[Regional pub/sub A]
GB --> BB[Regional pub/sub B]
BA <-->|interest-filtered bridge| BB
GA --> S[(Durable message store)]
GB --> S
Each design question
Regional affinity
A client connects to its nearest region and stays there for the life of the socket. Affinity is a performance preference, not a correctness requirement: nothing breaks if a user in Paris lands in Virginia after a failover; they just pay more latency. User data can still have a home region for storage, but the connection tier does not depend on it.
Global routing (GeoDNS / GSLB)
- Latency-based or geo DNS with health checks and a short TTL (time-to-live), such as 30 to 60 seconds. DNS failover is only as fast as clients honour the TTL, and some ignore it, so it is the coarse layer.
- Anycast in front of it (a managed anycast entry point a cloud provider runs for you, often marketed as a "global accelerator", or your own anycast edge) moves traffic off a dead region at the routing layer within seconds, without waiting for DNS caches.
- Recommendation: anycast entry points if the platform offers them, with GeoDNS as the fallback control. Either way, a failover means existing sockets break and clients reconnect, so the reconnect path (below) is the real failover mechanism.
Sticky-session implications
A WebSocket is inherently sticky: the TCP connection stays on one gateway until it closes. The question is whether a reconnect must return to the same gateway. My answer is no. If reconnects needed the same server, a gateway crash would strand its users and rolling deploys would need cookie routing. Instead, state needed to resume lives outside the gateway: the client sends last_seq on reconnect and the new gateway replays from the durable store. So the edge load balancer can be a plain L4 (TCP-level) balancer with no session affinity.
Edge proxies
L4 load balancers per region pass TCP through; TLS terminates on the gateways (or on a thin TLS-terminating proxy tier if you want certificate handling off the gateways). Two reasons to prefer L4: an L7 proxy (one that reads and understands the HTTP or WebSocket framing itself rather than just passing bytes through, the way the L4 balancer above does) terminating 10 million sockets doubles the socket count (one client side, one backend side), and every proxy restart would drop connections. Edge proxies do enforce the cheap protections: per-IP connection caps, handshake timeouts and idle timeouts aligned with application heartbeats (for example 30-second pings under a 60-second idle timeout).
How connection state is stored
- In gateway memory (ephemeral): the socket, its outbound queue, its subscription list. Lost on crash, by design.
- Regional registry (Redis or similar, replicated within the region):
user_id -> {gateway_id, connection_id}with a TTL refreshed by heartbeat, andchannel -> set of gateways with subscribers. Used to route a message to the right gateway. - Durable store (replicated database or log): per-user or per-conversation sequenced messages for replay on reconnect.
Sending to "any subset of users": for a list of user IDs, look the users up in the registry, group them by gateway, and send one batch per gateway. For channel-based audiences, publish once to the channel and let only gateways with subscribers receive it.
Consistency of presence state
Presence is eventually consistent, and I would say so explicitly in the design review. Per region, presence is a set of live connection IDs per user with a heartbeat TTL (for example 60 seconds). Across regions, each region publishes presence deltas asynchronously, and a user is "online" if any region holds a live connection for them. Modelling it as a set of connection IDs (rather than a single online/offline flag) means two regions changing it concurrently cannot overwrite each other: each region only ever adds its own connection IDs to the set or removes them, it never replaces the whole value, so merging two regions' updates is just a union of adds and removes and lands on the same result no matter what order the updates arrive in. A single online/offline flag has no such property: if region A's copy flips to offline at the same moment region B's copy flips to online, whichever write is applied last wins and silently erases the other. Concretely: a user on a phone in region A and a laptop in region B shows online until both sets are empty. Staleness of a few seconds is acceptable for presence; strong global consistency would cost cross-region round trips on every status change, for no user-visible benefit.
Cross-region replication of pub/sub metadata
Replicate interest, not messages: each region advertises the channels it has at least one subscriber for. When region A publishes to a channel, it forwards to region B only if B has declared interest. Interest changes rarely (when the first subscriber in a region joins or the last leaves), so this is low volume. Message payloads that must be durable go to the durable store, which has its own replication; the live cross-region bridge is best-effort, and the reconnect replay covers anything the bridge drops.
Latency budget for the 500 ms SLO
These are allocations to design against, not measurements:
| Hop | Budget |
|---|---|
| Sender to gateway (in-region network plus the TLS encryption overhead on an already-open connection, not a new handshake) | 50 ms |
| Gateway validation, auth check, sequencing write (writing the message to the durable store with its per-user sequence number, the state a later reconnect resumes from) | 50 ms |
| Registry lookup and regional pub/sub hop | 50 ms |
| Recipient gateway queue and send | 50 ms |
| Gateway to recipient network | 50 ms |
| Reserve for p99 (99th percentile) tail and GC (garbage collection) pauses, the stop-the-world moments a managed-memory runtime spends reclaiming memory | 250 ms |
Cross-region delivery adds the inter-region one-way latency, which is why the SLO is stated as regional.
Failover design
When region A dies, its 2.5 million clients disconnect together. Routing moves them elsewhere, and they reconnect with exponential backoff and jitter:
spread over 60 sspread over 120 s:2,500,000/60≈41,667 new connections per second:2,500,000/120≈20,833 new connections per secondAcross three surviving regions at 120 seconds, that is about 6,944 TLS handshakes per second per region, about 204 per gateway across 34 gateways. The design pieces: pre-provisioned headroom (the 3.33 million per region sizing), admission control (the gateway itself refusing new connections once it is at capacity, with a hint to retry later, instead of accepting every attempt and failing under load) that rejects with a retry hint when a gateway is above its handshake rate, and clients that resume from last_seq so no message is lost, only delayed.
Trade-offs and pitfalls
- What would change the design: if messages mostly go to one huge audience (a live event) the problem becomes broadcast fan-out and a hierarchical relay tree (a small number of relay nodes each fanning out to many more below them, the way a CDN distributes one piece of content, instead of one origin sending to every one of millions of users directly) beats per-user routing; if data residency requires EU users to stay in the EU, regional affinity becomes a hard rule and failover must stay inside the jurisdiction.
- Global strongly consistent presence is the classic over-build. Nobody notices a friend's dot turning green two seconds late.
- Sizing only for steady state is the classic under-build. The reconnect storm after a region loss is the peak you must survive.
- Replicating every message to every region multiplies cross-region egress (traffic leaving a region outbound, the direction cloud providers usually charge for) by the region count for messages nobody there wants. Filter by interest.
Compare HLS, DASH, and WebRTC for different streaming use cases, and recommend which protocol you'd choose for: VOD on heterogeneous devices, large-scale live events, and interactive real-time apps (e.g. auctions or low-latency chat). Consider latency, compatibility, CDN-friendliness, and DRM integration in your recommendations.
Sample Answer
Direct answer
HLS and DASH are the same idea (chop video into small HTTP files and list them in a manifest), so they ride on ordinary CDN caching and scale cheaply, at the cost of seconds of latency. WebRTC is a different idea (push media packets continuously over UDP to each viewer through media servers), which gets latency under a second but gives up HTTP caching and standard DRM. My picks: VOD on heterogeneous devices: CMAF packaged once, served as both HLS and DASH. Large-scale live events: low-latency HLS and DASH (LL-HLS, LL-DASH) over a CDN. Interactive real-time apps such as auctions: WebRTC for the media, with bids on a separate server-authoritative channel.
The three protocols in plain terms
- HLS (HTTP Live Streaming): Apple's format. A text playlist (
.m3u8) lists short media segments; the player downloads them one after another over HTTP and picks between several bitrate versions (adaptive bitrate, ABR). - DASH (Dynamic Adaptive Streaming over HTTP, MPEG-DASH): the ISO standard with the same model and an XML manifest (the MPD, Media Presentation Description).
- CMAF (Common Media Application Format): a shared fragmented-MP4 segment format (an MP4 file structured as small, independently parseable pieces, instead of one block that must be complete before playback can start) that both HLS and DASH can point at, so one set of media files serves both.
- WebRTC (Web Real-Time Communication): browser-native real-time media over UDP (User Datagram Protocol: packets sent without the connection setup, ordering or automatic retransmission that TCP, and so HTTP, guarantees; a lost packet is simply gone unless the application recovers it, trading reliability for speed) (RTP, the Real-time Transport Protocol, encrypted as SRTP, its secure variant), with congestion control (adjusting how much is sent based on measured network conditions, to avoid overwhelming the path) that drops quality rather than adding delay. Large audiences are served through an SFU (selective forwarding unit), a server that receives one stream and forwards copies to each viewer.
- DRM (digital rights management): encryption plus a licence server so only entitled players can decrypt. The three systems are Apple FairPlay, Google Widevine, Microsoft PlayReady; browsers reach them through EME (Encrypted Media Extensions).
Comparison on the four asked axes
| Axis | HLS | DASH | WebRTC |
|---|---|---|---|
| Latency | Classic: roughly 15 to 30 s behind live. LL-HLS: roughly 2 to 5 s | Classic: similar to HLS. LL-DASH (chunked CMAF): roughly 2 to 5 s | Sub-second, typically a few hundred ms |
| Compatibility | Native on Safari, iOS, iPadOS, tvOS; elsewhere via an MSE-based JavaScript player (hls.js, or Shaka Player, Google's open-source JavaScript player for the same purpose) or native Android players | Not native on Apple devices; Android (ExoPlayer/Media3, Google's standard Android media-playback library), smart TVs, and browsers via dash.js or Shaka | Every modern browser natively; weak on smart TVs and set-top boxes (the dedicated hardware boxes, often supplied by a cable or satellite provider, that decode and display TV service) |
| CDN-friendliness | Excellent: static files, cacheable at every edge | Excellent: same | Poor: each viewer is a stateful session on an SFU; needs a WebRTC-specific delivery network, not a cache |
| DRM integration | FairPlay natively; Widevine and PlayReady when packaged with CENC (Common Encryption, a shared encryption standard) in cbcs mode (one of CENC's two encryption patterns, the one all three DRM systems' current clients can decrypt) | Widevine and PlayReady natively; FairPlay-compatible via cbcs | No standard studio DRM path. Media is encrypted in transit (SRTP) but the viewer's endpoint decrypts it, so it is not licence-controlled |
Where the classic HLS figure comes from: the HLS spec makes a player stay at least three target durations behind the live edge (the most recently published moment of the stream). With 6 s segments that is 18 s before encode, packaging and network are even counted.
MSE (Media Source Extensions) is the browser API that lets JavaScript feed downloaded segments into the video element; it is what lets hls.js and dash.js play these formats in browsers that do not support them natively.
Recommendation per use case
VOD on heterogeneous devices: CMAF, exposed as HLS and DASH
- Apple devices want HLS (and FairPlay); Android, smart TVs and many set-top boxes are happiest with DASH (and Widevine or PlayReady). Serving only one of them forces JavaScript players or re-packaging on the other half of the device fleet.
- Package the video once as CMAF with
cbcsencryption, and generate an HLS playlist and a DASH MPD that reference the same files. You pay for one set of storage and one set of cache entries, and every DRM system can use the same encrypted bytes. - Latency is irrelevant for VOD, so use 4 to 6 s segments for efficient encoding and fewer requests.
- WebRTC is simply the wrong tool: no seeking model, no caching, no DRM.
Large-scale live events: LL-HLS and LL-DASH over a CDN (multi-CDN at scale)
- Hundreds of thousands to millions of viewers are only affordable if one cached chunk serves thousands of them. HTTP-based delivery keeps that property.
- Low-latency modes (partial segments in LL-HLS, small pieces of a segment published as each finishes encoding instead of waiting for the whole segment; chunked transfer of CMAF in LL-DASH, the segment sent as a stream of pieces usable as they arrive rather than delivered whole) bring latency to a few seconds, which is close enough to broadcast TV that social-media spoilers stop being a problem.
- DRM works exactly as for VOD.
- Choose classic (non-low-latency) HLS and DASH only when a 20 to 30 s delay is acceptable; you get larger buffers and fewer rebuffers (playback pauses where the player stalls to refill its buffer) in return.
Interactive real-time apps (auctions, low-latency chat with video): WebRTC
- In a live auction, a bidder who sees the gavel 6 s late is bidding blind. Only WebRTC reliably keeps glass-to-glass latency (camera sensor to the pixel on the viewer's screen) under a second.
- Put the bids on a separate channel (a WebSocket, a persistent two-way connection separate from the video path, or a WebRTC data channel) to an authoritative server that timestamps and orders them. The video is for humans; the server clock decides who bid first, so a viewer on a worse network is not ruled out by their video lag.
- Protect the content with signed, short-lived join tokens at the SFU and, if needed, forensic watermarking (an invisible, viewer-specific marker embedded in the video, so a leaked recording can be traced back to whose stream it came from), because studio-grade DRM is not available on this path. If licensed premium content must be DRM-protected, that is the signal to drop to LL-HLS or LL-DASH and accept 2 to 5 s.
- For a large passive audience watching an interactive event (few bidders, many watchers), a common hybrid is WebRTC for the participants and LL-HLS for everyone else.
Worked example: why "just use WebRTC for everything" gets expensive
An auction house streams one 1.5 Mbps video to 20,000 simultaneous viewers.
- Total egress either way: 20,000 x 1.5 Mbps = 30 Gbps.
- With LL-HLS on a CDN, the origin serves one copy of each chunk per edge cache; the 30 Gbps is spread over the CDN's existing edges and billed as bandwidth.
- With WebRTC, the SFU fleet itself must forward the 30 Gbps. If one SFU node can sustain, say, 10 Gbps of forwarding (an assumption to be measured on your hardware, not a vendor figure), you need at least 3 nodes, and about 5 with headroom for failure and bursts. Each viewer is also a stateful session with its own congestion-control loop, so a node failure means thousands of reconnects rather than a cache miss.
If only 200 of those viewers are actually bidding, putting just those 200 on WebRTC and the other 19,800 on LL-HLS cuts the SFU load by 99%.
Trade-offs and pitfalls
- Treating "HLS vs DASH" as a single choice in a mixed-device world. With CMAF the real question is which manifests to generate; generating both is cheap.
- Assuming low-latency modes work end to end. Any cache or proxy that buffers a whole object before forwarding it breaks chunked delivery and quietly adds seconds.
- Assuming WebRTC's encryption is DRM. It protects the wire, not the endpoint.
- Ignoring encryption scheme compatibility. Older devices only support
cencin CTR mode (Counter mode, the other common CENC encryption pattern besidescbcs) encryption, which FairPlay does not use; if your device fleet includes them you may still need two encrypted copies. - What would flip the recommendations: a hard requirement below about 1 s for everyone (WebRTC-based distribution despite the cost), or strict DRM requirements on an interactive format (LL-HLS and a slower, server-timed interaction model).
Design a multi-region streaming architecture for six global regions to provide low-latency playback and region-level failover. Explain traffic routing (DNS, Anycast), origin selection, how to replicate or cache metadata and entitlements, and how to keep control-plane operations consistent across regions.
Sample Answer
Direct answer
Split the platform into a data plane (delivering video: CDN, content delivery network, edges, regional origins, playback API) that runs active-active (every region accepts live production traffic at the same time, rather than one active region with idle standbys) in all six regions and keeps working when anything else fails, and a control plane (catalogue publishing, rights rules, DRM (digital rights management) key management, configuration) that has one writable primary region with a standby and pushes versioned, immutable snapshots to the regions. Route viewers with anycast at the edge and API front door and health-checked DNS / CDN origin groups (a named set of candidate origins a CDN's shield can fail over between automatically) for origin selection. Replicate metadata to every region as read replicas plus edge caches, and make entitlement checks local by issuing short-lived signed playback tokens.
Terms
- Anycast: the same IP address is announced from many sites; internet routing (BGP, the protocol networks use to exchange routes) delivers each user to the topologically nearest one. Withdrawing a site's announcement moves its users in seconds.
- DNS-based routing: the DNS answer differs by user location or latency, with health checks removing dead targets. Failover is limited by DNS caching: TTLs (time-to-live) of 30 to 60 s, and some resolvers hold answers longer.
- Entitlement: the record that says a user may watch a title (subscription tier, purchase, region rights).
- Static stability: the data plane keeps serving with its last known good state when the control plane is unreachable.
Architecture
flowchart TB
V[Viewer] --> AC[Anycast front door and CDN edge]
AC --> PA[Playback API in nearest region]
AC --> OG[Origin group: regional origin, then neighbour]
PA --> RR[(Regional read replica: catalogue, entitlements)]
PA --> TOK[Signed playback token and DRM licence]
CP[Control plane primary] --> SNAP[Versioned config snapshots]
SNAP --> PA
CP --> RR
CP -.-> CPS[Control plane standby]
1. Traffic routing
- Edge and API front door on anycast. Users reach the nearest healthy site without waiting on DNS caches, and a region can be drained (stop sending it any new traffic while its existing connections finish naturally, so it empties out instead of being cut off mid-request) by withdrawing its route.
- Behind the CDN, DNS or CDN origin groups for regional origins. The CDN's shield (a regional caching tier that sits between edge points of presence and the origin, absorbing their combined misses so the origin sees far fewer requests) uses a primary regional origin and fails over to a named neighbour on timeout or 5xx. This avoids depending on viewer-side DNS for origin failover.
- Why not DNS alone for viewers: resolver caching means some users keep hitting a dead region for minutes. DNS is fine for origin selection behind the CDN, where you control the resolver and TTL.
2. Origin selection and content placement
- Every region has an origin, but not every title lives in every region. Popular and new titles are replicated to all six; long-tail titles live in two regions (a home and a backup) and are fetched cross-region on a miss. Six full copies of a long-tail catalogue buys almost nothing.
- Live channels are packaged in two regions simultaneously so a regional failure does not interrupt the stream.
3. Metadata and entitlements
- Catalogue metadata: written in the control-plane primary, replicated asynchronously to a read replica in each region, cached at the edge for tens of seconds. Staleness of seconds is acceptable.
- Entitlements: the authoritative record lives in a strongly consistent store (every read is guaranteed to see the latest acknowledged write; no stale answers), built as a single primary with synchronous replication (the primary waits for the standby to confirm it has the write before acknowledging the caller, so a failover never loses that write) to a standby region. Each region holds a read replica.
- Playback check is local: the playback API in the viewer's region reads the local replica and issues a signed playback token (for example a JWT, JSON Web Token, valid for 5 minutes) and a DRM licence. Edges and origins validate the signature without calling any database.
- Read-your-writes after purchase: a purchase writes to the primary; the replica may lag by a second. The purchase response returns a receipt carrying the entitlement version; if the local replica is older than that version, the playback API reads from the primary for that one request.
- Revocation is bounded by the token lifetime: a cancelled subscription stops at the next token refresh (5 minutes), which is the product decision being made explicit.
4. Keeping control-plane operations consistent
- One writer. Publishing, rights windows, DRM key rotation and routing config are written only in the primary region, through a store that gives linearizable writes (every read after a write sees it). Two writable control planes would need conflict resolution for things like "is this title licensed in Brazil", which must never be last-writer-wins (automatically keeping whichever of two conflicting writes carries the later timestamp and silently discarding the other). For example: region A revokes Brazil licensing at 10:00:00.100, and region B, working from a slightly stale read and a clock running a little ahead, re-confirms the title as licensed at 10:00:00.150; last-writer-wins keeps region B's write, so the title stays visible in Brazil with no error raised and no operator aware anything was overwritten.
- Versioned, immutable snapshots. Each change produces version N+1; regions pull snapshots and report the version they have applied. Rollouts go region by region, with automatic halt if error rates rise.
- Monotonic apply. A region never moves backwards to an older snapshot, which prevents a delayed message re-enabling a revoked title.
- Control-plane outage: regions keep serving the last applied snapshot (static stability). Publishing pauses; playback does not. Promote the standby if the primary region is lost.
Worked example: capacity for regional failover
Six regions sharing load evenly carry 1/6 each. Losing one and spreading its load evenly across the other five raises each to 1/5: a factor of (1/5) ÷ (1/6) = 1.2, so each region can run at most 1 ÷ 1.2 = 83% of capacity in normal operation.
Geography rarely allows an even spread. If a failed region's traffic lands on one neighbour, that neighbour carries 2× its load and must run at 50%. If failover is split across two neighbours, each takes 1.5× and must run at 67%. I would configure each region's failover across two named neighbours and plan for 67% steady-state utilisation, which is cheaper than 50% and does not rely on an even spread that routing will not deliver.
Trade-offs and pitfalls
- Multi-writer control plane for "no single region dependency" trades a rare publishing pause for permanent conflict-resolution bugs in rights data. Not worth it for a system where playback does not depend on publishing.
- Entitlement checks calling a central database on every play turn a regional failure into a global outage. Local tokens remove the dependency.
- Failover never tested. Drain one region in production on a schedule; unexercised failover paths are where the outage happens.
- Cold caches on failover: a region absorbing a neighbour's viewers also takes its cache misses. Pre-warm the most popular titles in neighbour origins.
- What would change the design: if regulation requires entitlement data to stay in its country, the single-primary store becomes per-jurisdiction primaries, and the token must also carry the region of record.
Estimate and compare the per-connection server resource implications (file descriptors, memory per connection, CPU context switching or event-loop overhead) and request pattern characteristics for: (a) persistent WebSocket connections, (b) Server-Sent Events (SSE), and (c) short polling every 5 seconds. Provide a back-of-envelope calculation for 1M concurrent clients and discuss the main bottlenecks and mitigation strategies.
Sample Answer
Direct answer
For 1 million clients, WebSockets and SSE (Server-Sent Events, a long-lived one-way HTTP stream from server to browser) have almost the same server shape: 1 million long-lived sockets, so 1 million file descriptors and a memory cost per idle connection, but almost no CPU while nothing is happening. Short polling every 5 seconds flips the cost: few sockets need to be open at any instant, but the fleet must answer 200,000 HTTP requests per second, almost all of them empty, which burns CPU and bandwidth on headers. My recommendation: SSE when traffic is server-to-client only, WebSockets when clients also send frequently, and short polling only as a fallback or when updates are rare and 5 seconds of staleness is acceptable.
Terms: WebSocket is a single TCP connection that starts as a normal HTTP request and is upgraded (the client and server agree, via an Upgrade header, to reuse that same connection as a long-lived, two-way channel instead of closing it after one response) into a channel either side can write to at any time; SSE (Server-Sent Events) is one long HTTP response the server keeps open and writes events into; a file descriptor (FD) is the kernel handle a process holds per open socket; an event loop is a single thread that waits on thousands of sockets at once (epoll on Linux) and runs code only for sockets that have data, instead of dedicating a thread per connection.
Assumptions (stated so the numbers can be checked)
- 1,000,000 concurrent clients, updates are infrequent (the typical case for notifications).
- Idle memory per long-lived connection: 8 KiB across kernel socket, TLS and application state (a budget to validate by measurement, not a benchmark).
- Poll request headers about 800 bytes (cookies, auth token, user agent), empty response about 200 bytes including headers.
- Poll request service time about 20 ms.
- Heartbeat on long-lived connections every 30 seconds.
Back-of-envelope for 1M clients
Short polling at 5 seconds
request rateper dayheader bandwidthin-flight requests (Little’s law)=1,000,000/5=200,000 requests per second=200,000×86,400=17.28 billion requests=200,000×(800+200) B=200 MB/s=1.6 Gbps=200,000×0.020 s=4,000Little's law says the average number of requests in the system equals arrival rate times time each spends inside. So only about 4,000 requests are being processed at any moment, but 200,000 per second must be parsed, authenticated and answered.
WebSocket or SSE
open socketsidle memoryheartbeats=1,000,000 (one FD each)=1,000,000×8 KiB≈7.63 GiB=1,000,000/30≈33,333 small frames per secondAt 16 KiB per connection the memory doubles to about 15.26 GiB; at 4 KiB it is about 3.81 GiB. Either way it spreads across a handful of servers.
Side-by-side comparison
| Dimension | WebSocket | SSE | Short polling (5 s) |
|---|---|---|---|
| FDs held | 1M, permanent | 1M, permanent | About 4,000 in flight; up to 1M if clients keep idle keep-alive connections open |
| Memory per client | Socket, TLS, small app state | Same as WebSocket | Near zero between polls |
| CPU when idle | Near zero with an event loop; heartbeats only | Same; a comment line (: keepalive) as heartbeat | 200,000 requests per second of parsing, auth and empty responses, continuously |
| Context switching | Low with event loop; a thread per connection would be 1M threads, unworkable | Same | Normal request-per-thread or async servers handle it; the load is request volume |
| Request pattern | One upgrade, then small frames both ways | One request, then a stream server-to-client; client sends via separate HTTP calls | A full HTTP request and response every 5 s per client |
| Latency to deliver an update | One network trip | One network trip | About 2.5 s average, up to 5 s |
| Reconnect behaviour | Client code must implement backoff | Browser EventSource reconnects automatically and sends Last-Event-ID (a header naming the last event it received, so the server can resume the stream instead of starting over) | Each poll is independent |
| Proxy and browser concerns | Upgrade must be allowed; idle timeouts | Response buffering by proxies must be off; HTTP/1.1 limits browsers to about 6 connections per origin, so serve over HTTP/2, which multiplexes many logical streams over one TCP connection and removes that per-origin cap | Works everywhere |
Main bottlenecks and mitigations
WebSocket and SSE
- Memory per connection dominates. Allocate read and write buffers only while data is moving, keep app state compact, and cap connections per server from a measured budget.
- FD limits: raise
ulimit -nand the kernel limits; run more than one process per host past about a million. - Reconnect storms after a deploy or failure: 1M reconnects in 60 seconds is about 16,667 TLS handshakes per second. Clients use exponential backoff with jitter (each retry waits longer than the last, up to a cap, with a random amount added so many clients don't all retry at the same instant); servers drain gradually and enforce an accept-rate cap (a limit on new connections accepted per second, so the surge queues instead of overwhelming the process).
- Load balancer idle timeouts kill quiet connections; heartbeat more often than the shortest timeout in the path.
- Fan-out (turning one update into one outbound write per connected client): a broadcast to 1M sockets is 1M writes. Distribute across gateways via pub/sub (publish/subscribe: one publish reaches every subscriber without the publisher knowing who or how many they are) so each server writes only to its own clients.
Short polling
- Request rate and header overhead: 1.6 Gbps of mostly headers to say "nothing new". Mitigate with conditional requests (
ETagandIf-None-Match, which let the server answer "304 Not Modified" with no body), a CDN (content delivery network) or cache in front for shared data, and longer intervals for background tabs. - TLS handshakes: if clients do not reuse connections, 200,000 new TLS handshakes per second would dominate CPU; enforce keep-alive, which pushes FD counts back toward the long-lived models.
- Database load: each poll that checks "anything new since X?" against a database is 200,000 queries per second; serve polls from a cache holding the latest version per user or channel.
- Synchronised polling: clients started together poll together; add random jitter to the interval.
Worked choice
A sports app pushing score updates to 1M fans, with clients never sending: SSE. Same socket cost as WebSockets, automatic reconnection with resume in the browser, plain HTTP that works through ordinary proxies. A collaborative editor where every client sends keystrokes: WebSockets. A dashboard refreshed every few minutes: short polling at a longer interval with ETag, because 1M clients at 5 minutes is only about 3,333 requests per second and no persistent-connection infrastructure is needed.
Trade-offs and pitfalls
- Counting only memory misses that polling's cost is CPU and bandwidth, while persistent connections' cost is memory and FDs. Compare each on its own bottleneck.
- Assuming SSE is cheaper than WebSockets per connection: server-side it is roughly the same. SSE's gains are simplicity and built-in resume, not resources.
- Forgetting mobile: long-lived connections are killed when apps go to the background; mobile clients combine a socket while in the foreground with push notifications otherwise.
Unlock Full Question Bank
Get access to all 20 Real-Time and Streaming System Design interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.