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.
Estimate capacity and do rough calculations: how many WebSocket connections can a single 32GB, 8-core server handle given per-connection protocol state ~1KB and average application buffer 50KB per active connection, leaving 6GB headroom for OS/runtime and other services? Show your calculations and describe other factors that could change your estimate (file descriptor limits, TCP buffer sizes, TLS memory overhead).
Sample Answer
Direct answer
Memory is the binding constraint. After the 6 GB headroom, 26 GB is left for connections; at 1 KB of protocol state plus 50 KB of application buffer, each active connection costs 51 KB, so the server holds about 535,000 fully active connections (working in binary units: 26 GiB divided by 51 KiB). That figure moves a lot with three things the question flags: the fraction of connections that are actually active (idle ones do not need the 50 KB buffer), kernel TCP buffers that fill when clients stop reading, and TLS state. The 8 cores are not the limit for mostly idle connections, but they become the limit once message rates or TLS handshake storms rise.
The core calculation
I use binary units throughout (1 GiB = 1,024 MiB = 1,048,576 KiB) so the basis stays consistent; decimal units are shown at the end for comparison.
memory for connections26 GiBper active connectionconnections=32−6=26 GiB=26×1,048,576=27,262,976 KiB=1+50=51 KiB=27,262,976/51≈534,568Reading the same inputs as decimal (26 x 10^9 bytes divided by 51 x 10^3 bytes) gives about 509,803. The roughly 5% difference is only the unit convention, which is why an answer should state which one it used.
So the headline estimate is about 500,000 active connections, and I would plan to operate at 60 to 70% of it (300,000 to 350,000) to leave room for bursts.
Sensitivity 1: how many connections are really active
The 50 KiB buffer is described per active connection. Most real-time connections (chat, notifications, live scores) are idle most of the time, and a server that allocates buffers only while data is moving pays the 50 KiB only for active ones. If a fraction a is active at any moment:
| Active fraction | Average KiB per connection | Connections in 26 GiB |
|---|---|---|
| 100% | 51 | 534,568 |
| 50% | 26 | 1,048,576 |
| 20% | 11 | 2,478,452 |
| 10% | 6 | 4,543,829 |
This is the single biggest lever, and it is a software design choice (pooled buffers on demand versus a permanent buffer per connection). Note the 50% row lands at exactly 1,048,576, which collides with the next factor.
Sensitivity 2: file descriptor limits
Every socket is a file descriptor (FD), the integer handle the kernel gives a process for an open file or socket. Limits stack at three levels:
- Per-process soft and hard limit (
ulimit -n): often 1,024 by default for a login shell, far too low. A service unit or container sets its own; the container used for this answer reported 20,480. - Per-process kernel maximum
fs.nr_open, commonly 1,048,576 by default, which caps how highulimit -ncan be raised. - System-wide
fs.file-max.
So beyond about a million connections per process, FD limits need raising (or you run several processes per host). The server also needs FDs for its own outbound connections to the broker, logs and databases.
Sensitivity 3: TCP buffer sizes
The 1 KiB "protocol state" underestimates the kernel. On the Linux kernel checked for this answer, sysctl net.ipv4.tcp_rmem reported 4096 131072 33554432 and tcp_wmem reported 4096 16384 4194304 (minimum, default and maximum bytes per socket). Those buffers are consumed only when data sits in them, so idle sockets cost little. The danger is a slow or stalled client: the kernel send buffer for that socket fills toward its limit while your application also queues. Worked example: if 10% of 500,000 connections each had 64 KiB stuck in kernel buffers, that is 50,000 x 64 KiB, about 3.05 GiB, invisible to an application-level budget. Mitigations: lower the per-socket defaults for this workload, cap total TCP memory system-wide with net.ipv4.tcp_mem (a page-count budget for the whole host's TCP buffers, unlike tcp_rmem/tcp_wmem above, which set the min/default/max bytes for one socket), and disconnect slow consumers.
Sensitivity 4: TLS memory overhead
TLS adds per-connection state (keys, session data) and, depending on the library, its own record buffers for reading and writing. I would not assume a number: measure it with the slope method: record resident memory (RSS, resident set size, the RAM the process actually occupies) at 100,000 and at 200,000 idle TLS connections, and divide the difference by 100,000. The effect of an extra x KiB per connection on the 100%-active case:
| Extra per connection | Connections in 26 GiB |
|---|---|
| +0 KiB | 534,568 |
| +4 KiB | 495,690 |
| +16 KiB | 406,910 |
| +32 KiB | 328,469 |
| +64 KiB | 237,069 |
Terminating TLS on a separate proxy tier moves this cost off the connection servers; it does not remove it from the system.
Other factors that change the estimate
- CPU (8 cores). Idle connections cost almost no CPU with an event loop. Load comes from messages and handshakes. Example: 500,000 connections each sending one message every 10 seconds is 50,000 messages per second, about 6,250 per core, reasonable. A reconnect storm where all 500,000 reconnect within 60 seconds is about 8,333 TLS handshakes per second, and handshakes cost far more CPU than steady messages. That storm, not steady state, often sets the real per-server cap.
- Heartbeats. Pinging 500,000 connections every 30 seconds is about 16,667 small frames per second, a background load worth counting.
- Runtime overhead. A thread per connection would need a stack per thread, which makes 500,000 impossible; the estimate assumes an event loop or lightweight tasks (an event loop is a single thread that services many sockets by reacting to readiness, for example epoll on Linux; lightweight tasks are a runtime's cheap user-space threads, such as Go's goroutines, that let thousands run without one OS thread each). A garbage-collected runtime also needs headroom above live data.
- Connection tracking. If a firewall or NAT (network address translation) on the host tracks connections (conntrack, the kernel's table of in-progress connections, used to route return traffic and enforce rules), its table size is another limit.
- Ports are not the limit on the server side. Clients connecting to one listening port are distinguished by their own IP and port, so the server does not run out of ports. Ephemeral ports (the local port numbers the OS hands out for an outbound connection; Linux defaults to about 28,000 of them,
net.ipv4.ip_local_port_rangetypically 32768-60999) do limit load-test clients: one machine opening many connections to the same destination IP and port can use at most one ephemeral port per connection, so a single load-test host tops out around 28,000 connections to one target unless it spreads the load across several source IPs. The same limit applies to any proxy opening backend connections from a single source IP.
Trade-offs and pitfalls
- Quoting 535,000 without assumptions is the weak answer. The strong one says "about half a million at 100% active, several million if buffers are allocated only while active, and I would cap at 60 to 70% of whichever number measurement confirms."
- Mixing GB and GiB silently shifts the answer by about 5% (and more at larger scales); pick one and say it.
- Planning to the memory ceiling leaves nothing for a reconnect storm, which is precisely when memory and CPU spike together.
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.
That is every published Real-Time and Streaming System Design question for Systems Engineer so far. Browse the other topics in this category, or practice this one interactively.