Live Sports Scores to Millions
Worked solution

Live Sports Scores to Millions — a worked solution

Goal scored → 30M phones updated in <2s. Pub/sub fanout, edge delivery, thundering herd.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Live Sports Scores to Millions workspace

The problem

A goal is scored at minute 78. Within two seconds, thirty million phones — across continents, on 3G, on Wi-Fi, on the subway — must show the new score before the user feels it as late. The reads are the product. The writes are sparse (a goal happens, not 60 times a second), but the fanout factor is the design pressure: one publish must reach tens of millions of subscribers, and the on-field event rate spikes from 0 to ~5 events/second the instant something happens.

This is not the URL-shortener shape. The hot path is not a key-value lookup against a cache; it is sustaining a 30M-concurrent persistent-connection fleet, holding per-match subscriber lists in memory at the right shard, and absorbing reconnect storms when a PoP flaps. The single most important design decision is "where does the per-match fanout live" — and the answer (per-match actor sharded by hash(matchId), with a tree-relay topology activated for hot matches) determines the rest of the architecture.

The reference architecture

Reference architecture for Live Sports Scores to Millions: 15 components — Mobile / Web client, Snapshot API, GSLB / Geo-DNS, Control Plane, Anycast L4 LB, WSS Edge Terminator, Auth / JWT Issuer, Per-match Fanout Shard, Score Event Spine (Kafka tiered), Ingest + Normalizer, Provider Feed (Sportradar / Genius), Match Catalog, Per-match Replay Buffer, Tracing Collector, Region Lease Coordinator — connected by 23 flows.Mobile / Web clientiOS / Android / browser…Snapshot APIGo service behind CDN; …GSLB / Geo-DNSAWS Route53 traffic pol…Control PlaneGo service + RBAC + aud…Anycast L4 LBBGP anycast + Cloudflar…WSS Edge TerminatorCentrifugo (Go) on bare…Auth / JWT IssuerGo service + JWKS endpo…Per-match Fanout ShardElixir/Erlang GenServer…Score Event Spine (Ka…Apache Kafka 3.7 with t…Ingest + NormalizerGo service + Flink job …Provider Feed (Sportr…Sportradar UOF, Genius …Match CatalogPostgres 15 + Patroni (…Per-match Replay Buff…Redis Cluster with Stre…Tracing CollectorGrafana Tempo + OTel Co…Region Lease Coordina…etcd 5-node Raft cluste…
15 components, 23 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Mobile / Web clientiOS / Android / browser SDK

Holds one persistent WebSocket (or SSE / long-poll fallback) to the nearest edge PoP, subscribes to one or more matches, renders incoming score deltas applying last-revision-wins for the same eventId, persists last_event_id to disk so it survives app kills, and falls back to a snapshot fetch when the replay gap exceeds the buffer or when ws-edge emits a feed_paused control frame.

Why it exists. The user-visible product. Everything else exists to deliver one Protobuf delta into this socket within 2 seconds of the on-field event. Modeled as a first-class node because the *client's* state — its last-event-id, its sticky-cookie session, its decorrelated reconnect backoff, its snapshot-fetch fallback policy — is what determines whether the system survives reconnect storms and failover windows.

When it fails. Naive reconnect on PoP flap (no jitter, fixed delay) turns one 90s blip into a 5-minute self-DOS on the failover PoP. App without persistent last-event-id silently drops the goal that happened during a 6-second 4G handoff. Default UI that fails silent on snapshot-gap (rather than re-rendering the current score from snapshot-api) leaves the user staring at 1-0 forever after they actually saw it become 1-1 on Twitter. UI without feed_paused handling sits silent for the full 30s failover window.

Snapshot APIGo service behind CDN; reads from catalog + replay current-state

Serves GET /v1/matches/{id}/snapshot — returns the authoritative current scoreline for a match. Sized for two access patterns: (a) every cold app-open fetches one snapshot for each subscribed match before opening the WS (1-2 calls per app-open); (b) reconnects whose Last-Event-ID falls outside the replay buffer fetch a snapshot before re-subscribing. Reads match meta from catalog, current-score from replay (XREVRANGE for the latest event), assembles the response in <30 ms p99.

Why it exists. The system has TWO entry points for the question 'what is the current score?': the live WS stream (the dominant path) and the snapshot (for cold connects + replay-gap fallback). Folding the snapshot endpoint into ws-edge would couple snapshot capacity to socket capacity — bad, because the snapshot traffic spikes 20× when the replay cluster cold-restarts (the exact moment ws-edge is also under pressure). Separating it as its own service lets snapshot scale on its own schedule and lets a CDN absorb 90%+ of snapshot traffic with 5-15 second cache TTL.

When it fails. Redis (replay) cold restart during a kickoff: snapshot rate spikes 20×, CDN absorbs most, origin still sees 5× normal. If the snapshot-api is undersized, cold-connecting users see 503 → exponential retry → cascade. Mitigation: aggressive CDN TTL + pre-warmed pod fleet at gameday-forecast count. If the catalog is also down (worst case), snapshot falls back to a stale-in-memory snapshot served as 200-OK with X-Snapshot-Stale: true header so the UI can warn the user.

GSLB / Geo-DNSAWS Route53 traffic policy or NS1 (managed)

Resolves scores.example.com to the IP of the nearest healthy region's anycast prefix. Honors geographic policy (US users → us-east-1, EU → eu-west-1, APAC → ap-south-1) on the steady-state path; during operator-driven failover, control-plane flips the GSLB routing rule to redirect new resolutions away from the draining region within the DNS TTL window.

Why it exists. Anycast handles steady-state routing in microseconds, but anycast cannot do *graceful* regional drain — once a PoP advertises a prefix, every TCP SYN within its catchment lands there. The control-plane failover playbook step 3 (region resume) depends on the GSLB to redirect *new* connections without affecting existing sockets. Without a GSLB layer, the only way to drain a region is BGP withdrawal — which is correct but blunt, and triggers the anycast-hairpin failure mode at the failover PoP.

When it fails. GSLB outage (the whole DNS provider goes down) silently freezes regional steering: new sockets continue to land on their cached-resolution region; if that region is the failed one, connects fail. Mitigated by the 4-provider redundancy. DNS TTL too long (e.g. 5 min) makes failover take effect 5 min late and breaks the 30s RTO claim. DNS TTL too short (e.g. 5s) puts the GSLB on every WS handshake and self-DOSes the DNS provider's edge.

Control PlaneGo service + RBAC + audit log; hosts /admin endpoints

Hosts the POST /admin/lease/transfer, score-ops region drain, score-ops region resume, and the deploy-blackout-gate endpoints. Every operator action is RBAC-gated (operator-IAM-role required), audit-logged with operator identity, justification, and runbook reference, and broadcast via spine to the audit-log consumer. Talks to lease (etcd CAS), anycast-lb (drain command), and gslb (DNS policy flip).

Why it exists. Operator actions during a high-stakes failover (mid-match cutover) require the same rigor as data-plane writes: authenticated, authorized, audited, ideally one-command. A kubectl edit against etcd directly has no audit trail and no RBAC for 'who is allowed to transfer a lease' as distinct from 'who can read etcd'. Separating this as a control-plane service is what makes monthly failover drills repeatable and what keeps a mid-incident cutover from degenerating into un-audited, precondition-free manual etcd writes.

When it fails. Control-plane unavailable during a real incident means manual etcd writes via kubectl — possible but slow, no audit, no preconditions. Mitigated by keeping control-plane on a separate deploy cadence + IAM + cluster from the data-plane, so the same incident is unlikely to take both down. Buggy precondition check that incorrectly says 'cannot drain' during a real cutover is the worst case — must be overrideable with a typed --force flag that is itself audited.

Anycast L4 LBBGP anycast + Cloudflare Magic Transit / AWS Global Accelerator

Terminates the TCP connection at the geographically nearest PoP via BGP anycast, runs an L4 admission-control token bucket on Upgrade: websocket requests, then forwards the raw TCP stream to a co-located WS-edge process. Issues Retry-After with decorrelated jitter when the local PoP is over its admission ceiling. Receives drain commands from control-plane to stop accepting new Upgrade traffic without affecting existing sockets.

Why it exists. Pinning the TLS handshake at the PoP cuts ~80 ms of cross-continent RTT off cold-connect latency, and the L4 admission cap is the *only* defense against a reconnect storm landing on a flapped-failover PoP at 5× normal rate. Putting the admission gate at L4 (before WS upgrade) means we shed at 200 µs/decision instead of 20 ms of TLS + Lua decision cost.

When it fails. BGP withdrawal in one PoP (Cloudflare Jul 17 2020) re-routes 5M sockets through one neighbor PoP that immediately saturates — anycast hairpin. Mitigation is pre-announce backup /24s and per-PoP admission caps that 503 cleanly. If the admission cap is too high, you cascade to ws-edge OOM (see node ws-edge). If the admission cap is too low, you self-DOS during legitimate kickoff spikes — the cap MUST be sized to 90-percentile gameday concurrency, not steady-state.

WSS Edge TerminatorCentrifugo (Go) on bare metal — primary path (Cloudflare DO + WS Hibernation is the rejected alternative; see tradeoffs)

Terminates TLS, completes the WebSocket Upgrade, validates the bearer JWT against a local JWKS cache, looks up the match entitlement in the catalog, registers the subscription on the right per-match relay shard, optionally replays missed events from the per-match Redis Stream (if the client supplied Last-Event-ID), and then streams every inbound publish from its subscribed match-shards into the open socket. On lease-loss notification from a match-shard, emits a feed_paused WS control-frame to all affected subscribers so the UI can render 'Reconnecting…' instead of silent.

Why it exists. Splitting this from the match-fanout layer (match-shard) is the load-bearing decoupling: per-match fanout has to live wherever the *publishing* topic lives (one node per match), but TLS termination has to live wherever the *subscriber* lives (200+ PoPs globally). Folding them together would force each match's fanout shard to terminate sockets from every continent, which is exactly the topology that took Discord's pre-Manifold cluster to 900 ms–2.1 s hot-guild fanout.

When it fails. Hits the fd ceiling (default 1024 — must raise fs.nr_open to 4M and net.ipv4.ip_local_port_range accordingly) and accept_eagain_total climbs while new connects silently spin. TLS-resumption stampede after key rotation forces full handshakes and CPU melts. Slow-reader backpressure without the 32 KB cap OOM-kills the box and takes 200 K sockets with it (then triggers a reconnect storm — see anycast-lb). The coalescing window dropping a VAR-correction (the SEV-1 bug class) is what the (matchId, eventId, revision) keying defends against; without that keying, revision_collapsed_at_edge_total would silently violate the 'wrong score = error budget 0' promise.

Auth / JWT IssuerGo service + JWKS endpoint + KMS for signing keys

Issues 15-minute RS256 JWTs when the mobile app authenticates (OAuth, Apple/Google sign-in, anonymous device id). Publishes the public JWKS endpoint that every ws-edge instance caches with TTL 15 min. Also serves a *signed JWKS pin* (the JWKS document signed by a long-lived root key) so ws-edge can verify offline for up to 24 h if auth is unreachable — closes the AWS Builders' Library 'avoid fallback' anti-pattern.

Why it exists. Existing as a *separate, off-the-WS-path* tier is the load-bearing choice. If ws-edge called back to auth on every WS open, kickoff time (when 5M people simultaneously open the app) would self-DOS the auth tier, which would cascade to every other product sharing it. Short-lived JWT + local JWKS cache means auth only sees traffic when (a) a user first signs in or (b) JWKS keys rotate.

When it fails. JWKS rotation collides with edge cache invalidation (rare, but the 15-min TTL is the worst-case staleness window): some edge boxes see signatures fail until they re-fetch. Auth outage during JWKS rotation extends that window — the signed JWKS pin lets ws-edge continue verifying any non-rotated token for 24h. If ws-edge is misconfigured to fall back to remote verification on JWKS miss with exp-backoff and no circuit-breaker, an auth outage during a kickoff cascades — keyChoice (e) on ws-edge (fixed-retries=1, fast-fail) is the design's defense.

Per-match Fanout ShardElixir/Erlang GenServer per matchId (Discord-style) primary path

Owns the in-memory subscriber list for ONE matchId (or one slice of a hot matchId), consumes that match's partition from the Kafka spine, appends every event to the per-match Redis replay stream, checks the lease coordinator (etcd) to confirm it still owns the active-region token, and pushes the event into the WS-edge processes that have registered subscribers. For hot matches (predicted concurrent > 250K), the shard becomes the ROOT of a 2-tier relay tree: root publishes to ceil(concurrent / 15K) leaf relays (grouped by region), each relay fanning its ~15K subscribers' WS-edge connections — ~100 relays for a 1.5M-subscriber marquee match, exactly Discord's Manifold pattern.

Why it exists. Per-match consolidation is the only way the cross-PoP fanout-vs-receive ratio stays sane. Without it, every edge PoP would consume the full publish stream (60 events/s × 40 matches = 2.4 K events/s) and filter locally to its subscribers — that wastes 99% of the inter-PoP bandwidth and pins the spine consumer-group fan-out by N(PoPs). With per-match shards, only PoPs that have a subscriber for matchId-X get a connection to its shard.

When it fails. Hot match (the India-vs-Pakistan / WC final case) takes a single shard's CPU to 100%; without relay-tree activation per-event fanout climbs from 30 ms to 1.5 s — Discord pre-Manifold. Crashed shard takes its in-memory subscriber list with it; WS-edges must re-register subscriptions (handled gracefully via spine re-consume + replay). If lease-check is fail-OPEN instead of fail-CLOSED, a botched region failover produces duplicate sends. If consumer-offset commit is misordered (committed before publish-ack), lease-loss silently skips events.

Score Event Spine (Kafka tiered)Apache Kafka 3.7 with tiered storage to S3 (KIP-405); MirrorMaker 2 cross-region

Carries two topics on one cluster: score-events (the canonical event stream, partitioned by matchId so per-match ordering is the load-bearing primitive) and audit.ops (control-plane action log). Hot tier (NVMe) keeps the last 6 hours; tiered storage to S3 retains score-events 14 days and audit.ops 7 years (compliance). MirrorMaker 2 replicates score-events cross-region into the standby cluster — not modeled as separate spine nodes to keep the diagram readable, but the lag SLO mm2_replication_lag_p99_seconds > 10s is the gating signal for failover safety.

Why it exists. Without a durable backbone, every consumer that wants score events (match-shards, audit log) would need a direct connection to ingest, with N² coupling and no replay story. The dedupe boundary is the Flink ingest stage, not the broker: ingest collapses duplicate provider events by the canonical (matchId, eventId, revision) before producing. Kafka's enable.idempotence=true only suppresses producer-retry duplicates (per PID+sequence within one session) — the broker does NOT dedupe by an application key. So business-level dedupe happens once, in ingest, and every consumer reads an already-deduped stream.

When it fails. Consumer lag on hot match's partition climbs past 30 s during a goal-spike if a downstream consumer restarts; alert on kafka_consumer_lag_seconds{topic='score-events'} > 5s. Poison message (corrupted provider payload) crash-loops the consumer-group and halts fanout silently — DLQ + schema validation at ingest is the safety net. Stretched-cluster failure mode is why we use MirrorMaker, not stretched. MM2 lag tail > 30s during a network blip means RPO 0 via spine fails over to RPO 0 via provider replay (see provider-feed.keyChoices).

Ingest + NormalizerGo service + Flink job for windowed dedupe

Maintains long-lived WSS connections to two redundant provider feeds (Sportradar primary, Genius Sports / Stats Perform secondary), normalizes each provider's wire format into the canonical Protobuf score event (assigning a provider-independent eventId), dedupes by (matchId, eventId, revision) over a 500 ms tumbling window in Flink so whichever provider delivers a given (eventId, revision) first wins and the other's copy is dropped — that IS 'take the faster of the two providers'. (Keying on providerId would defeat it: the two providers' copies would carry different keys and never dedupe against each other.) Validates against the catalog (matchId must exist + match must be in live state), and produces to the Kafka spine with idempotent-producer guarantees.

Why it exists. Provider feeds are messy. Sportradar can emit the same goal under two different eventIds when its scout corrects an early call; Stats Perform delivers via a redundant pair that occasionally double-publishes. Doing dedupe + normalization HERE, once, behind the spine means every downstream consumer reads a single authoritative event stream.

When it fails. Both providers go down → 503 from ingest → spine has no new events → users see frozen scores. Detection: provider_event_age_seconds_p99 / expected_cadence_for_sport > 3 for 60s (sport-aware threshold; flat 8s threshold over-fires on football, under-fires on tennis). Mitigation: surface 'live feed delayed' in the UI. Watermark too tight drops legitimate late events from the slower provider. Wrong-schema event corrupts the Flink job mid-stream — DLQ at the Flink boundary, not crash the worker.

Provider Feed (Sportradar / Genius)Sportradar UOF, Genius Sports Live, Stats Perform OPTA

Streams normalized in-venue events (goal, card, substitution, clock tick, possession change) to our ingest via a persistent WebSocket or AMQP feed. Pricing is per-event-per-licensee; the provider's SLA covers latency from venue-scout-tap to provider-edge, typically 1-3 s.

Why it exists. We don't have scouts in every Premier League / NFL / NBA stadium and we never will. Sportradar + Genius + Stats Perform collectively cover ~all top-tier sports globally and re-license. Modeled as kind: external so the engine routes it through third-party-boundary semantics: circuit-break it, cache it, never block on its tail.

When it fails. Stadium uplink flap (weather, primary cellular carrier outage) silently delays the feed for 5-20 s. Provider's own backend incident can take an entire sport offline globally — happened in published Sportradar status incidents 3-5× / year. VAR-rollback or human-scout correction arrives out of order and the dedupe key MUST honor (eventId, revision) semantics or we silently lose the correction. Provider-replay catch-up exceeding the 500ms Flink watermark needs recovery-replay-mode or events get dedupe-dropped.

Match CatalogPostgres 15 + Patroni (HA primary + 2 sync followers)

Authoritative store for: (a) match schedule + status (matchId → start, sport, league, teams, home_region, current_score materialized column updated by a spine consumer for snapshot-api fallback); (b) entitlements (this user can watch Premier League but not La Liga); (c) geographic blackout rules; (d) canonical channel-name registry. The current-score column makes snapshot-api resilient to a replay outage.

Why it exists. Small-but-load-bearing relational data; read on every WS-open + every ingest publish. Single source of truth for entitlement is the only defense against geographic-blackout litigation.

When it fails. Primary failure during a kickoff: new WS-opens stall the entitlement check until Patroni promotes — 30 s gap during which new connects fail with the coordinated Retry-After. Existing sockets unaffected. If entitlement check fails-OPEN on timeout (the anti-pattern), unblacked-out matches serve to blackout regions = legal liability. Always fail-CLOSED.

Per-match Replay BufferRedis Cluster with Streams (XADD/XRANGE), AOF persistence

One Redis Stream per matchId, keyed replay:{matchId}, holding every score event for the match's lifetime. When a client reconnects with Last-Event-ID=X, WS-edge runs XRANGE replay:{matchId} (X + to fetch the gap and replays the deltas before subscribing to live updates. If the client's Last-Event-ID is older than the buffer head (very long disconnect or out-of-buffer match), the edge returns a snapshot-required signal so the client refetches the current score from snapshot-api.

Why it exists. The 'mobile loses signal for 30 s' case is by far the most common. Without a replay buffer, the user reconnects and silently misses every event during the gap, including potentially the goal they cared about. Redis Streams over a custom buffer because XRANGE + consumer-group semantics are exactly the shape we need.

When it fails. Cluster cold restart wipes every replay buffer — every reconnecting client during warmup has to snapshot-fetch (correct, but snapshot-api takes 20× spike — sized accordingly). AOF replay on restart is fast (~2 min) but during that window XADD waits. Per-match key evicted under memory pressure if maxmemory-policy != noeviction — set noeviction + alert on used_memory_pct > 80. Hot-match XRANGE during a reconnect storm pins one Redis shard's CPU; mitigated by ws-edge in-process LRU of last 200 events (most reconnects don't go to Redis at all).

Tracing CollectorGrafana Tempo + OTel Collector; tail-based sampling at the collector

Receives OTel spans from ws-edge and match-shard over gRPC, applies tail-based sampling at the collector (Discord-published strategy: 100% sample for fanout=1, ~0.1% for fanout > 1M to keep the volume sane), and stores in Tempo. Critically, supports an always-sample-on-incident override triggered by SEV-1 declaration: during an incident window the collector forces 100% sampling for the affected matchId(s) so on-call can do per-subscriber forensics.

Why it exists. Bundling tracing into a data-plane flow means an outage in that flow takes tracing with it — exactly the wrong moment. Separating tracing as its own node, on its own collector + storage, means observability survives data-plane outages. Inverse-fanout sampling is non-trivial control-plane code that has to live in the collector, not in every service.

When it fails. Collector queue overflows during a SEV-1 (100% sampling × goal-spike × hot-match fanout); spans drop on the floor exactly when on-call needs them. Mitigation: per-tenant span budget + DLQ + sample reduction fallback (start dropping low-priority service spans before the hot-match's). Storage backend (S3 / Tempo) outage means current incident is debugged from logs + metrics only — degraded but not blind.

Region Lease Coordinatoretcd 5-node Raft cluster spanning regions (3 regions × ≥1 member)

Holds the linearizable (matchId → active-region, fencing-token) mapping. Every match-shard polls this on lease.refresh.interval = 5 s and refuses to publish into the fanout tree if its lease is not current. Operator-driven failover (via control-plane) atomically transfers the lease via a Raft-committed Compare-And-Swap; old region's match-shards see lease-loss on the next refresh and drop in-flight events.

Why it exists. Without a CP fence-token store, the failover ramp is unsafe by construction. GitHub Oct 21 2018 is the canonical post-mortem: a 43-second network partition, an automated cross-DC MySQL failover with no fence, and both datacenters accepting writes for the same data. The lease coordinator is *the* component that makes our active-active topology survivable.

When it fails. Quorum loss (3 of 5 nodes down) → lease cannot be renewed → all match-shards fail-CLOSED within 30 s. Detection: etcd_has_leader == 0 > 10 s → P0; lease_renewal_lag_seconds > 10s → P1 (early signal that next refresh will miss TTL and dual-active is imminent). Network partition between regions: lease has already been transferred via Raft (which survives by losing quorum on the minority side); old region fails-CLOSED correctly. Named anti-pattern: using etcd v2 API or TTL-only leases without lease grants.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Clarifications worth surfacing on day 1:

  • Scope of "live scores": just the scoreline (1-0, then 1-1), or also play-by-play (commentary, possession ticks, expected-goals, win-probability)? Drives the publish QPS by 10×.
  • Latency target: end-to-end (goal-line to phone glass) or our-system-only (our ingest to our edge)? The provider feed adds 1-3 s we cannot eliminate.
  • Audience: is the marquee-match concurrent peak 5M (Super Bowl, US-only), 30M (IPL final, India-heavy), or 100M (World Cup final, global)? Different topologies flip at different points.
  • Multi-region: active-active or active-standby? Active-active multiplies the lease/fencing complexity but is the only honest answer at >5M global concurrent.
  • Provider feed: licensed (Sportradar / Genius / Stats Perform) or in-house scouts? In-house puts us in the venue with all the cellular-uplink failure modes.
  • Replay window: how long can a user be offline and still expect to "see what happened"? 30 s? 6 min? Full match? (Note: differs by sport — Test cricket runs for days.)
  • Geographic blackouts: NFL blackouts, EPL UK 3pm rule — entitlement check on every connect is non-negotiable if these are in scope.
  • Correction protocol: the venue scout or VAR rolls back the goal 5 s later — how do we render "GOAL!" then "Goal disallowed"? Or do we delay all renders by 5 s as a safety margin?

02Functional reqs

What must this system actually do?

  • Subscribe to one or more matches via a persistent socket (WS primary, SSE fallback, long-poll legacy).
  • Receive every score event for a subscribed match within 2 s of the on-field event.
  • Resume after a network blip via Last-Event-ID — no missed goals up to the replay window (match-duration-aware: 2h for football, 12h+ for cricket).
  • Receive event_retracted (VAR rollback, scout correction) and render last-revision-wins.
  • Honor geographic blackouts (NFL local-market rules, EPL UK 3pm-Saturday rule, etc.).
  • Operator can drain a region and fail over the active publisher within 30 s without sending duplicate events.
  • During a regional failover, the user UI receives a feed_paused control frame and renders a 'Reconnecting…' state rather than sitting silent.

03Non-functional

What must it promise about speed, uptime and correctness?

  • End-to-end latency (goal-line → phone glass): P99 < 2.5 s, P99.9 < 4 s. Our-system-only (ingest → edge flush): P99 < 800 ms.
  • Availability on the read/subscribe path: 99.95% monthly (21.6 min budget). Asymmetric: stale-by-30s = warn; wrong score broadcast = SEV-1 with error budget = 0 (fail-closed, drop the event, never lie).
  • Durability of the spine event log: 14 days hot-tiered (S3 KIP-405); zero in-flight loss with acks=all, min.insync.replicas=2.
  • Throughput at peak (reconciled):
  • Publish QPS, system-wide: 60 events/s (40 concurrent matches × 1.5 events/s/match).
  • Steady-state aggregate outbound msg rate: 90M msgs/s (60 × 1.5M recipients/event).
  • Goal-spike (single event in the marquee match): 1.5M outbound messages in a ~5s window (30M concurrent × 2 subs/user / 40 matches).
  • Egress (reconciled):
  • Steady-state aggregate: 130 Gbps (90M msgs/s × 180 B Protobuf delta).
  • Goal-spike burst: ~54 MB/s ≈ 0.43 Gbps added on top of baseline for the 5s window (1.5M msgs × 180 B).
  • With JSON instead of Protobuf-delta this would be 360 Gbps steady-state — Protobuf is a ~2.8× egress cost reduction.
  • Reconnect storm resilience: survive a PoP flap that dumps 4M sockets onto the failover PoP without that PoP cascading.
  • Multi-region: active-active with one home-region-per-match (consistent-hash on matchId), fenced via etcd lease — RTO 5 s on home-region match-shard process failure, RTO 30 s on full home-region failover, RPO 0 on score events (provider replay covers the MirrorMaker lag gap; see §multi-region).
  • Security: TLS 1.3 everywhere; short-lived bearer JWT for client auth; mTLS between internal services; geographic blackout enforcement is fail-closed (never serve when in doubt).

04Capacity estimation

How much load and data does this have to hold?

Inputs (defaults — tune for the specific gameday):

  • Peak concurrent (global, marquee match): 30M
  • Avg matches subscribed per user: 2
  • Publish events/s/match (mixed sport avg): 1.5
  • Concurrent live matches at peak: 40 (NFL Sunday 1pm slate)
  • Avg event payload (Protobuf delta): 180 B
  • Replay window: 7200 s (default for 90-min sports; match-duration-aware per-key)

Derivations (every step shown — the canonical's numerical claims must be internally consistent):

1. Publish QPS: 1.5 events/s/match × 40 matches = 60 events/s system-wide. Trivial — Kafka eats this. (The publish rate is small; the fan-out is huge.)

2. Per-event audience: average user subscribes to 2 of 40 matches → per-event audience = 30M × 2 / 40 = 1.5M recipients/event. This is what hits the WS-edge fleet per event.

3. Steady-state aggregate fan-out msg rate: 60 events/s × 1.5M recipients/event = 90M msgs/s. This is the LOAD on the WS-edge fleet.

4. Steady-state egress: 90M msgs/s × 180 B × 8 = 130 Gbps. With JSON (500 B) this would be 360 Gbps — Protobuf is a 2.8× cost reduction.

5. Goal-spike (single-event broadcast): one goal in the marquee match → audience = 1.5M users (those subscribed to that match). Delivered in a ~5 s window → 1.5M × 180 B / 5 s = ~54 MB/s = 0.43 Gbps added burst. The WS-edge per-socket coalescing buffer prevents this from compounding within the same client.

6. PoP-level egress headroom: 130 Gbps × 0.8 / 30 top PoPs = 3.5 Gbps / top PoP. A 25 Gbps NIC handles this; even 10 Gbps with N+1 instances per PoP works. Socket cap follows the same 80/30 split: 30M × 0.8 / 30 = 800K sockets / top PoP (≈4 ws-edge boxes). The per-PoP cap published as pop_admission_cap = 800K sockets, 3.5 Gbps egress, 50K new-opens/s.

7. WS-edge box count: Centrifugo benchmarks at ~200-250K WS/box on 32-core 64 GB. At 200K/box: 30M / 200K = 150 boxes minimum; run 200 for N+1 + 30% headroom. Spread across ~280 PoPs ≠ 280 boxes — top-30 carry 80%.

8. Match-shard count: active shards = N(live matches) + N(hot-match relays). Cold match = 1 shard. Hot match (>250K concurrent) activates a 2-tier relay tree: 1 root + ceil(concurrent / 15K) relays — the matchShardCount derivation gives 1.5M / 15K = ~100 relays for the marquee match. NFL Sunday: 40 matches × 1.2 = ~50 base shards + ~100 relay shards = ~150 total. Node has replicas: 160 for headroom.

9. Replay buffer footprint: per-match per 2h = 1.5 × 7200 × 180 = ~2 MB. 40 concurrent matches × 2 MB = 80 MB total for 90-min sports. Cricket (Test, 3 days) adds ~250 MB per concurrent Test match — still trivial. Drives per-sport MAXLEN (= replay TTL × event rate × 2 headroom: ~20K for football, ~1.4M for a Test) — a flat 10K would trim football mid-window at 10.8K events / 2h.

10. Spine partition count = 256 with a custom partitioner that direct-routes matchId → partition from the live-set table in catalog. Naive hash(matchId) % 256 with 40 matches has 95% birthday-paradox collision probability — would destroy per-match ordering isolation. Custom partitioner is a 50-line Go function reading from catalog.

11. etcd lease store sizing: one key per matchId, ~256 B/key × ~256 matches = 64 KB total. Trivial.

Architectural flip-points (with the rejected alternative):

  • At ~250K WS/box Nginx+Lua tops out — flip to Centrifugo (Go) or Erlang/Phoenix.
  • At ~250K concurrent per match the single fanout-shard saturates — flip to the 2-tier relay tree (Discord Manifold).
  • At >5M global concurrent single-region stops working — flip to multi-region active-active with per-match home-region fencing.
  • At >100M msgs/s outbound per-recipient enqueue is the wrong shape — flip from per-recipient broadcast to topic-broadcast (we are at 90M/s — at the threshold).

05API design

What does the outside world call, and what comes back?

# Subscribe via WSS (the dominant entry point)
GET /v1/stream
Host: scores.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: ...
Authorization: Bearer <15-min JWT>
Last-Event-ID: 78145327                 # optional, for reconnect

101 Switching Protocols
Set-Cookie: edge_session=<signed cookie>; SameSite=Strict; Secure; HttpOnly

After the upgrade, the client sends control frames:

// Client → server (subscribe to a match)
{ "op": "subscribe", "match_id": "epl:LIV-ARS-20260517", "since_event_id": 78145327 }

// Server → client (score event)
{ "type": "score_event", "match_id": "epl:LIV-ARS-20260517",
  "event_id": 78145328, "revision": 1,
  "minute": "78:14", "score": { "home": 1, "away": 1 },
  "kind": "goal", "scorer": { "team": "ARS", "player": "Saka" } }

// Server → client (correction — VAR rollback, ALWAYS emitted immediately
// regardless of coalescing window for the prior revision)
{ "type": "score_event", "match_id": "epl:LIV-ARS-20260517",
  "event_id": 78145328, "revision": 2,
  "kind": "goal_retracted", "reason": "var_offside" }

// Server → client (control — emitted when match-shard fails CLOSED on lease loss)
{ "type": "feed_paused", "match_id": "epl:LIV-ARS-20260517",
  "reason": "regional_failover_in_progress",
  "snapshot_url": "/v1/matches/epl:LIV-ARS-20260517/snapshot",
  "retry_after_seconds": 5 }

// Server → client (replay overflow — client's last_event_id older than buffer)
{ "type": "snapshot_required", "match_id": "epl:LIV-ARS-20260517",
  "snapshot_url": "/v1/matches/epl:LIV-ARS-20260517/snapshot" }
# Snapshot fallback (served by snapshot-api; the CDN edge worker validates the
# JWT + derives the geo-IP region BEFORE cache lookup; key = (match_id, geo_region))
GET /v1/matches/epl:LIV-ARS-20260517/snapshot
Authorization: Bearer <15-min JWT>

200 OK
Cache-Control: private, max-age=10, stale-while-revalidate=30
{ "match_id": "...", "as_of_event_id": 78145400,
  "score": { "home": 1, "away": 1 }, "minute": "78:14", "status": "live" }
# Operator failover (control-plane only, RBAC-gated)
POST /admin/region/failover
Authorization: Bearer <operator JWT, role: region-ops>
{ "from_region": "us-east-1", "to_region": "eu-west-1",
  "justification": "INC-2026-0042 — us-east-1 DR drill" }

# Wraps drain → lease transfer → DNS flip in one transaction.

Why WS over SSE primary: WS is bidirectional (client subscribe / unsubscribe control frames; server score-event push), more efficient header overhead, better mobile-OS support. SSE is the corporate-proxy fallback.

06Data model

What gets stored, and what is it looked up by?

The data model is deliberately small. The system is a fanout fabric, not a database.

1. Match catalog (Postgres):

fieldtypenotes
match_idvarchar(64)PK, e.g. epl:LIV-ARS-20260517
sportvarchar(16)football, nba, nfl, cricket, tennis
home_team_id / away_team_idvarchar(16)
start_tstimestamp
statusenumscheduled / live / final / postponed
home_regionvarchar(16)consistent-hash home for the matchId
current_scorejsonbmaterialized by a spine consumer; for snapshot-api fallback when replay is down
replay_ttl_secondsintper-sport: 7200 (football), 14400 (NBA), 86400+ (Test cricket)

Plus tables for entitlement (tenant × sport × region) and blackout (geographic rules).

2. Score event (Kafka spine, Protobuf):

message ScoreEvent {
  string event_id = 1;          // canonical, monotonic per match (assigned at ingest)
  string match_id = 2;          // home-region-routable
  uint32 revision = 3;          // bumped on correction; client renders last-revision-wins
  string provider_id = 4;       // sportradar / genius / stats — provenance, NOT the dedupe key
  uint64 provider_seq = 5;      // provider-native sequence — ordering/gap-detection only
  google.protobuf.Timestamp venue_ts = 6;  // PTP-disciplined at venue
  google.protobuf.Timestamp ingest_ts = 7; // our ingest timestamp
  oneof body {
    ScoreUpdate score = 10;
    ClockTick clock = 11;
    GoalRetracted retracted = 12;
    LineupChange lineup = 13;
  }
}

Dedupe key = (match_id, event_id, revision) at ingest (the canonical event_id is provider-independent, so both providers' copies of the same on-field event collide and the faster wins); same key at ws-edge for the coalescing window (NEVER collapses across revisions).

3. Replay stream (Redis Streams): replay:{match_id} → list of {event_id, payload} keyed by stream-id; per-sport MAXLEN (replay TTL × event rate × 2 headroom); per-key TTL set from catalog.replay_ttl_seconds on first XADD.

4. Lease (etcd): /leases/match/{match_id} → {active_region, fencing_token, lease_ttl_30s}. Linearizable Raft writes.

07High-level design

Which components handle a request, and in what order?

Architecture summary — the event-flow direction first, subscribe-flow second:

Event flow (provider → user):

  1. Provider feed (Sportradar + Genius, dual-source) pushes the on-field event over a persistent WSS to ingest.
  2. Ingest normalizes to the canonical Protobuf (assigning a provider-independent eventId), runs Flink-based dedupe over a 500 ms window keyed on (matchId, eventId, revision) (so the faster provider wins — NOT keyed on providerId), validates the matchId against the catalog, and publishes to the spine with enable.idempotence=true, acks=all. Per-provider correction semantics are handled here (Genius events without explicit revision get manufactured-revision).
  3. The spine uses a custom partitioner that direct-routes matchId → partition (avoiding the 95% birthday-paradox collision rate of naive hash).
  4. Match-shard consumes its match's partition, appends every event to the per-match replay stream, checks the lease coordinator to confirm it owns the active-region token (fail-CLOSED on loss), and pushes the event downstream into the WS-edge processes that have registered subscribers. The consumer offset is committed ONLY after publish-ack from ws-edge AND with idempotent downstream fanout ((matchId, eventId, revision) dedupe at ws-edge). For matches above 250K concurrent, the shard activates a 2-tier relay tree.
  5. WS-edge holds the persistent subscriber socket, applies the (matchId, eventId, revision)-keyed coalescing window (corrections NEVER collapse), and flushes the Protobuf delta as a WS frame.
  6. Client renders the delta with last-revision-wins, persists last_event_id to disk.

Subscribe flow (cold connect):

  1. Client resolves scores.example.com via gslb (DNS TTL 30s) to the nearest healthy region's anycast prefix.
  2. Client opens WSS to anycast-lb (BGP anycast to the geographically nearest PoP within the region).
  3. Anycast-lb runs L4 admission control on WS-upgrade, then forwards raw TCP to a co-located ws-edge.
  4. WS-edge validates the bearer JWT against its auth-fetched JWKS local cache + signed JWKS pin (24h offline tolerance); on cache miss falls through to auth with fixed-retries=1 (fast-fail, no exp-backoff).
  5. WS-edge reads catalog for entitlement + blackout check (fail-CLOSED with coordinated Retry-After).
  6. WS-edge registers the subscription on the correct match-shard.
  7. If client supplied Last-Event-ID, ws-edge does XRANGE against replay for the gap; if the gap is older than the buffer, returns snapshot-required → client fetches from snapshot-api.

Control plane (failover, drain, deploy gates):

  1. Operator hits control-plane /admin/region/failover (RBAC-gated, audit-logged).
  2. Control-plane issues drain to anycast-lb, transfers lease via Raft CAS to lease, flips DNS policy on gslb.
  3. Match-shard sees lease-loss on next 5s refresh, fail-CLOSED, emits feed_paused to all subscribed ws-edges, which propagates to clients.

Observability: ws-edge and match-shard push OTel spans to tracing (tail-based sampling + always-sample-on-incident).

Multi-region: active-active. Per-match home-region by consistent-hash. MirrorMaker 2 replicates spine cross-region. Lease coordinator's etcd cluster spans 3 regions for cross-region quorum.

08Deep dives

Which part breaks first, and what do you do about it?

1. The Bieber match — single-shard saturation.

When the marquee match concentrates 60-80% of concurrent subscribers, the per-match-shard pattern breaks because a single shard's CPU saturates at ~250K-500K direct subscribers (Discord saw 900 ms–2.1 s on hot guilds pre-Manifold). Our analog: the 2-tier relay tree activated when predicted concurrent > 250K. Root shard → ceil(concurrent / 15K) region-grouped relays at ~15K subscribers each — ~100 relays for a 1.5M-subscriber marquee match (the matchShardCount derivation), ~1,200 for the 60%-of-fleet India-Pakistan case. The TRIGGER for activation is the leading indicator predicted_subscriber_count_from_fixture_calendar > 250K from the gameday-forecast service (a separate batch ML job emitting predictions to a known endpoint at kickoff-10min), NOT the reactive shard_fanout_p99 > 1s that fires after the shard is already saturated. Without the gameday-forecast service, "pre-warm" is wishful aspiration — explicitly named as an upstream dependency.

2. Reconnect storms — the failover-PoP cascade.

A PoP flaps. Its 5M sockets retry on the standard EventSource 3-second default — no jitter. All 5M reconnects arrive at the failover PoP within a 100 ms window. The failover PoP's TLS-handshake CPU saturates, accept queue overflows, new connects fail, and the failover PoP joins the flap. This is the SRE Workbook Ch.22 cascade.

The fix is in three places: (a) decorrelated-jitter in the client reconnect (base 1s, cap 30s, AWS Builders' Library); (b) L4 admission-control token bucket at anycast-lb sized to 1.2× the p90 gameday-forecast peak (steady-state sizing would self-DOS legitimate kickoff spikes); (c) server-issued Retry-After with jitter computed at the server, so even misconfigured clients honor it.

3. VAR-rollback — the coalescing-window invariant.

The provider fires GOAL at T=0 (eventId=X, rev=1). Inside ws-edge's 200ms coalescing window — say at T=180ms — the provider fires CORRECTION (eventId=X, rev=2, kind=goal_retracted). The naive "same-event coalescing" would silently collapse them; whichever direction the collapse goes, one event is silently dropped — violating the "wrong score = error budget 0" SEV-1 promise. The design's defense is keying the coalescing window by (matchId, eventId, revision) and emitting rev N+1 IMMEDIATELY even mid-window for rev N. The SLO revision_collapsed_at_edge_total > 0 in 5min is P0 — it verifies the invariant holds.

4. Consumer-offset commit on lease loss — the invariant that prevents the dual-active footgun.

Match-shard consumes from spine continuously. When the shard sees lease-loss mid-batch, when does it commit its consumer-group offset? Both naive answers are wrong differently: commit-before-publish → new-region's shard skips unpublished events (silent loss); commit-without-idempotent-downstream → shard restart re-publishes (duplicates). The correct answer is commit only after publish-ack from downstream ws-edge AND with idempotent downstream fanout via (matchId, eventId, revision). This is exactly the GitHub Oct 21 2018 failure shape — a cross-DC failover with no fencing left both sides accepting writes — and the design names the explicit invariant to defend against it.

5. Provider feed delayed — the latency-budget gap.

The provider's SLA is 1-3 s from venue-scout-tap to provider-edge. Our budget is 2.5 s end-to-end. That's a structural conflict — we can never deliver "live" faster than our provider. Mitigation: (a) dual-provider feeds with take-faster semantics; (b) surfacing feed_delayed=true in UI when provider_event_age_seconds_p99 / expected_event_cadence_for_sport > 3 (sport-aware threshold — flat 8s threshold over-fires on football's 30-60s cadence and under-fires on tennis's 5-10s cadence).

6. The feed_paused UX during failover — eliminating the 30s silent blackout.

When match-shard loses its lease during operator failover, it stops publishing. WS-edge sees no events for 30s. Without the feed_paused control frame, connected subscribers see nothing — UI says "Live: 1-0" while Twitter is updating to 1-1. This is exactly the "Twitter beat us by 30s" P0 threshold. The design's defense: match-shard emits a lease-loss notification to its connected ws-edges; ws-edge emits feed_paused WS frames; client UI renders "Reconnecting…" and arms a snapshot-fetch from snapshot-api. No silent blackouts.

09Trade-offs

What did this design cost, and what breaks at 10×?

What we accepted, with the rejected alternative:

DecisionRejected alternativeWhy
Per-match actor (GenServer / DO) sharded by hash(matchId)Single Kafka consumer group per region, filter-at-edgeFilter-at-edge multiplies inter-PoP bandwidth by ~280; per-match actor is the only shape with sane fanout-vs-receive ratio.
Centrifugo (Go) on bare metal for ws-edgeCloudflare Durable Objects + WS HibernationDO is cheaper on spiky cost but operationally is a different world (we're a tenant, not an operator). Cost model differs 5×, failure modes differ entirely. We pick one as primary; DO is a strategic future-look, not a daily decision.
Custom partitioner direct-routing matchId → partitionhash(matchId) % 25640 matches × 256 slots has 95% birthday-paradox collision under hash — would destroy per-match ordering isolation. Custom partitioner is a 50-line Go function.
2-tier relay tree for hot matches, activated by gameday-forecastSingle shard at higher CPU; or reactive activation on shard saturationReactive activation fires AFTER users already see the slowness. Proactive needs a forecast service (named as an upstream dependency).
Protobuf + delta encodingJSON full state2.8× egress cost reduction at our fanout factor; egress is 90% of variable cost.
Coalescing window keyed by (matchId, eventId, revision)Naive same-event coalescingNaive coalescing silently drops VAR corrections — SEV-1 risk. The keying makes corrections never collapse.
Match-shard commits consumer offset AFTER publish-ack + idempotent downstreamCommit-before-publish OR commit-after-without-dedupeBoth naive answers are wrong differently (silent skip vs duplicates).
Bearer JWT + signed JWKS pin + ws-edge fixed-retries=1 on cache missAuth callback every WS open, or exp-backoff cache-missAuth callback = self-DOS at kickoff. Exp-backoff cache-miss = the AWS Builders' Library 'avoid fallback' anti-pattern.
etcd Raft 5-node spanning 3 regions; lease fail-CLOSEDSingle-region etcd; or fail-OPENCross-region etcd survives a region loss; fail-CLOSED costs 30s of paused scores during full etcd outage but avoids the dual-active footgun.
Active-active multi-region per matchSingle-region with multi-region read-replicasAt >5M global concurrent, cross-continent latency breaks SLO; active-active is the only honest answer.
Dual provider feeds (Sportradar + Genius)Single providerSingle-provider outage takes the whole product down.
Kafka with tiered storage + custom partitionerNATS JetStream / PulsarWe keep the KIP-429 cooperative-sticky rebalance machinery and the rich tooling.
Redis Streams for replay, match-duration-aware TTLEmbedded ring buffer per shard, or Kafka subsetPer-match TTL trivial; embedded loses replay on shard restart.
snapshot-api as separate node from ws-edgeBundled in ws-edgeSnapshot traffic spikes 20× when replay cold-restarts; coupling capacity is wrong.
Tracing on separate node (Tempo)Bundled into a data-plane serviceA data-plane outage would take tracing down with it during an incident — wrong.
WS score-delivery spine only(deferred) background push + delivery-telemetry OLAPBoth are additive consumers of the same score-event spine — a push-dispatcher fanning out to APNs/FCM and a ClickHouse telemetry sink layer back on as new consumers without touching the core fanout path. Deliberately scoped out here to keep the focus on the socket-fanout problem.

Primary sources

  • Hotstar — 25.3M concurrent (Disney+Hotstar engineering)
  • Discord — Scaling Elixir to 5M concurrent (Manifold)
  • Cloudflare Durable Objects WebSocket Hibernation
  • Slack Flannel — application-level edge cache
  • SRE Workbook Ch.22 — Addressing cascading failures

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Live Sports Scores to Millions yourself