Skip to content
Owais Barkati
BackendData

Broadcasting to 10 million concurrent users

A Kafka fan-out system that got the database out of the notification path

Role
Designed and implemented the broadcast system
Period
–
Concurrent users supported
10M+
Database write load
−70%

baseline → 30% of baseline

Delivery guarantee
Durable

TTL policies and log compaction

One produced event is read independently by every consumer group — adding recipients adds reads, not writes.
Read this diagram as text

A producer publishes a broadcast event once. The event is appended to a Kafka topic, and because a log is read rather than consumed destructively, every consumer group reads that same single append independently at its own offset. This is the asymmetry that makes ten million concurrent users tractable: fan-out costs one write regardless of audience size, where a per-recipient database row would cost ten million. A dynamic filtering stage decides which consumer groups should receive which events, so clients receive only their relevant slice instead of receiving everything and discarding most of it client-side. Filtering dynamically rather than creating a topic per audience segment keeps routing rules changeable without reshaping the cluster topology. Retention policies and log compaction preserve durability: a consumer that restarts resumes from its committed offset and replays what it missed, and for state-like events compaction keeps the latest value per key so a reconnecting consumer can rebuild without querying the database. Database write load fell 70%.

The problem

A broadcast is deceptively simple to describe and punishing to build. One event — an announcement, a status change, a price update — has to reach millions of connected clients, and each client should receive only what applies to it. The naive implementation writes a notification row per recipient. At 10 million recipients, one broadcast becomes 10 million database writes, and the database becomes the system.

That write amplification is the actual problem. It is not that the database is slow; it is that fan-out is being modelled as persistence when it is really routing. Every recipient-row write is paying the cost of durability, indexing and transactional guarantees to answer a question — who should see this? — that does not need any of them.

Constraints

The system had to support 10 million or more concurrent users, deliver in real time rather than on a polling interval, and survive consumer restarts without silently dropping messages. Clients subscribe to different slices of the event space, so delivery had to be selective: broadcasting everything to everyone and filtering client-side would have traded a database bottleneck for a bandwidth one.

Architecture

The design moves fan-out off the database and onto Apache Kafka, with dynamic topic filtering deciding which consumers see which events. Producers publish once. Kafka’s partition and consumer-group model handles the multiplication, which is what a log is good at — a single append is read independently by every consumer group at its own offset, so adding recipients adds reads, not writes.

Filtering dynamically rather than provisioning a topic per audience segment is the part that makes this scale operationally. Topic-per-segment is tidy until segments multiply; then every new audience is a new topic, new partitions, and metadata pressure on the cluster. Filtering lets the routing rules change without reshaping the topology underneath.

Durability without a database

Dropping the per-recipient write raises the obvious objection: what guarantees delivery now? Two Kafka mechanisms carry that weight.

TTL (retention) policies bound how long events remain replayable. A consumer that restarts resumes from its committed offset and replays what it missed, so a brief outage does not become lost messages — while retention stops the log growing without limit.

Log compaction changes what retention means for state-like events. Rather than discarding by age alone, compaction keeps the most recent value per key. For anything representing current state rather than a discrete occurrence, a reconnecting consumer can rebuild from the compacted log instead of querying the database for a snapshot. That is a second path removed from the database.

Results

Database write load fell by 70%. The notification path stopped scaling with recipient count and started scaling with event count, which is the asymmetry that makes 10 million concurrent users tractable: a broadcast to ten million people is still one write.

What I would do differently

I would instrument consumer lag from the start rather than treating it as an operational afterthought. Lag is the honest health signal for a system like this — throughput and error rate both look fine while a consumer group falls progressively further behind, and by the time users report stale data the backlog is already large. I would also pin down the compaction and retention policy per topic class explicitly in configuration, because “durable” means genuinely different things for a transient announcement and for a state update.

Stack

  • Apache Kafka
  • Event streaming
  • Log compaction