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.
How would you implement authentication and authorization for long-lived WebSocket connections in a multi-tenant SaaS product? Cover initial token exchange during the handshake, token refresh for long-lived sessions, per-message authorization checks where needed, handling token revocation or tenant deprovisioning (including third-party API consumers, not just browser clients), and techniques to minimize re-auth overhead while ensuring security.
Sample Answer
Direct answer
Authenticate once at the handshake with a short-lived credential, then treat the connection as a session whose authority expires with the token: the client sends a fresh token over the socket before expiry, and the server closes the connection if it does not. Every message is authorized against the tenant bound to the connection at handshake, never a tenant named in the payload, and subscribe or publish actions on specific resources are checked against a cached permission set. Revocation is pushed, not polled: a revocation event reaches every gateway, which closes matching connections within seconds (immediate revocation), while routine cases just refuse the next refresh (soft revocation). Third-party API consumers use the same model with OAuth client credentials and are revocable by client ID.
Key terms
- Authentication (authn): who is this? Authorization (authz): may they do this, to this resource, in this tenant?
- Multi-tenant: many customer organisations (tenants) share the same servers; the worst failure is one tenant seeing another's data.
- JWT (JSON Web Token): a signed token carrying claims, meaning the individual pieces of information encoded inside it (user, tenant, expiry, token ID
jti), that a server can verify locally with the issuer's public key, with no network call. - Token introspection: asking the identity provider (the service, yours or an external one like Okta or Auth0, that authenticates principals and issues tokens for them) "is this token still valid?" over the network; accurate but costly per call.
- Revocation: withdrawing access before a token naturally expires.
- Principal: whoever the token represents, a human user or a third-party client, the thing being authenticated and authorized.
- Scopes: the specific permissions baked into a token (for example "read messages in tenant t9"), narrower than everything that principal could ever do.
- OAuth 2.0 client-credentials flow: the OAuth grant (the authorization a client has been issued) used when the caller is a backend service, not a human. The service authenticates directly with its own client ID and secret, no user or browser in the loop, and gets back an access token scoped to what that service is allowed to do.
- Key set / key ID: the issuer publishes a set of public keys it may sign tokens with; each token names which one it used via a key ID (
kid), so a verifier fetches the matching key instead of guessing.
1. Initial token exchange during the handshake
Browser clients cannot set an Authorization header on the browser WebSocket API (the browser's built-in interface for opening a persistent, full-duplex connection upgraded from HTTP), so pick one of these:
| Option | How | Verdict |
|---|---|---|
| Single-use ticket | Client calls an authenticated REST endpoint, gets a random ticket valid for 30 s and one use, connects to wss://.../ws?ticket=... | My default for browsers: short-lived and useless once redeemed, so a URL that lands in a log is harmless |
| Session cookie | Browser sends cookies on the handshake automatically | Fine for first-party apps, but you must check the Origin header, otherwise any website can open a socket with the user's cookies (cross-site WebSocket hijacking) |
Token in Sec-WebSocket-Protocol | Smuggle the token as a subprotocol value (this header exists so client and server can agree on an application-level protocol name, for example chat.v2, not to carry credentials; some clients abuse it because it is one of the few headers the browser API lets you set) | Works but abuses the header; avoid |
| Long-lived token in query string | ?token=eyJ... | Avoid: URLs end up in proxy and access logs |
Third-party API consumers (a partner's backend, not a browser) can set headers, so they send Authorization: Bearer <access token> obtained through the OAuth 2.0 client-credentials flow, scoped to one tenant and specific permissions. On success the gateway builds a connection context: user or client ID, tenant ID, token ID, scopes, token expiry, and a permission-set version. Everything later reads from this context.
2. Token refresh for long-lived sessions
- Access tokens live for 5 to 15 minutes; the connection may live for hours.
- About a minute before expiry, the client obtains a new token (browsers through their normal refresh flow, third parties by repeating client credentials) and sends it as an in-band
reauthmessage. - The gateway verifies it, checks it names the same principal and tenant (a different tenant is rejected, not swapped), and updates the context's expiry and scopes.
- If no valid token arrives by expiry plus a short grace period (say 30 s), the gateway closes with an application close code in the 4000 to 4999 range, which RFC 6455 reserves for application use (for example 4001 for "reauthentication required" and 4003 for "access revoked"). The client reconnects through the normal handshake.
3. Per-message authorization
- Tenant isolation is structural: channel names and resource IDs are resolved inside the connection's tenant; a message that names another tenant's resource simply finds nothing. Never read the tenant from the payload.
- Subscribe and publish are authorization points: "join channel
project:88" is checked against the user's permissions for project 88. Subsequent messages on that subscription are pre-authorized until permissions change. - Sensitive actions (admin operations, payments, deletes) are re-checked per message, or better, performed through a normal HTTP API where standard authz and audit already exist.
- Outbound filtering: the fan-out path (where one published event is delivered out to every subscribed connection) checks that each recipient connection's tenant matches the event's tenant before sending, as a second line of defence.
4. Revocation and tenant deprovisioning
Every gateway indexes its live connections by user ID, tenant ID, client ID and token ID. Revocation events travel on a control channel (a pub/sub topic every gateway subscribes to).
| Case | Mode | What happens |
|---|---|---|
| Password reset, user logs out everywhere, suspected compromise | Immediate | Event {user: u1} published; each gateway closes u1's connections now; token IDs added to a deny list until their natural expiry so a reconnect with an old token fails |
| Role changed (lost access to one project) | Immediate, narrow | Permission-version bump; gateways drop the affected subscriptions and reload permissions without closing the socket |
| Plan downgrade, routine offboarding | Soft | Mark the principal; the next refresh is refused; access ends within one token lifetime (at most 15 minutes) |
| Tenant deprovisioned | Immediate | Stop issuing tickets and tokens for the tenant; publish {tenant: t9}; every gateway closes all t9 connections, including third-party clients; queued outbound messages for t9 are discarded |
| Third-party app's access revoked by the tenant admin | Immediate | Revoke the OAuth grant (the permission that client-credentials flow issued to that client, independent of any single still-valid token); publish {client: acme-sync, tenant: t9}; close that client's connections in that tenant only, leaving the same app's connections in other tenants alone |
Revocation events must be durable and replayable (a gateway that restarts or was partitioned must read recent revocations before accepting connections), and the deny list must outlive the longest token lifetime.
sequenceDiagram
participant C as Client
participant A as Auth API
participant G as Gateway
participant R as Revocation topic
C->>A: Request ticket with session
A-->>C: Ticket valid 30 s, single use
C->>G: Handshake with ticket
G-->>C: Welcome, context tenant t9
C->>G: reauth with new token before expiry
R-->>G: Revoke tenant t9
G-->>C: Close 4003 access revoked
5. Minimising re-auth overhead
- Verify JWTs locally: signature check against cached public keys (refetch the issuer's key set on unknown key ID), no network call per refresh.
- Authorize at subscribe time and cache the permission set in the connection context, invalidated by version bumps rather than time-based polling.
- Push revocations instead of re-validating every connection periodically.
- Refresh in band on the existing socket rather than reconnecting.
Worked example: why per-message introspection does not scale
A product has 2,000,000 connections and 15-minute tokens.
- In-band refresh load: 2,000,000 / 900 s = about 2,222 token verifications per second across the fleet, each a local signature check.
- If instead every message were checked by introspection at a traffic level of 50,000 messages per second, the identity provider would take 50,000 calls per second, over 20 times the refresh load, and every message would wait on a network round trip.
- Revocation for a tenant with 3,000 connections spread across 40 gateway nodes: one event broadcast to 40 nodes, each closing its local share from an in-memory index. The cost is proportional to the tenant's connections, not to the fleet.
Pitfalls
- Authenticating the handshake and never again, so a fired employee's socket lives for days.
- Trusting a tenant ID inside the message payload.
- Cookie auth without an
Origincheck. - A revocation broadcast that is fire-and-forget, so a gateway that missed it keeps a revoked session alive.
- Forgetting non-browser API consumers, which often hold the longest-lived connections of all.
Compare WebRTC data channels and WebSocket for real-time peer-to-peer data synchronization use cases (e.g., multi-user collaborative whiteboard). Include signaling architecture, NAT traversal (STUN/TURN), expected throughput/latency differences, privacy/security trade-offs, and operational costs/complexity (TURN servers, scaling signaling). When would you choose one over the other?
Sample Answer
Direct answer
For a multi-user collaborative whiteboard I would use WebSocket (a persistent, two-way connection between browser and server that either side can push messages over at any time, unlike a normal request that closes after one reply) through a server as the backbone, because the server must persist the board, serve late joiners, enforce permissions and scale beyond a handful of users. I would add WebRTC data channels only for ephemeral, high-frequency data (live cursor positions) in small rooms, where peer-to-peer delivery and unordered, unreliable mode genuinely help. WebRTC is the right primary choice when rooms are small, latency matters more than persistence, and the data should not pass through your servers (for example, peer-to-peer file transfer or games).
The two technologies
- WebSocket: one long-lived TCP connection between the browser and your server, carrying messages both ways. Everything goes through the server.
- WebRTC data channel: a browser-to-browser channel. It runs on SCTP (a message-oriented transport) inside DTLS (TLS for UDP, so it is always encrypted) over UDP (User Datagram Protocol: packets are sent without any guarantee they arrive, arrive once, or arrive in order, unlike TCP, Transmission Control Protocol, the protocol under WebSocket, which guarantees all of that by resending anything lost, at the cost of waiting for it). Each channel can be reliable and ordered (like TCP) or configured unordered and with limited retransmits (
ordered: false,maxRetransmitsormaxPacketLifeTime), so a lost cursor packet is skipped instead of delaying newer ones.
Signaling architecture
WebRTC has no built-in way for two browsers to find each other. Signaling is the exchange of connection descriptions before the peer connection exists: an SDP offer and answer (Session Description Protocol: what each side supports) and ICE candidates (addresses at which each peer might be reachable). You build signaling yourself, and in practice it is usually a WebSocket server. So choosing WebRTC does not remove the WebSocket server; it adds a second system.
sequenceDiagram
participant A as Browser A
participant S as Signaling server (WebSocket)
participant B as Browser B
participant T as STUN / TURN
A->>T: STUN: what is my public address?
A->>S: SDP offer and ICE candidates
S->>B: forward offer and candidates
B->>S: SDP answer and ICE candidates
S->>A: forward answer and candidates
A-->>B: direct data channel if NAT allows
A-->>T: otherwise relay through TURN
NAT traversal: STUN and TURN
Most devices sit behind NAT (network address translation: a router shares one public IP among many private devices), so a peer's private address is unreachable from outside.
- STUN (Session Traversal Utilities for NAT): a lightweight server that tells a browser its public IP and port. Cheap: a few small packets per session.
- TURN (Traversal Using Relays around NAT): when a direct path is impossible (symmetric NATs, a NAT type that hands out a different public port for every destination a device talks to, so the address a peer discovers via STUN for one connection cannot be reused for another; or strict corporate firewalls), a TURN server relays all the traffic. At that point the connection is no longer peer-to-peer, and you pay for the egress (the cost of data leaving your servers to the internet, billed per GB by most cloud and hosting providers).
- ICE (Interactive Connectivity Establishment) is the procedure that tries direct candidates first and falls back to TURN.
Throughput and latency
- Latency: peer-to-peer saves the server hop only when the peers are close to each other and far from the server. Two colleagues in one office with a server on another continent gain a lot; two users in different countries with a nearby server region gain little. Relayed (TURN) paths add a hop comparable to the WebSocket path.
- Head-of-line blocking: TCP delivers bytes in order, so one lost packet stalls every message behind it until it is retransmitted. An unordered, unreliable data channel avoids that, which is the real advantage for streams where only the newest value matters (cursors).
- Throughput: in a mesh (every peer connected to every other), each peer uploads its data once per other peer, so upload grows with room size; with WebSocket each client uploads once and the server fans out.
- Connection setup: WebSocket is ready after one TCP and TLS handshake; WebRTC needs signaling plus ICE checks, typically noticeably longer, which matters for short sessions.
Privacy and security trade-offs
| Concern | WebSocket via server | WebRTC data channel |
|---|---|---|
| Encryption in transit | TLS (wss://) to the server; server sees plaintext | DTLS mandatory; peers see plaintext, a TURN relay only sees encrypted packets |
| IP address exposure | Peers see only the server | ICE candidates reveal each peer's IP addresses to the other peers; mitigate with iceTransportPolicy: "relay" (TURN only) at the cost of relay bandwidth |
| Authorization and moderation | Server checks every message against permissions | Server is out of the data path: a malicious peer can send anything; each client must validate |
| Persistence, audit, late join | Natural: server stores the operations | Needs a separate path to a server anyway |
Operational cost and complexity
- TURN servers are bandwidth-heavy relays you must run or buy, in several regions, with credentials that expire (short-lived TURN credentials, so a leaked one cannot be reused as a free relay). Their cost is dominated by egress.
- Scaling signaling is a WebSocket fan-out problem you have anyway.
- Debugging WebRTC connectivity failures (which NAT, which candidate failed) is significantly harder than debugging a WebSocket.
- WebSocket at scale costs connection-holding servers and egress for fan-out, both familiar and well-tooled.
Worked example: a 10-person whiteboard
Cursor updates at 30 per second, 40 bytes each.
- WebSocket: each client uploads 30 x 40 = 1,200 bytes/s; the server sends each update to the other 9 users: 10 x 30 x 9 x 40 = 108,000 bytes/s of server egress per room (about 0.86 Mbit/s).
- WebRTC mesh: 10 x 9 / 2 = 45 peer connections per room; each peer uploads 30 x 9 x 40 = 10,800 bytes/s, nine times the WebSocket upload. Fine for 10 people on broadband; at 30 people it becomes 435 connections and 29x upload per peer, which is why meshes stop scaling around small rooms and larger rooms use a server relay (an SFU, selective forwarding unit: a server that receives each peer's stream once and forwards, i.e. selectively relays, a copy to every other peer without decoding or mixing it, so upload no longer grows with room size) anyway.
- If some sessions need TURN, each relayed peer's traffic passes through your TURN server, so its cost scales with the relayed fraction times the mesh traffic. As an illustrative order of magnitude (not a benchmark): production WebRTC deployments commonly report something like one session in five to one in four falling back to TURN, but the real fraction depends heavily on your users' networks, corporate firewalls and mobile carriers. Measure that fraction from ICE statistics in your own user base rather than assuming it.
The drawing strokes themselves (which must persist and be seen by late joiners) go over WebSocket in both designs.
When to choose which
- WebSocket: persistent shared state, more than a few participants, need for authorization, moderation, history or server-side integrations; users behind strict networks.
- WebRTC data channels: small rooms, latency-critical or bulky data between peers (games, file transfer, cursor streams), data you do not want on your servers, or an app already using WebRTC for audio and video.
Trade-offs and pitfalls
- "Peer-to-peer is free" is false once TURN relays and signaling are counted.
- Using WebRTC for the authoritative board state forces you to reinvent ordering and persistence, and a single malicious peer can corrupt others' state.
- Forgetting that ICE candidates leak IP addresses is a real privacy issue in apps with strangers in the same room.
- WebTransport (a newer browser API over HTTP/3 with unreliable datagrams to a server) addresses head-of-line blocking without peer-to-peer; worth evaluating where browser support fits your audience.
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.
Design an end-to-end realtime system to broadcast order-book updates for a high-frequency trading platform to thousands of subscribers. Requirements: maintain strict ordering per instrument, sub-10ms tail latency within region, support replay of missed updates for up to 1 minute, guarantee at-most-once delivery for downstream processing, and handle 100k updates/sec. Detail components (ingest, sequencer, transport), storage for replay, sequencing, backpressure handling, and operational considerations.
Sample Answer
Direct answer
Put a single-writer sequencer per instrument at the centre: every update for an instrument passes through one process that stamps it with a gap-free, per-instrument sequence number, which is what "strict ordering" means in practice. From the sequencer, updates go out on UDP multicast (one packet copied by the network switches to every subscriber) inside the data centre, with a second, identical A/B feed on separate network paths so one lost packet is usually recovered from the other feed with no round trip. Missed updates are recovered from retransmission servers holding a 60-second in-memory ring buffer, and from snapshots when a subscriber is too far behind. Slow subscribers never slow the sequencer: they drop and recover. The "at-most-once" requirement (downstream applies each update no more than once, even though the transport itself may redeliver it) is met on the subscriber side, by applying each sequence number at most once.
Terms
- Instrument: a tradable security, such as a stock, option, or futures contract. Every message in this system is scoped to one instrument, which is what makes the system partitionable.
- Order book: the live list of buy and sell orders waiting to trade for one instrument, ordered by price. It is what a matching engine maintains and what this system broadcasts changes to.
- Matching engine: the system that receives buy and sell orders for an instrument and pairs (matches) them into trades, updating the order book.
- Book change: one update to the order book: an order was added, modified (price or quantity changed), deleted (cancelled), or resulted in a trade (matched against another order).
- Ring buffer: a fixed-size buffer that overwrites its oldest entries once full, used here to hold the last 60 seconds of messages for replay without unbounded memory growth.
- Multicast group: a network address that, once a device subscribes to it, causes network switches to deliver a copy of every packet sent to that address, without the sender addressing each recipient individually.
- Hot standby: a second instance that runs the same computation as the primary in real time but does not publish, so it can take over instantly if the primary fails.
- A/B feed / arbitration: sending the same data twice over two independent network paths (A and B); arbitration is picking whichever copy of a given message arrives first and discarding the duplicate.
Requirements and the one that needs interpreting
- Strict ordering per instrument (not globally: ordering across instruments is not required, which is what makes the system partitionable).
- Sub-10 ms tail latency (p99, the 99th percentile) within a region.
- Replay of missed updates for up to 1 minute.
- At-most-once delivery for downstream processing, and 100,000 updates per second.
"At-most-once" normally means "may be lost, never duplicated", while "replay missed updates" means "not lost". Taken together, the sensible reading is: the transport may deliver a message more than once (A/B feeds and replays guarantee duplicates), but downstream processing must apply each update at most once, and gaps must be repairable within 60 seconds. Each subscriber keeps last_applied_seq per instrument and discards anything at or below it. I would confirm this reading with the interviewer; if they truly meant "fire and forget, never replay", the retransmission tier disappears and the design gets simpler, but the order book would silently diverge after any loss.
Components
flowchart LR
ME[Matching engines] --> SEQ[Sequencers by instrument partition]
SEQ --> MCA[Multicast feed A]
SEQ --> MCB[Multicast feed B]
SEQ --> RING[(Retransmission ring buffer 60 s)]
SEQ --> SNAP[(Snapshot service)]
MCA --> SUB[Subscribers]
MCB --> SUB
SUB -. gap request .-> RING
SUB -. resync .-> SNAP
MCA --> EDGE[TCP fan-out tier for remote subscribers]
- Ingest: matching engines (which already process orders per instrument) emit book changes (add, modify, delete, trade). Each engine owns a set of instruments.
- Sequencer: instruments are partitioned (for example by instrument id) across sequencer processes. Within a partition, one thread stamps
(instrument_id, seq)and writes the message to the feeds and to the ring buffer. Single-writer means no locks and no coordination on the hot path. - Transport: UDP multicast, one multicast group per partition, so a subscriber that only cares about some instruments joins only those groups.
- Retransmission servers: listen to the feed, store the last 60 seconds, answer "resend instrument X, sequences 1001 to 1010".
- Snapshot service: maintains the full book per instrument with the sequence number it reflects, so a subscriber that joins late or falls far behind can load the book and continue from the stream.
- Fan-out tier (TCP or WebSocket) for subscribers outside the data centre that cannot receive multicast.
Why multicast: the bandwidth arithmetic
With 64-byte messages (a compact binary format: message type, instrument id, sequence, price, quantity, side, timestamp):
- One subscriber taking the full feed: 100,000 x 64 bytes = 6.4 MB/s = 51.2 Mbit/s.
- 2,000 subscribers over unicast (a separate copy per subscriber): 51.2 Mbit/s x 2,000 = 102.4 Gbit/s leaving the publisher. That is a hardware problem and a latency problem, because the last copy is sent long after the first.
- Multicast: the publisher sends 51.2 Mbit/s once; switches replicate it. Latency is the same for every subscriber.
Averages understate market data: bursts at the open or on news can be several times the average. I would provision the feed and every subscriber path for at least 5x (500,000 updates per second, 256 Mbit/s), and treat that multiplier as an assumption to confirm with historical peak data.
Sequencing and gap detection
- Each message carries
(instrument_id, seq); sequences are per instrument, gap-free, starting at 1 per trading session. - Packets also carry a partition-level sequence (a per-instrument sequence only reveals a gap when the next message for that same instrument arrives, and if that instrument stays quiet there may be no next message for a long time) so a subscriber detects a lost packet even for an instrument that then goes quiet. A periodic heartbeat per partition carries the latest sequence for the same reason.
- Subscriber logic: receive from A and B, take whichever copy arrives first, discard the duplicate by sequence. If a gap persists beyond a short window (say 1 ms, to absorb A/B skew), request the range from a retransmission server. If the gap is older than 60 s, or the request fails, load a snapshot and replay from its sequence.
Storage for replay
60 seconds x 100,000 updates per second = 6,000,000 messages x 64 bytes = 384 MB. At the 5x burst assumption it is 1.92 GB. This is an in-memory ring buffer, indexed by (partition, seq) to an offset, on at least two retransmission servers. No disk is on the replay path. Separately, the full session is written to durable storage asynchronously for audit and regulatory record-keeping, which is off the latency path.
Latency budget for the p99 under 10 ms
| Hop | Budget |
|---|---|
| Matching engine to sequencer (same data centre) | 1 ms |
| Sequencing and serialization | 0.5 ms |
| Multicast to subscribers in the same region | 2 ms |
| Subscriber decode and apply | 1 ms |
| Headroom for bursts and retransmission of a single lost packet | 5.5 ms |
These are allocations, not measurements: each hop gets a histogram and an alert on its p99, so a regression is attributable to a hop. Keeping headroom matters because a gap repair costs one round trip to a retransmission server, and it must fit within the 10 ms target for the p99 to hold. Timestamps at each hop need clock synchronization much finer than the budget; PTP (Precision Time Protocol, hardware-assisted clock sync) is the standard choice.
Backpressure
The sequencer and the multicast feed never wait for any subscriber. Backpressure from one slow reader must not become latency for thousands of others.
- A slow multicast subscriber overflows its own socket buffer (the OS's fixed-size holding area for incoming network data before the application reads it) and loses packets (the OS drops new packets once that buffer is full, since there is nowhere to put them); its gap logic recovers via retransmission, and if it is persistently slow, it is disconnected from retransmission (rate-limited per subscriber) and told to use snapshots.
- In the TCP fan-out tier, each subscriber has a bounded outbound queue. On overflow, the tier either conflates (sends the latest book state for each instrument instead of every intermediate change, acceptable for display clients) or disconnects the subscriber with a "resync from snapshot" instruction (required for trading clients that need every change).
- Retransmission servers are rate-limited per subscriber, so one broken client cannot use the replay tier to hurt everyone else's recovery.
Operational considerations
- Sequencer failover: a hot standby consumes the same inputs and computes the same sequence numbers deterministically, but does not publish. On failure, it takes over publishing from the next sequence. A paused primary is dangerous precisely because it can resume: a garbage-collection pause, a frozen VM, or a network partition can make the standby correctly promote itself while the old primary is still running, and once unpaused the old primary still believes it is the leader and keeps publishing. Fencing (a monotonically increasing leadership epoch in every packet, so subscribers ignore the old primary if it wakes up) prevents two sequencers publishing the same numbers.
- Start of day: sequences reset per session; subscribers load the opening snapshot.
- Monitoring: per-hop p99, gap rate per subscriber, retransmission requests per second, A/B arbitration stats (how often B saved you), snapshot requests.
- Capacity tests with recorded peak days replayed at 5x.
Trade-offs and pitfalls
- Kafka or a similar log broker gives replay and ordering per partition for free, but adds a disk-backed, replicated write before delivery. That is fine for downstream analytics consumers (they can subscribe to a copy) and wrong for the sub-10 ms path.
- Global ordering across instruments would force one sequencer for everything and cap throughput; per-instrument ordering is what the requirement asks for.
- Treating duplicates as the enemy at the transport: A/B feeds deliberately create duplicates; dedup belongs at the subscriber, by sequence.
- What would change the design: if subscribers were mostly across the internet (retail apps), multicast is unavailable and the design becomes a TCP/WebSocket fan-out tree with conflation, and "sub-10 ms" would have to be redefined as within our network.
Design strategies to manage memory usage for per-connection state in servers handling millions of persistent connections, so growth in connection count does not cause unpredictable memory spikes, OOM kills, or long GC pauses. Explain how you would represent, store, and evict that state, and how you would validate your approach holds at scale.
Sample Answer
Direct answer
Treat memory per connection as a budget you design to, not a number you discover in an outage. I would (1) measure what one connection actually costs across kernel, TLS, runtime and application; (2) shrink the steady-state cost by keeping idle connections tiny, allocating read and write buffers only while data is moving and returning them to a shared pool; (3) keep only routing-critical state in the process and push everything else to an external store; (4) put hard bounds on everything that can grow (outbound queues, subscriptions, total connections) with eviction rules for each; and (5) enforce a connection cap derived from the budget so the server refuses new connections before the kernel's out-of-memory (OOM) killer chooses for you. Then prove it with a load test that measures bytes per connection as a slope, not a snapshot.
Terms: OOM kill is the operating system terminating a process that exceeded its memory limit; GC pause is the time a garbage-collected runtime (Java, Go, Node.js) stops or slows the program to reclaim memory, which grows with the number of live objects and pointers it must scan.
Where per-connection memory goes
| Layer | What it holds | How it behaves |
|---|---|---|
| Kernel socket | Socket structure plus receive and send buffers | Buffers are sized dynamically; memory is consumed when data is queued, bounded by net.ipv4.tcp_rmem and tcp_wmem |
| TLS | Session keys and record buffers | Library-dependent; can be a large share if read and write buffers are allocated per connection |
| Runtime | Thread stack or goroutine stack, event-loop registration | A thread per connection is the worst case; one goroutine (Go's lightweight, runtime-scheduled thread; many goroutines share a handful of OS threads) per connection starts small but grows |
| Application | User ID, subscriptions, sequence numbers, outbound queue | Fully under your control, and the part that spikes |
On the Linux kernel checked for this answer, sysctl net.ipv4.tcp_rmem reports 4096 131072 33554432 (minimum, default, maximum bytes). That default of 128 KiB is a ceiling the kernel grows into when data is waiting, not memory spent on an idle socket. It explains a common spike, though: if a million clients all stop reading at once, the kernel can queue far more than your application budget assumes.
Strategies
Represent state compactly
- Store IDs, not objects: a subscription is a 4-byte channel index, not a copy of a channel object with its name string.
- Use bitsets for subscriptions to a small fixed set of channel types.
- Intern repeated strings (region, device type, tenant) so a million connections share one copy: interning means storing one copy of a repeated string and pointing every reference at it, instead of each connection holding its own copy of the same bytes.
- In garbage-collected languages, prefer flat structs or parallel arrays over one heap object per field: fewer objects means less GC scanning.
The demo below measures three ways to hold the same four fields per connection: a plain dict (a hash table of key/value pairs, one per connection), a __slots__ object (a Python object with a fixed set of fields and no per-instance dict, so it skips that hash-table overhead), and columnar arrays (one flat array('q'), a compact array of 64-bit integers, per field, shared across all connections instead of one container per connection). tracemalloc is Python's built-in tool for measuring how many bytes a block of code actually allocated. A runnable illustration of how much representation alone changes the per-connection cost for the same four integer fields:
import sys, tracemalloc
from array import array
N = 100_000
class ConnSlots:
__slots__ = ("user_id", "last_seen", "sub_mask", "seq")
def __init__(self, i):
self.user_id, self.last_seen, self.sub_mask, self.seq = i, 1_700_000_000 + i, 0b1011, 0
def as_dict(i):
return {"user_id": i, "last_seen": 1_700_000_000 + i,
"sub_mask": 0b1011, "seq": 0}
def measure(build):
tracemalloc.start()
obj = build()
size, _ = tracemalloc.get_traced_memory()
tracemalloc.stop()
return size / N
layouts = {
"dict per connection": lambda: [as_dict(i) for i in range(N)],
"__slots__ object per conn": lambda: [ConnSlots(i) for i in range(N)],
"columnar arrays (4 x int64)": lambda: [array("q", range(N)),
array("q", (1_700_000_000 + i for i in range(N))),
array("q", [0b1011]) * N,
array("q", [0]) * N],
}
print(f"Python {sys.version_info.major}.{sys.version_info.minor}, {N:,} connections")
for name, build in layouts.items():
print(f" {name:<28} {measure(build):6.1f} bytes/connection")
Output (exact bytes vary by interpreter version):
Python 3.14, 100,000 connections
dict per connection 259.9 bytes/connection
__slots__ object per conn 139.9 bytes/connection
columnar arrays (4 x int64) 32.3 bytes/connection
Same information, an 8x difference, and the columnar layout also gives the garbage collector 4 objects to track instead of hundreds of thousands. The same principle applies in Java (primitive arrays versus boxed objects) and Go (structs without pointers are not scanned by the GC).
Allocate buffers on demand, from a pool
The biggest per-connection cost is usually read and write buffers held permanently for sockets that are idle 99% of the time. With an event loop (a single thread that asks the kernel which of thousands of sockets have data waiting, instead of dedicating a thread to each one; epoll is the Linux system call that answers that question), the server learns a socket is readable before it reads, so it can borrow a buffer from a shared pool, read, process, and return it. A million idle connections then hold no application buffer at all; only the few thousand active at any moment do. Frameworks built for very high connection counts do exactly this instead of one goroutine or thread with its own buffer per connection.
Externalize what is not needed per message
Keep in process only what the send path needs: the socket, subscription indices, the last sequence number. Profile data, message history and presence details live in Redis or a database and are fetched when needed. This also makes the gateway disposable, which helps deploys and failover.
Bound and evict everything that grows
| Structure | Bound | Eviction rule |
|---|---|---|
| Outbound queue per connection | For example 256 KiB or 1,000 messages | Slow consumer: drop oldest non-critical messages, then disconnect with a "reconnect and resync" code |
| Subscriptions per connection | For example 500 | Reject further subscribes |
| Idle connections | Heartbeat every 30 s | Close after 2 missed heartbeats |
| Caches (user permissions, channel metadata) | Fixed entry count | LRU (least-recently-used) eviction |
| Total connections per process | Derived from budget, below | Refuse new handshakes; balancer routes elsewhere |
The slow-consumer rule matters most: one client on a bad mobile network, subscribed to a busy channel, is how a single connection grows from kilobytes to hundreds of megabytes.
Tame the runtime
- Set the runtime's memory ceiling below the container limit so the GC works harder before the kernel kills the process: Go's
GOMEMLIMIT(a soft cap that tells the Go runtime to collect more aggressively as it is approached), Java's-Xmxwith a low-pause collector such as ZGC (a Java garbage collector designed to keep pause times low even on multi-gigabyte heaps), Node.js--max-old-space-size. - Keep headroom for a reconnect storm: TLS handshakes allocate temporary memory, and 100,000 handshakes in a minute is a burst the steady-state figure does not show.
- Avoid a thread per connection: a million threads means a million stacks and heavy context switching.
Worked example: deriving the connection cap
Budget for a 16 GiB container:
reserved for runtime, handshake bursts, cachesavailable for connectionsmeasured cost per connection (target)cap=4 GiB=16−4=12 GiB=12,582,912 KiB=8 KiB=12,582,912/8=1,572,864 connectionsI would set the enforced cap to about 70% of that, roughly 1.1 million, leaving room for the fraction of connections that are active with full buffers at any moment. That 30% reserve is $12{,}582{,}912 \times 0.3 \approx 3{,}774{,}874$ KiB, about 3.6 GiB. If a connection actively moving data holds roughly 32 KiB of read and write buffer beyond the 8 KiB baseline (two 16 KiB buffers), that reserve covers about $3{,}774{,}874 / 32 \approx 117{,}965$ connections being simultaneously active with full buffers, about 10.7% of the 1.1 million cap. Given that idle connections were established earlier as the 99% case, a working assumption of roughly one in ten connections active at once is the kind of number the slope test below should confirm or correct, not a fixed target. The cap is enforced in code; the balancer's health check reports "full" so new handshakes go elsewhere.
Validating it at scale
- Slope test. Ramp synthetic clients from 0 to 1 million in steps (connection count per client machine is limited by source ports, so use several client hosts or several source IPs). At each step record resident memory (RSS). The slope of RSS versus connections is the true bytes per connection; compare it to the budget.
- Activity mix. Repeat with 1%, 10% and 50% of connections actively sending. Idle cost and active cost are different numbers and both belong in the budget.
- Slow-consumer test. Make 5% of clients stop reading. Memory must plateau at the queue bound, and those clients must be disconnected.
- Soak test. Run 24 hours or more at target load with churn (connects and disconnects). Flat memory proves no leak; a rising line shows a structure that is never evicted.
- GC observation. Record pause-time percentiles (p99, the 99th percentile) during the tests, not just averages.
- Production guardrails. Alert on bytes per connection (RSS divided by live connections), outbound queue depth percentiles and slow-consumer disconnect counts.
Trade-offs and pitfalls
- Pooling buffers adds complexity (use-after-return bugs in manual-memory languages, where code keeps writing to a buffer after returning it to the pool and it has already been handed to a different connection, corrupting both); it pays off at hundreds of thousands of connections, not at a few thousand.
- Externalizing state adds a network hop on paths that need it. Keep the send path fully in memory; externalize only what the send path does not touch.
- Disconnecting slow consumers is a product decision. For a chat app the client resyncs; for a trading feed you might prefer conflation (sending only the latest value per symbol) over disconnection.
- A snapshot measurement lies. Memory at 10,000 connections divided by 10,000 includes fixed process overhead and overstates per-connection cost; the slope is the honest number.
Unlock Full Question Bank
Get access to all Real-Time and Streaming System Design interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.