Skip to main content

Facebook Live Comments

Sharpened prompt. Design the live comment system for a stream with 5M concurrent viewers producing 100k comments/sec, delivering comments within 2 seconds, running safety classification inline without adding lag, keeping mobile clients from drowning, and anchoring comments to the point in the video they were written about.

The naive design — broadcast every comment to every viewer — is 100,000 × 5,000,000 = 500 billion messages/sec. Say that number early. It establishes that the design is fundamentally about deciding what not to send.

1. Problem framing​

Functional requirements​

  • Post a comment on a live stream; see comments from others in near real time.
  • Show recent comments when a viewer joins mid-stream.
  • Pinned comments, creator replies, reactions.
  • Moderation: block abusive content before it is broadcast.

Non-functional requirements​

PropertyTargetConsequence
Delivery latencyP99 < 2s from post to displayNo batch pipeline in the path
Client load≤ 20–30 comments/sec displayedServer must sample; the client cannot render more anyway
Ingestion100k comments/sec at peakWrite path must be append-only and partitioned
ModerationInline, < 100ms addedSmall models at the edge, not a service call to a GPU cluster
AvailabilityComments degrade before video doesComments are secondary to the stream itself

Back-of-the-envelope​

Viewers: 5M concurrent on one stream, 50M across all streams
Comments in: 100k/sec at peak on the hot stream
Naive fan-out: 100k × 5M = 5 × 10^11 msg/sec -> impossible by ~6 orders of magnitude
Sampled out: 20 comments/sec × 5M viewers = 100M msg/sec -> still enormous,
but tractable with a fan-out tree: 5M / 50k per node = 100 edge nodes
Bandwidth: 100M msg/s × 150 B = 15 GB/s egress on one stream
Connections: 5M WebSockets × ~10 KB state = 50 GB RAM across the edge tier

The second calculation is the one that matters: sampling to 20/sec turns an impossible number into a merely large one, and everything else in the design supports that decision.

2. High-level architecture​

Three structural decisions carry the design: moderation happens before Kafka (so blocked content never enters the distribution system at all), sampling happens once centrally rather than per connection, and the fan-out is a tree rather than a star.

3. Component inventory​

ComponentConcrete choiceWhy this one
IngestStateless HTTP service with per-user rate limitsWrites are independent; no need for the write path to be stateful
ModerationONNX Runtime / Triton with distilled classifiers, co-located with ingestA remote call to a large model would add 200ms+; a distilled model runs in ~5ms on CPU
BusKafka partitioned by stream_idAll comments for one stream are ordered on one partition
Sampler/rankerFlink, 200ms tumbling windowsOne place decides what the world sees; per-connection sampling would be 5M duplicate decisions
HubOne actor per stream (Akka/Pekko, Elixir, or Durable Objects)Serialises ordering and sequence numbering per stream
RelaysGo/Erlang WebSocket nodes, ~50k connections eachCheap concurrency; the tree's leaf tier
Hot bufferRedis list, last 200 commentsJoiners need instant context without touching Cassandra
ArchiveCassandra by (stream_id, ts)VOD replay needs the full ordered log

4. The toughest parts​

4.1 Fan-out: 500 billion messages per second is not a scaling problem​

Why it's hard. No amount of hardware makes the naive number work, and even if it did, the client could not use it. A phone rendering 100k comments/sec would drop frames, burn battery, and show an unreadable blur. The bottleneck is not the server — it is human reading speed, roughly 5 comments/second of comprehension.

Solution — sample centrally, then distribute through a tree.

Sampling first, because it is the bigger lever:

// Flink: one sampling decision per stream per 200ms window, applied to everyone.
comments
.keyBy(c -> c.streamId)
.window(TumblingProcessingTimeWindows.of(Time.milliseconds(200)))
.process(new SelectVisible());

class SelectVisible extends ProcessWindowFunction<Comment, Batch, Long, TimeWindow> {
public void process(Long streamId, Context ctx, Iterable<Comment> in, Collector<Batch> out) {
List<Comment> all = Lists.newArrayList(in); // up to ~20,000 in 200ms
List<Comment> picked = new ArrayList<>();

// Always include high-signal comments regardless of volume.
picked.addAll(filter(all, c -> c.isCreator || c.isPinned || c.isFriendOfManyViewers));

// Fill the rest of the budget by score, with reservoir sampling for fairness
// so an ordinary viewer's comment can still surface.
int budget = 4 - picked.size(); // 4 per 200ms = 20/sec
picked.addAll(topKByScore(all, budget * 2 / 3)); // engagement-ranked
picked.addAll(reservoirSample(all, budget / 3)); // random, keeps it fair

out.collect(new Batch(streamId, picked, all.size())); // include the true total
}
}

Two design points worth defending. Mixing ranked and random selection matters: pure ranking means only already-popular commenters are ever seen, which kills participation. And sending all.size() lets the client render "12,483 comments" alongside the 4 it shows — users get an accurate sense of scale without receiving the volume.

Then the tree, because even 20/sec × 5M viewers is 100M messages/sec:

Hub (1 per stream)
└─► Regional aggregator (×10)
└─► Relay node (×10 per region, 50k connections each)
└─► Viewer WebSocket (×50,000)

The hub sends 10 messages. Each aggregator sends 10. Each relay sends 50,000.
Total hops from hub to viewer: 3. Hub egress: trivial. Work pushed to the leaves,
which is exactly where it can be scaled horizontally.

A star topology would require the hub to write 5M sockets per batch — impossible from one process. The tree makes the hub's job O(number of regions).

4.2 Personalising what each viewer sees, without per-viewer computation​

Why it's hard. Central sampling gives everyone the same 20 comments/sec. But a viewer wants to see their friends' comments and the creator's replies — comments that would almost never survive a global ranking against 100k competitors. Doing per-viewer selection means 5M independent ranking decisions per batch, which reintroduces the original problem.

Solution — two streams, merged at the client.

Stream 1 (global, one computation, tree-broadcast): ~20 comments/sec
Stream 2 (personal, tiny, targeted): friends + creator replies + own

The personal stream is cheap precisely because it is sparse: a viewer typically has a handful of friends watching, and their comment rate is low. Route it by maintaining, per relay node, an index of friend_id → local connections, so when a comment arrives the relay can push it only to the sockets that care.

// Relay node: the global batch goes to everyone; personal routing is a map lookup.
func (r *Relay) onBatch(b Batch) {
frame := encode(b) // encode ONCE, not per connection
for _, conn := range r.conns { conn.WriteRaw(frame) } // zero-copy broadcast
}

func (r *Relay) onPersonal(c Comment) {
for _, conn := range r.friendIndex[c.authorID] { // usually 0–5 connections
conn.Write(encode(c))
}
}

Encoding the global frame once and writing the same bytes to every socket is a meaningful optimisation at 50k connections per node — serialising per connection would dominate CPU.

4.3 Moderation that does not add latency​

Why it's hard. Every comment needs safety classification: harassment, spam, self-harm signals, and — for live content especially — coordinated brigading. A large transformer takes 100–300ms on a GPU and requires a network hop. At 100k comments/sec that is both too slow (blowing the 2s budget when queued) and too expensive (thousands of GPUs for comment classification alone).

Solution — a cascade, cheapest first, with the expensive tier off the critical path.

async def moderate(comment: str, ctx: Context) -> Verdict:
# Tier 1: exact and fuzzy matching. Microseconds. Catches the bulk of spam.
if bloom_blocklist.contains(normalize(comment)):
return Verdict.BLOCK # ~40% of bad content, ~0 cost

# Tier 2: distilled classifier, INT8-quantised, running locally on CPU.
# A 6-layer distilled BERT at ~5ms handles 100k/sec on a few hundred cores.
score = await local_onnx.infer(comment) # p(violating)
if score > 0.95: return Verdict.BLOCK
if score < 0.30: return Verdict.ALLOW # ~95% of traffic decided here

# Tier 3: the grey zone (~5%). Allow it through with reduced distribution,
# and let the heavy model adjudicate asynchronously.
await kafka.send("moderation.review", comment)
return Verdict.ALLOW_LIMITED # excluded from the global sample,
# visible only to the author's friends

ALLOW_LIMITED is the interesting state and it is what makes the cascade work: uncertain content is not blocked (which would generate false positives at enormous volume) and not fully broadcast (which would risk harm). It is shown narrowly while the heavy model decides, and retracted within seconds if the verdict is negative. Retraction is possible because comments are addressed by ID and clients handle a delete message.

Two more layers worth naming: author reputation as a prior (a first-time commenter's grey-zone content is treated more conservatively than a long-standing account's), and brigading detection as a stream-level signal — a sudden spike in similar comments from unconnected accounts is a coordinated attack, which Flink can detect over a 10-second window and respond to by tightening thresholds for that stream specifically.

4.4 Anchoring comments to the video, not the clock​

Why it's hard. Live streams do not arrive simultaneously. A viewer on a low-latency HLS path is 3 seconds behind the broadcaster; one on standard HLS is 30 seconds behind; a DVR viewer who rewound is minutes behind. If comments are ordered by wall clock, the 30-second-behind viewer sees reactions to a goal 30 seconds before the goal happens. Comments become spoilers.

Solution — tag every comment with the video presentation timestamp its author was watching, and render by PTS.

// Client, on submit: attach the playhead, not the clock.
async function postComment(text) {
await api.post('/comments', {
stream_id: streamId,
text,
pts_ms: player.currentTime() * 1000, // where THIS viewer is in the video
client_ts: Date.now(), // kept only for abuse analysis
});
}

// Client, on receive: buffer and release in playback time.
class CommentScheduler {
constructor(player) { this.buf = []; this.player = player; }
ingest(c) { this.buf.push(c); this.buf.sort((a, b) => a.pts_ms - b.pts_ms); }
tick() {
const now = this.player.currentTime() * 1000;
// Show comments whose PTS has passed, with a small display window.
while (this.buf.length && this.buf[0].pts_ms <= now) this.render(this.buf.shift());
}
}

This gives every viewer a coherent experience regardless of their latency, and it makes VOD replay free: after the stream ends, the recorded comments replay against the same PTS timeline, so a viewer watching the archive sees comments appear exactly where they did live.

The edge case to raise: a viewer far behind live receives comments with PTS values in their future, which they must buffer. A viewer ahead (rare, but possible with different CDN paths) receives comments already past — render those immediately in a small catch-up burst rather than dropping them.

4.5 Five million people joining, leaving, and reconnecting​

Why it's hard. A stream's viewer count is not stable — a celebrity going live adds a million viewers in a minute. Each needs a WebSocket, an authorisation check, and a backfill of recent comments. Worse, when a relay node dies, 50,000 clients reconnect simultaneously, and if they all reconnect immediately they take down the next node too.

Solution.

// Client reconnect: full jitter, or a node failure cascades through the tier.
let attempt = 0;
function reconnect() {
const delay = Math.random() * Math.min(30000, 500 * 2 ** attempt++);
setTimeout(async () => {
const node = await discovery.pick(); // load-aware assignment, not round-robin
try { await connect(node); attempt = 0; }
catch { reconnect(); }
}, delay);
}
  • Full jitter on reconnect, capped at 30 seconds. Without it, node loss is contagious.
  • Load-aware node assignment. The discovery service returns the least-loaded relay, and relays advertise a "draining" state so a node under pressure stops accepting new connections rather than degrading everyone on it.
  • Backfill from Redis, not Cassandra. A joiner needs the last ~50 comments; serve them from the hot buffer in one operation. 5M joins in a minute against Cassandra would be a self-inflicted outage.
  • Connection budget per stream. If a stream exceeds capacity, new viewers get comments in polling mode (a GET every 3 seconds against the Redis buffer, cacheable at the CDN) rather than a socket. Degraded but functional, and the video keeps playing — which is the property that actually matters.

That last point generalises: comments must fail before video does. The comment tier should have its own capacity ceiling and shed into a polling fallback rather than competing for resources with stream delivery.

4.6 Ordering and idempotency in a lossy world​

Why it's hard. Viewers see comments in different orders depending on which relay they are attached to and when they connected. Reconnecting clients re-request recent comments and re-render duplicates. And a comment retracted by the moderation cascade must disappear on 5M clients that may or may not have received it.

Solution — a per-stream monotonic sequence number assigned by the hub.

# The hub is the single serialisation point per stream, so sequencing is trivial.
class StreamHub:
def __init__(self, stream_id):
self.seq = 0
def publish(self, comment):
self.seq += 1
comment.seq = self.seq # total order for this stream
self.broadcast(comment)

# Client: dedupe and gap-detect on seq.
def on_comment(c):
if c.seq <= self.last_seq: return # duplicate from a reconnect
if c.seq > self.last_seq + 1:
self.request_range(self.last_seq + 1, c.seq - 1) # fill the gap from Redis
self.last_seq = c.seq
self.render(c)

Sequence numbers give you three things at once: deduplication (a client ignores anything it has already seen), gap detection (a client knows what it missed and can ask for it), and a stable resume point for reconnection. Retraction is then just another sequenced message — {"type": "delete", "id": ..., "seq": ...} — that flows through the same path and is naturally ordered against the original.

Note that ordering is guaranteed within a stream only. Global ordering across streams is meaningless and would require coordination you should refuse to pay for.

5. What breaks first​

EventFirst failureMitigation
Celebrity goes liveConnection storm on relay tierLoad-aware assignment, drain signalling, polling fallback
Comment rate spike (a goal is scored)Sampler window overflowsBudget is fixed per window by construction; only the reported total changes
Relay node loss50k simultaneous reconnectsFull-jitter backoff; capacity headroom of at least one node per region
Moderation model latency spikeIngest backpressureCascade degrades to tier-1 blocklist only; grey zone widens to ALLOW_LIMITED
Redis hot-buffer lossJoiners get no backfillRebuild from Kafka; joiners see comments from now on, which is acceptable
Coordinated brigadingSample fills with abuseStream-scoped threshold tightening; similarity clustering over a 10s window

6. Cheat sheet​

  • The number: 100k × 5M = 5 × 10^11 msg/sec naive. Sampling to 20/sec is the design, not an optimisation.
  • Sample once, centrally — mixed ranked + reservoir so ordinary viewers can still surface; send the true total for display.
  • Fan-out tree: hub → regional aggregators → relays (~50k conns each). Encode the frame once per node.
  • Personalisation: a second sparse stream (friends, creator, self) merged client-side.
  • Moderation cascade: blocklist → distilled ONNX on CPU → ALLOW_LIMITED grey zone adjudicated async and retractable.
  • Timing: anchor to video PTS, not wall clock. Free VOD replay falls out of it.
  • Ordering: per-stream monotonic seq from the hub gives dedupe, gap-fill, and resume.
  • The one-liner: "You cannot deliver 100k comments/sec to 5M people, and nobody could read them anyway — so the system is a central sampler feeding a fan-out tree, and every other decision protects that budget."