View Count on a Video/Post
Worked solution

View Count on a Video/Post — a worked solution

Dedup, bot-filter, batched aggregation.

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 View Count on a Video/Post workspace

The problem

Count the views on a video or post. It sounds like count++; it is one of the harder write-heavy distributed-systems problems in practice. The work is dominated by three pressures the naive version ignores: dedup (the same viewer reloading must not inflate the number), bot / invalid-traffic filtering (a large fraction of raw "views" are not humans, and the public number must be reported net of them), and batched aggregation (you cannot synchronously increment a row per view at 250k views/sec — you log an idempotent event and roll it up asynchronously). The system is also read-dominated for the display of the count and hot-keyed on the write of a viral video, so it pulls in opposite directions at once.

This canonical models the main flows: the ingest/write path, the read path that displays the count, the speed-layer aggregation that produces the fast approximate number, and the batch "verified-views" layer that produces the authoritative one — the layer behind YouTube's famous 301+ freeze and the reason counts sometimes go down.

The reference architecture

Reference architecture for View Count on a Video/Post: 16 components — Viewer & Creator Apps, Edge CDN, Beacon & API Gateway, View Ingest Service, Dedup & Uniques Store, View Event Log, Stream Aggregator, Counter Store, Count Cache, Count Read Service, Raw Event Lake, Verified-Views Auditor, IVT Signal Feed, Analytics OLAP, Dead-Letter Queue, Synthetic View Canary — connected by 22 flows.Viewer & Creator AppsclientEdge CDNCloudflare / FastlyBeacon & API GatewayEnvoy + WAFView Ingest ServiceGo (stateless)Dedup & Uniques StoreRedis (cluster, noevict…View Event LogKafkaStream AggregatorFlink (RocksDB state)Counter StoreCassandra / ScyllaDB (r…Count CacheRedis / EVCacheCount Read ServiceGo (single-flight)Raw Event LakeS3 (Parquet, date-parti…Verified-Views AuditorSpark (scheduled batch)IVT Signal Feed3rd-party invalid-traff…Analytics OLAPClickHouse / Druid / Pi…Dead-Letter QueueKafka DLQ topicSynthetic View Canaryout-of-band prober
16 components, 22 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Viewer & Creator Apps

Two roles. As a VIEWER it fires a fire-and-forget view beacon once a play crosses the 'counts as a view' threshold (a few seconds of real watch time) and, separately, reads the public display count for a video. As a CREATOR it reads time-bucketed analytics (views over time, by geo / device / traffic source) from the Studio dashboard.

Why it exists. It is the traffic source and the only place the two numbers that matter are observed: the cheap public count every viewer sees, and the verified analytics a creator is paid against. We model it because the ~8:1 count-read:view ratio, the 'a view is watch-time, not a page-load' rule, and the expectation that a public count never visibly flickers downward are all stated from the client's point of view.

When it fails. A buggy or malicious client firing beacons in a loop is the front line of view fraud and the easiest way to inflate one video's count. The event_id dedup, the watch-time gate, and the gateway's per-device rate limit bound a single client's contribution; whatever slips past inline is clawed back by the batch verifier.

Edge CDNCloudflare / Fastly

Caches the tiny public display-count JSON at the edge with a short TTL and serves it to viewers worldwide without touching the origin. View beacons (POST) and authenticated Studio reads pass straight through to the gateway.

Why it exists. Reads outnumber view events ~8:1 and a viral video's count is read millions of times a second — the count is the single most cacheable object in the system because it is tiny, public, and allowed to be a few seconds stale. The edge collapses >99% of read QPS so the origin only ever sees cache-fill traffic. Folding this into the read service would dump 1.5–2.5M req/s on the origin for a number that changes once per rollup window.

When it fails. An edge outage shifts the full 1.5–2.5M read QPS onto the origin read tier and count cache — survivable only because the read tier is sized for it and single-flights cache fills. A mis-set long TTL freezes a count visibly stale for minutes (the bug users notice on a breaking-news video). Detection: edge hit-ratio and origin-offload ratio.

Beacon & API GatewayEnvoy + WAF

L7 ingress for both paths. Terminates TLS, runs the WAF, applies the cheap first layer of bot filtering (known-bot user-agents, datacenter / VPN IP ranges, missing or forged signals — 'General Invalid Traffic'), enforces per-device / per-IP token-bucket rate limits, and routes by intent: view beacons to the ingest service, count reads (on edge miss) and authed Studio reads to the read service.

Why it exists. Centralizes the cheap, high-volume defenses — TLS, rate-limit, GIVT drop — at one tier so a beacon flood is shed before it can reach the durable log and the aggregation pipeline. GIVT is filtered HERE, pre-ingest, because it is decidable from the request alone and there is no point paying to log and aggregate traffic you already know is a crawler (the GIVT/SIVT split ad systems are accredited on).

When it fails. Gateway saturation reads to users as 'views aren't updating' (beacons shed) and 'the count won't load' at once. A too-loose rate limit lets a fraud campaign inflate counts faster than the verifier claws them back; a too-tight one drops legitimate views off a genuinely viral video (false-positive undercount). Detection: beacon accept/reject rate, rate-limit-reject rate, downstream 5xx.

View Ingest ServiceGo (stateless)

Stateless collector for view beacons. For each beacon it (1) checks the dedup store to see whether this (viewer, video) was already counted inside the suppression window and whether this exact event_id was already seen, (2) records the viewer in the per-video unique-count sketch, and (3) if the view is new, emits an idempotent view event to the durable log. It acks the client immediately — the beacon is fire-and-forget.

Why it exists. The hot path must do the minimum that cannot be deferred — dedup and idempotency — and nothing that can wait. It deliberately does NOT increment any counter (the golden rule from every production write-up: never put a non-idempotent increment against a strongly-consistent store on the hot path); it appends an idempotent event and lets the async layer aggregate. Splitting ingest from the read tier lets it scale independently to the 250k/s peak.

When it fails. If ingest can't reach the dedup store it fails OPEN — but not blindly. Every event still carries its event_id, and the speed-layer aggregator dedups by event_id WITHIN its window state, so a retried/duplicate beacon during a dedup outage is still caught in near-real-time (inside the ~10s window) and the daily verifier catches the rare cross-window duplicate. What actually leaks is only the per-(viewer,video) SUPPRESSION (a genuine re-view inside the window may count) — bounded and recoverable — chosen over fail-CLOSED because dropping views is a silent, unrecoverable undercount on a creator-paid number. We do NOT count blind: failing open without that downstream event_id dedup would manufacture exactly the inflation-then-purge that later forces the public count to drop. If ingest can't reach the log, the beacon is lost (best-effort) — an undercount, never a 5xx to the viewer. Detection: dedup-store error rate, publish error rate, ingest→log lag.

Dedup & Uniques StoreRedis (cluster, noeviction)

Two bounded-memory jobs in one Redis tier. (1) A short-TTL key per (viewer, video) and per event_id so repeat views inside the suppression window and retried beacons are recognized and dropped. (2) A per-video HyperLogLog sketch that estimates the number of UNIQUE viewers without storing the viewer set.

Why it exists. Exact dedup would require remembering every (viewer, video) pair forever — unbounded memory and a privacy liability. A TTL'd key set bounds the dedup window to what actually matters (a day or two), and HyperLogLog bounds unique-counting to a fixed ~12 KB per video at ~0.81% error instead of the 400–600 MB an exact set of 10M viewer-ids would cost (Reddit's exact design: Kafka → Redis HLL per post → Cassandra totals). It is a SEPARATE tier from the count cache so a flood of cache traffic can't evict dedup state and silently start inflating counts.

When it fails. Eviction (or OOM-with-LRU) silently drops dedup keys → already-counted views count again → fleet-wide INFLATION during the busiest window. That is why the policy is noeviction + page-on-ANY-eviction, not 'size it and hope'. On leader loss the async follower loses only the un-replicated tail (RPO ≈ replication lag, ~5s); losing BOTH replicas of a shard loses that shard's whole in-memory dedup window (no persistence) — bounded re-count for those keys, caught and corrected by the verifier. Detection: evicted_keys > 0 (page), used/max memory > 80%, dedup-miss-rate trend.

View Event LogKafka

The durable commit log of view events. Every accepted view is appended here as an immutable event; the speed-layer aggregator and the raw-event sink both consume it. Partitioned by video_id so all events for a video are ordered on one partition and handled by one consumer.

Why it exists. It is the buffer that absorbs viral spikes and decouples the fast ingest path from the slower, consistent aggregation. Because the log is durable and replayable, the counter is a PROJECTION that can be rebuilt from it — this is the event-sourced answer to 'counts went backwards / got corrupted': replay from an offset. A fire-and-forget multicast you can't replay couldn't do that.

When it fails. Consumer lag here is the user-visible 'counts are hours behind' failure — and it is usually ONE partition (a viral video) lagging while the rest are fine, so the alert is on per-partition lag VARIANCE, not aggregate. A poison/schema-incompatible event can crash-loop a consumer on one offset and stall its whole partition — mitigated by the DLQ + skip. Partition leader election (~10s, RF=3) queues but doesn't lose events.

Stream AggregatorFlink (RocksDB state)

Consumes view events, keys by video_id, and maintains a short tumbling window (~10s) of counts and merged unique-sketches in RocksDB state. At each window tick it emits ONE batched rollup per (video, region) into the counter store, a fresh approximate count into the count cache, and a time-bucketed rollup into the analytics OLAP store. Bad events are routed to the DLQ.

Why it exists. This is the hot-key absorber. Instead of N increments racing on one counter, N events in a window collapse to ONE write — a video taking 100k views/s becomes a single '+~1,000,000 this window' write, the only way a single counter survives a viral spike (Twitter TSAR, Cloudflare's merge-time rollups, and Netflix's interval rollup all do this). It also does the windowed HLL merge for uniques. The window length is the accuracy/latency dial: shorter = fresher + more writes; longer = cheaper + staler.

When it fails. A checkpoint failure or restart from a stale offset replays an OLDER state, so the emitted count can DECREASE — the public 'view count went down' bug. Mitigated by the read tier's monotonic guard (never serve below the last-served speed value) and by treating the batch-verified count as the only legitimate source of a decrease. A backed-up window (slow downstream) grows checkpoint size and stalls the pipeline. Detection: checkpoint duration/failures, watermark lag, per-key window backlog, and a monotonicity alarm on the emitted count.

Counter StoreCassandra / ScyllaDB (rollup cells)

The durable serving store for the displayed count. Holds the rolled-up count as idempotent cells keyed by (video_id, region, shard) — the total for a video is the SUM of its cells. The speed layer writes recent windows; the batch verifier overwrites with authoritative verified totals.

Why it exists. Counts must survive a cache flush and a region loss, and they must be writable from every region at once (a view in Tokyo and a view in Frankfurt both happen 'now'). Cells keyed by (video, region) make increments COMMUTATIVE — each region owns its slot and the global count is their sum — so the store can be active-active across regions with no conflict to reconcile. This is why view counting is multi-region-friendly where a single mutable balance would not be.

When it fails. A hot video's cells can still concentrate on one shard if salting is mis-tuned → that shard browns out while the cluster idles (the dominant failure mode here). Read amplification grows with shard count (N cell-reads per displayed count) — the cost of surviving the hot key. Losing a shard makes 1/N of videos' counts unavailable until restore; reads degrade (the count cache serves last-good) rather than fail. RPO is 0 for in-region node loss (quorum) but ~the cross-region merge lag (a few seconds) for a full REGION loss — a dead region's most-recent un-merged increments are simply absent from the global SUM until it recovers and replays from its log. Detection: per-shard CPU/latency skew, hot-partition sampling.

Count CacheRedis / EVCache

Caches the current displayed count per video so the read service answers an origin (CDN-miss) request in ~1ms without summing counter shards. The aggregator writes the fresh count here on each window tick; the read service back-fills it on a miss.

Why it exists. Even after the CDN, the origin still sees cold-key and TTL-expiry reads, and SUMming N counter cells per read (the read amplification of the sharded counter) is too slow to do on every origin read. The cache turns the common origin read into a single GET and is where the read tier's single-flight collapses a herd of concurrent misses for a freshly-viral video into one shard-SUM.

When it fails. A cold cache (deploy / restart) shifts every origin read to a counter-shard SUM — survivable because single-flight collapses the herd, but a latency cliff. A briefly-stale generation served is acceptable (counts are eventual). Detection: hit-ratio cliff, counter-store read-QPS spike. RTO ~15s/host, RPO 0 (no durable data).

Count Read ServiceGo (single-flight)

Serves origin (CDN-miss) count reads and authenticated creator-analytics reads. For a count it reads the count cache; on a miss it single-flights a SUM of the video's counter cells and back-fills the cache. For analytics it queries the OLAP store for time-bucketed / dimensional rollups. It applies the monotonic guard so a public count never visibly decreases due to speed-layer wobble.

Why it exists. A dedicated read tier is where single-flight / request coalescing lives — all concurrent misses for one hot video collapse to one counter-store SUM (the textbook stampede fix). It also owns the two product invariants the stores below it don't: monotonic serving (counts only go up, except a deliberate verified correction) and the read-amplification SUM logic for sharded counters. Keeping it stateless lets it absorb a CDN outage's worth of read QPS.

When it fails. If the monotonic guard is wrong, users see the count flicker down (a reputational bug). If single-flight breaks (e.g. a deploy resets the in-flight map) a hot-key miss can stampede the counter store. Detection: count-decrease alarm, counter-store QPS vs cache-miss correlation. Stateless → failover is just capacity.

Raw Event LakeS3 (Parquet, date-partitioned)

Durably archives every raw view event (via a Kafka→S3 sink) as date-partitioned Parquet. It is the long-horizon source the batch verifier scans to re-score fraud and recompute authoritative counts, and the backfill source for analytics.

Why it exists. The stream's 72h retention is enough for the speed layer to catch up but not for a DAILY fraud audit or a multi-week backfill — those need cheap, durable, replayable raw events. The lake is the batch layer's input in the lambda split: the speed layer approximates from the stream, the batch layer corrects from the lake. It's object storage because it's write-once, read-rarely, huge (~500 GB/day) and must be cheap.

When it fails. Lake unavailability doesn't touch the live count — it only stalls the batch verifier, so fraud corrections and backfills are delayed (the verified count lags further behind the raw count). A sink-connector lag means recent events aren't archived yet, narrowing how far back the verifier can audit. Detection: sink lag, object-write error rate. Durability ~11 nines; effectively no RPO.

Verified-Views AuditorSpark (scheduled batch)

The batch / correctness layer. On a schedule (the daily audit) it scans raw events from the lake, joins them against an external invalid-traffic feed, re-scores sophisticated fraud (SIVT) that inline GIVT filtering can't catch, computes the AUTHORITATIVE verified count, and overwrites the counter store (and analytics) with it — which can move the public number DOWN when fake views are purged.

Why it exists. Inline filtering at the gateway catches only the cheap, obvious bots (GIVT); sophisticated fraud needs cross-event ML and third-party intelligence far too expensive to run on the 250k/s hot path. So the system runs two numbers on two latencies: a fast approximate count (speed layer) and a slow verified count (this). This is literally the YouTube model — the '301+' freeze and the later up-or-DOWN correction are this batch audit reconciling the raw count (Twitter TSAR's batch layer reconciling its speed layer is the same idea).

When it fails. A verifier bug can purge legitimate views (false-positive undercount → creator-revenue disputes) or miss a fraud campaign (false-negative inflation → advertiser fraud). Because its correction is the only sanctioned count decrease, a bad run is directly user-visible as 'my views dropped'. Mitigated by shadow-scoring before enforcing and a manual-review queue for high-velocity assets. A stalled verifier just delays corrections (raw and verified diverge further). Detection: audit lag, correction-magnitude distribution, sampled FP/FN review.

IVT Signal Feed3rd-party invalid-traffic intelligence

A third-party intelligence feed (IP reputation, datacenter/VPN ranges, known-bot fingerprints, device-attestation verdicts) that the batch verifier joins against to classify sophisticated invalid traffic it can't decide from the raw events alone.

Why it exists. Sophisticated fraud detection benefits from cross-customer intelligence no single platform sees in isolation, so teams buy an IVT feed rather than build every signal in-house (the MRC-accredited ad-fraud vendors exist for exactly this). It is modeled as external because its availability, latency, and SLA are outside our control — which dictates how we call it.

When it fails. Feed outage or rate-limiting stalls or degrades the fraud audit — the verified count then lags or is computed on stale signals, but the live (raw) count is untouched. A feed that returns bad verdicts causes mass false-positive purges. Detection: circuit-breaker state, feed error/latency, verdict-distribution drift. The classic third-party-dependency risk — bounded by the circuit breaker and by being off the hot path.

Analytics OLAPClickHouse / Druid / Pinot

Columnar OLAP store holding view rollups by time bucket × dimensions (geo, device, traffic source, watch-time band) plus approximate unique-viewer sketches. Powers creator Studio dashboards and internal reporting with sub-second aggregate queries over hundreds of billions of rows.

Why it exists. The counter store answers one question fast — 'what's the total for this video' — but a creator dashboard asks 'views per hour for the last 28 days, by country, by device', a scan / group-by no point-lookup store can serve. That's a columnar OLAP workload. Separating it from the counter store keeps the hot count path simple and lets the analytics store optimize for wide scans (LinkedIn serves 'Who Viewed Your Profile' from Pinot for exactly this reason).

When it fails. Compaction/ingestion lag stalls fresh segments → dashboards show stale data; replica lag shows different numbers to different viewers of the same dashboard. None of this touches the public count (separate store). Detection: segment-handoff / compaction lag, replica seconds-behind, a cross-replica divergence canary. Mitigation: route reads to caught-up replicas; bound compaction backlog.

Dead-Letter QueueKafka DLQ topic

Holds view events the aggregator can't parse or process — malformed payloads, schema-incompatible records, repeated processing failures — so the main partition keeps flowing and a human or a repair job can triage and replay them.

Why it exists. Head-of-line protection. Without it, one poison event crash-loops the consumer on a single offset and stalls that whole partition — meaning one bad message freezes one (often viral) video's counting. The DLQ turns 'a bad event stalls a partition' into 'a bad event is parked and alerted'.

When it fails. A growing DLQ means events are silently NOT being counted — an accumulating undercount. A DLQ that fills faster than it drains is an incident, not a backlog. Detection: DLQ depth + age-of-oldest-message.

Synthetic View Canaryout-of-band prober

Continuously generates K KNOWN probe views through the REAL front door (gateway → ingest → log), then reads the affected counter back out-of-band and asserts the count rose by exactly K (net of the dedup it deliberately exercises) within the freshness SLO, and that the verified-vs-raw divergence stays within bound. It pages on violation regardless of what any internal metric says.

Why it exists. Every OTHER 'is the count wrong' signal is computed FROM the pipeline that would be wrong: the monotonicity alarm watches the aggregator's own emitted value, and the inflation canary divides two outputs of the same dedup/HLL pipeline — so a smoothly-inflating or confidently-wrong pipeline keeps them green. The system's worst failure is silent (counts wrong, every endpoint returns 200). A watchdog cannot share a failure domain with the thing it watches, so this prober exercises the end-to-end path itself and reads ground truth directly — the only signal with independent ground truth.

When it fails. If the canary is down you lose the only independent signal, so it runs redundantly across zones and its liveness is externally alarmed. It is deliberately the dumbest, most independent component in the system precisely so it almost never fails for the same reason the pipeline does.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Typical clarifications to surface:

  • What is a "view"? A page load, or watch-time past a threshold? Drives the client gate and the dedup window. We assume watch-time-gated.
  • Unique viewers or total views? Both are usually wanted; uniques are approximate (HyperLogLog). Drives the dedup store design.
  • How exact must the public count be? Approximate + eventually consistent (seconds of lag) is universally accepted for a display count; advertiser/creator-paid counts need a verified, auditable number on a slower cadence.
  • Must the count be monotonic? Users hate counts going down — but fraud purges legitimately lower it. We serve monotonic for speed-layer wobble and allow the verified layer to correct.
  • Retention? Raw events: weeks (then drop / cold). Aggregates: forever. Legal/PII constraints on raw viewer data.
  • Multi-region? Counts are commutative, so active-active is natural — unlike a mutable balance.

Assumptions to state:

  • 5B view events/day (~58k/s avg, ~174k/s diurnal peak, +100k/s on a single viral video).
  • Count reads ≈ 8× view events → 1.5–2.5M reads/s, almost all CDN-absorbable.
  • A view = watch-time past a threshold, deduped per (viewer, video) over a 24–48h window.

02Functional reqs

What must this system actually do?

  • Record a view: accept a fire-and-forget beacon, dedup it, and durably log an idempotent event.
  • Display the total count: serve a tiny, cacheable, approximate count fast and globally.
  • Count unique viewers (approximate, via HyperLogLog).
  • Filter invalid traffic: cheap inline GIVT at ingest; sophisticated SIVT in a batch audit that produces a verified count.
  • Creator analytics: views over time by geo / device / source for a dashboard.
  • Reconcile: a batch layer recomputes the authoritative count from the durable event log and corrects the serving store.

03Non-functional

What must it promise about speed, uptime and correctness?

SLOs:

DimensionTarget
Ingest a view (beacon enqueue + ack)p50 5 ms, p99 30 ms, p99.9 80 ms
Read display count (CDN/cache hit)p50 2 ms, p99 15 ms, p99.9 50 ms
Read display count (origin, cache miss → SUM cells)p99 ~60 ms
Count freshness (view → reflected in display)p50 ~10 s, p99 ~60 s (eventual, by design)
Verified-count cadencedaily audit (raw vs verified diverge until then)
Read availability (the count)99.99%
Ingest completeness (beacon-accept → log-append)gap < 0.1% — drops are a silent undercount on a paid number, so it's budgeted + alarmed
Event-log durabilityacks=all, RF=3 — we refuse to lose an accepted view

Error budget. Read availability 99.99% ≈ 52 min/yr. But the failures that actually govern on-call here are not unavailability — they are wrong numbers: count inflation (dedup failed), count staleness (consumer lag), count going backwards (checkpoint replay), and silent undercount (dropped beacons). The headline burn-rate alerts are the count-staleness SLO (display lag past 60s) and a monotonicity alarm (any served public count that decreases for a non-verified reason pages immediately). But those two — and the inflation canary (ingested ÷ unique-estimate) — are all computed FROM the pipeline that would be wrong, so the load-bearing signal is the independent synthetic canary: a prober that writes K known views through the real front door and reads the counter back out-of-band, paging if the count didn't move by K within SLO. It shares no failure domain with the pipeline, so it catches the silent gray failure (counts wrong while every endpoint returns 200) that the derived SLIs can miss. Multi-window multi-burn-rate per the SRE Workbook.

04Capacity estimation

How much load and data does this have to hold?

Anchor on a YouTube/TikTok-tier platform and derive every shard/partition/cell number rather than guessing.

1. Throughput. YouTube serves ~1B watch-hours/day ≈ 5B views/day; TikTok ~1B/day. Take V = 5e9 views/day. ingest avg = 5e9 / 86,400 ≈ 57,900 view-events/s. Diurnal peak ×3 ≈ 174k/s; a single viral video adds +100k/s on top → design for ~250k view-events/s ingest. Count reads ≈ 8× ingest → ~1.4M/s at diurnal peak, ~1.5–2.5M/s once viral-read headroom is folded in (a trending video is re-read far more than 8× its view rate), ~99% absorbed by the edge (counts are tiny, public, and stale-tolerant).

2. Storage asymmetry — the whole point. raw event ≈ 100 B (video_id, viewer/session, ts, geo/device, event_id) → 5e9 × 100 B = 500 GB/day (~3.5 TB/week, ~180 TB/yr if kept). vs. aggregated counter ≈ 24 B × 1e9 videos ≈ 24 GB total, kept forever. Raw events are huge and ephemeral (30-day hot in the lake, then cold/expire); the aggregates are tiny and permanent. This asymmetry is why raw events go to cheap object storage and counters live in a small, fast store.

3. The hot key — the central scaling problem. Views are power-law: one viral video can take 100k views/s — i.e. one counter key taking 100k writes/s.

  • A single Redis node does ~100k–180k INCR/s; a single DB row serializes on its lock at ~1,000 writes/s before it falls over.
  • Fix = sharded counter: split one logical counter into N cells, write key:{rand(N)}, read = SUM(cells). For a DB-row-class store, counterShardsHotKey = 100,000 / 1,000 = 100 cells; read amplification = N reads per displayed count (the trade).
  • And pre-aggregation: the Flink window collapses 100k events/s into one batched UPSERT per (video, region) per window, so the counter sees ~1 write/window, not 100k/s.

4. Log partitions. streamPartitionsMin = 174k/s ÷ ~1k events/s/partition ≈ 174, rounded up to 256 (next power of two, for headroom + clean rebalancing). A viral video still maps to one partition (hot), mitigated by key-salting the hottest keys.

5. Dedup / uniques state. Exact unique-set ≈ 40–60 MB for a 1M-viewer video in a Redis set (per-member overhead dominates; the raw 16-byte ids alone are 16 MB) — and it grows without bound. HyperLogLog = 12 KB/video at ~0.81% error, regardless of cardinality. Active working set ≈ 10M videos × 12 KB ≈ 120 GB sharded Redis; sketches age out with a 24–48h window.

6. Where the architecture flips.

ChoiceWorks toThen switch to
Synchronous count-on-write~10k writes/sasync: log → batch-aggregate → counter
Single DB-row counter~1k writes/s (row-lock wall)sharded counter cells, read = SUM
Single Redis INCR~100k writes/ssharded counters / probabilistic counting
Exact unique settens of MB/hot videoHyperLogLog (12 KB, 0.81% error)
Counter-column INCRnever (non-idempotent trap)idempotent windowed UPSERTs
Single-region counterone region(video, region) cells, commutative, active-active

05API design

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

Record a view (fire-and-forget beacon):

POST /v1/beacon
Content-Type: application/json
{ "video_id": "vid_123", "event_id": "9f1c-...", "event_time": 1717795200, "session": "s_88", "ms_watched": 7400 }

204 No Content        # acked immediately; dedup + log happen before ack, aggregation is async

event_id + event_time are the idempotency key: a retried or double-fired beacon is deduped, not counted (Netflix's counter-event design). ms_watched is re-checked server-side against the watch-time threshold.

Read the display count (cacheable):

GET /v1/videos/vid_123/count
200 OK
Cache-Control: public, max-age=30, stale-while-revalidate=30
{ "video_id": "vid_123", "views": 1048217, "as_of": "2026-06-08T12:00:30Z", "verified": false }

verified:false is the fast speed-layer number; the nightly audit replaces it with verified:true (which may be lower). The public count is views only — the unique-viewer estimate is not on this path. The HLL sketch lives in the dedup Redis (off the read path) and as a Theta sketch in the analytics OLAP store; uniques are served only by the authed Studio endpoint below, which reads that OLAP store. Putting unique_estimate here would require the aggregator to also project a merged per-video sketch into the count cache / counter store — see Open questions.

Creator analytics (authed, not cacheable):

GET /v1/studio/vid_123/views?from=...&to=...&by=country,device   (Bearer)
200 OK   { "series": [ ... per-bucket rollups ... ] }
403 Forbidden   # bearer principal does not own / is not authorized for vid_123

The bearer token is authentication; the read service additionally runs a per-resource authorization check — the principal must own (or be granted access to) video_id — before serving, and returns 403 otherwise. Without it any authenticated creator could read any other creator's private, competitively-sensitive analytics. The principal→video ownership lookup is part of the analytics read flow (see # hld Security).

06Data model

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

View event (immutable, in the log + lake): video_id, event_id (idempotency), event_time, viewer/session, region, coarse device/geo/source, ms_watched. The event_id makes the whole pipeline replayable and dedupable.

Dedup state (Redis, TTL'd): seen:{viewer}:{video} and seen:{event_id} keys (24–48h TTL); hll:{video} HyperLogLog sketch for uniques. Bounded, reconstructable, noeviction.

Counter cells (Cassandra/Scylla): (video_id, region, shard) → rolled_up_count written as idempotent windowed UPSERTs; total = SUM. A separate verified_count column the batch layer overwrites.

Analytics rollups (ClickHouse/Druid/Pinot): (date, video_id, country, device, ...) → views, theta_sketch(uniques) materialized at merge time.

Database choice — recommended. Kafka for the event log, Redis for dedup + HLL + the count cache, Cassandra/Scylla (idempotent rollup cells, not counter columns) for the durable serving counter, ClickHouse/Druid/Pinot for analytics, S3 for the raw-event lake. The non-negotiable: the hot path appends an idempotent event to a durable log and the count is a projection aggregated asynchronously — never a synchronous non-idempotent increment against a strongly-consistent store.

07High-level design

Which components handle a request, and in what order?

Ingest / write path. Client → (CDN pass-through) → Gateway → View Ingest Service → Dedup Store, then Ingest → View Event Log. The gateway sheds GIVT and rate-limits; ingest dedups (SETNX + event_id), updates the HLL, and publishes an idempotent event to Kafka. The client gets a 204 immediately — nothing on this path increments a counter.

Read path (the product, ~99% edge-served). Client → CDN. On an edge miss: → Gateway → Count Read Service → Count Cache; on a cache miss the read service single-flights a SUM of the video's counter cells and back-fills the cache. The read service applies the monotonic guard.

Speed layer (the hot-key absorber). View Event Log → Stream Aggregator (Flink). Per video, a tumbling window collapses N events into ONE batched, idempotent UPSERT into the Counter Store, a versioned SET into the Count Cache, and a rollup into the Analytics OLAP store. This is what lets a 100k-views/s video survive on one logical counter.

Batch / verified layer (the YouTube 301). View Event Log → Raw Event Lake (S3 sink). On a schedule, the Verified-Views Auditor scans the lake, enriches against the external IVT Signal Feed (circuit-broken, off the hot path), re-scores sophisticated fraud, and writes the authoritative verified count back to the counter store and analytics — the one sanctioned source of a count decrease. A DLQ parks poison events so one bad message can't stall a partition.

Independent watchdog. The Synthetic View Canary sits outside the pipeline entirely: it writes K known probe views through the real front door and reads the counter back directly (bypassing the cache and CDN), paging if the count didn't move by K within the freshness SLO. Because every other "is the count wrong" signal is derived from the pipeline that would be failing, this out-of-band prober is the only one with independent ground truth — the dead-man's-switch for the silent gray failure where counts are wrong but every endpoint returns 200.

Why two numbers. The speed layer gives a fast, approximate count (verified:false); the batch layer gives a slow, authoritative one (verified:true). Lambda architecture: the batch layer reconciles what the speed layer approximated (Twitter TSAR/Summingbird, Cloudflare's ClickHouse rollups, Netflix's interval rollup).

08Deep dives

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

The hot key (one viral video)

A single video at 100k views/s is 100k writes to one counter key. A DB row dies at ~1k/s; a Redis node at ~100k/s. Four mechanisms stack: (1) the log absorbs the spike and decouples writes from aggregation; (2) window pre-aggregation collapses N events → 1 UPSERT/window; (3) sharded counter cells spread even that one write across N cells (read = SUM); (4) key-salting the hottest log partitions so one video doesn't pin one partition. The edge absorbs the read side entirely.

Idempotency & the counter-column trap

The stream is at-least-once and Flink's exactly-once is internal-state only — an external counter write re-fires on checkpoint recovery. So the external write must be idempotent: a (video, region, window_id)-keyed UPSERT, not an INCR. Cassandra counter columns are the classic trap — non-idempotent, so a timed-out retry overcounts and a dropped one undercounts; "not even eventually consistent." Ably's post-mortem: a client retry loop on counter columns drove every node to 100% CPU across all regions, recurring at 24h intervals, until they removed counter columns entirely. The dedup store (event_id) is the upstream defense; idempotent UPSERTs are the downstream one. And the idempotency hinges on window_id being derived from event-time, not wall-clock, so a checkpoint replay re-derives the identical key instead of re-bucketing a late event into a new window and double-counting.

The txn.group "view-rollup" (mode event-sourced) tags the two edges that define the count's correctness — ingest→log (participant event) and aggregator→counter (participant projection) — because the log is the system of record and the counter is its replayable projection. This is event-sourced, not outbox: ingest has no local DB write to make atomic with the publish (the dedup SETNX is best-effort), so there is nothing for an outbox to atomically commit. The raw-event sink (log→lake) and the OLAP rollup (aggregator→analytics) are also projections of the same log, but are intentionally left out of the group so it stays focused on the serving-count correctness contract.

Why the count goes backwards (two causes)

Bug: a Flink checkpoint replays an older state → the emitted count decreases. Fix: the read tier's monotonic guard serves max(value, last-served) for the speed-layer path so a regression never reaches users. By design: the batch verifier purges confirmed-fake views → the verified count is legitimately lower. This is YouTube's 301+ freeze productized: the count froze pending a bot audit, then jumped — or dropped — when fakes were purged (now continuous AI validation + daily batch audits). The serving tier distinguishes them: speed-layer regressions are clamped; a verified write is the one sanctioned decrease.

How the verified decrease survives the guard AND the speed layer (the load-bearing detail). Two traps the mechanism must avoid, and both must be handled or "counts go down" silently doesn't work. (1) The monotonic guard could clamp the legitimate decrease back up — so every count value carries a verified bit; the guard exempts a verified value from the max() floor AND resets last_served down to it, so subsequent speed-layer reads clamp to the post-purge floor, not the pre-purge number. (2) The next speed-layer window could re-add the purged views — so the verifier writes a SEPARATE authoritative column the read-SUM prefers over the speed-layer cells, and advances the speed layer's event-time watermark past the audited day so the aggregator can't re-emit those events. Without (1) the verifier is a silent no-op (the count never drops); without (2) the correction is erased on the next ~10s tick. The verified:false/true bit already on the API response is exactly this stamp, threaded from the verifier through the counter column into the read tier's guard.

Dedup, uniques, and silent inflation

Dedup is a TTL'd SETNX per (viewer, video) + event_id; uniques are a HyperLogLog (12 KB/video, 0.81% error vs 400–600 MB for an exact set — Reddit's design). The killer failure is silent: if the dedup Redis uses LRU and hits maxmemory, it quietly evicts dedup keys → already-counted views count again → fleet-wide inflation exactly during peak. The fix is a one-liner: noeviction + page on ANY eviction (fail loud), and keep the dedup Redis physically separate from the count cache so cache churn can't evict dedup state.

Approximate vs exact, speed vs batch

The display total is approximate (summed cells, seconds stale) and that is fine — nobody needs the millionth view to be exact in real time. Uniques are approximate by construction (HLL/Theta). The verified count is exact-ish and slow. Choosing what may be approximate is the core latency/cost lever: it is what lets the hot path be a fire-and-forget append.

Observability

SLIs (RED/USE): Ingest — beacon accept/reject, publish error, ingest→log lag (RED). Log — produce/consume rate, per-partition consumer-lag VARIANCE (one viral partition lagging is the signal, not aggregate lag). Aggregator — checkpoint duration/failures, watermark lag, window backlog. Counter — per-shard CPU/latency skew (hot-key detector), read amplification. Dedup — evicted_keys (page on any), memory %, dedup-miss trend. The headline alarms: (1) count-staleness (display lag > 60s); (2) monotonicity (a served public count decreased for a non-verified reason — pages immediately); (3) inflation ratio (ingested-events ÷ unique-estimate above baseline). But (1)–(3) are all computed FROM the pipeline that would be failing — a confidently-wrong or smoothly-inflating pipeline keeps them green, and the inflation ratio divides two outputs of the same dedup/HLL pipeline so a shared bug can move both together — so the load-bearing alarm is (4) the independent synthetic canary: K known probe views written through the real front door, the counter read back out-of-band (bypassing cache + CDN), paging if it didn't move by K within SLO. It is the only signal with ground truth independent of the pipeline, paired with a 'canary-not-reporting' alarm so a dead canary can't hide behind a flatlined-green SLI. Dashboards on-call open: display-count staleness by region; per-partition log lag (variance); counter hot-shard skew; dedup memory/eviction; verified-vs-raw divergence + verifier audit lag; DLQ depth/age; canary probe-delta + liveness (the independent signal). Tracing: the ingest path is head-sampled (~1%); the verified-recount path and the canary's write-then-read are 100% sampled so a wrong correction or a missed probe always has a trace.

Deploy & runtime

Aggregator deploys use cooperative-sticky + static consumer-group membership so a rollout doesn't trigger a rebalance storm (each rebalance pauses consumption and feeds lag); externalized/verified checkpoints so a deploy can't replay a stale state into a regression. Counter store: rolling, in-shard followers so a node restart doesn't drop a cell. Count cache: rolling restart + single-flight so a cold cache doesn't stampede the counter store. Schema: the view-event schema is registry-managed with BACKWARD compatibility, rolled out consumer-first; incompatible events hit the DLQ, not a crash-loop. Secrets: the IVT feed's API key and inter-service mTLS certs auto-rotate with expiry-burn alarms.

Multi-region / DR

Posture: active-active. Counts are commutative — each region ingests views locally, appends to its regional log, aggregates into its own (video, region) counter cells, and serves reads locally; the global count is the SUM of per-region cells, merged async. There is no single-writer bottleneck and nothing to reconcile on a partition (unlike a mutable balance) — the deepest reason view counting is multi-region-friendly. DR: losing a region loses only that region's in-flight, un-replicated increments (RPO ≈ the cross-region merge lag, seconds) and its share of ingest until traffic reroutes; other regions' cells are unaffected and the global SUM simply omits the dead region's most recent delta until it recovers and replays from its durable log. Be honest about the two RTOs: the counter store's leaderless quorum gives RTO ≈ 0 for node loss, but losing a whole region's ingest tier has an RTO of minutes — the traffic-reroute window bounded by DNS TTLs, resolver caching, and sticky clients, during which that region's views are simply dropped (best-effort → a bounded undercount, recoverable at the next verified audit from the durable log) and its un-merged cells are absent from the global SUM until it recovers. Claiming a single sub-second number here is how a postmortem ends up explaining why the count "lagged for 8 minutes." The caveat on scope: this canvas draws ONE region's ingest+counter tier and bakes the active-active posture into the counter store's internals (leaderless, (video, region, shard) partitioning, async cross-region merge) rather than duplicating every node per region — see Open questions.

Security

AuthN/Z: beacons and public count reads are anonymous (auth:none, edge-served); Studio reads require a bearer token AND a per-resource authorization check — the read service resolves the bearer principal's identity, looks up ownership/ACL for the requested video_id, and serves only if the principal owns or is granted that video, else 403. Authentication alone is not enough: without the ownership check any authenticated creator could read any other creator's private analytics. Every internal hop is mTLS; the IVT feed uses an API key. PII: raw events carry viewer/geo signals, so the lake is encrypted (SSE-KMS), access-controlled, and lifecycle-expired — aggregates (the only thing kept forever) carry no per-viewer data. Threat headlines: (1) view fraud / inflation — bot farms inflating a count for ranking or ad revenue; mitigated by GIVT at the edge, watch-time gating, per-device rate limits, and the SIVT batch audit; (2) beacon spoofing — forged beacons; mitigated by rate limits + the dedup/idempotency key + batch re-scoring; (3) IVT-feed compromise / bad verdicts — mitigated by shadow-scoring before enforcing and a manual-review queue.

09Trade-offs

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

Trade-offs we accepted (and the alternative we rejected)
  • Fire-and-forget event log + async rollup over synchronous count-on-write. Rejected sync increment because it dies at ~10k/s and makes the hot path a non-idempotent write to a consistent store. Cost: the count is eventually consistent (seconds of lag).
  • Idempotent windowed UPSERTs over Cassandra counter columns. Rejected counter columns as a documented production trap (Ably). Cost: the aggregator must carry window_id keys and dedup logic.
  • Sharded counter cells (read = SUM) over a single counter row. Rejected the single row because a viral key serializes on one lock. Cost: read amplification (N cell reads per display).
  • Two numbers (fast approximate + slow verified) over one exact real-time count. Rejected exact-real-time because fraud detection can't run on the hot path. Cost: raw and verified diverge until the audit, and the public count can legitimately drop.
  • HyperLogLog uniques over exact sets. Rejected exact sets (400–600 MB/hot video). Cost: ~0.81% error on uniques.
  • noeviction dedup Redis over LRU. Rejected LRU because silent eviction = silent inflation. Cost: must size for peak + page on eviction (fail loud, not graceful).
  • Active-active commutative counters over single-master. Rejected single-master because counts don't need it and it adds cross-region write latency. Cost: the global count is a sum that's briefly behind by the merge lag.
  • Fail-open ingest with downstream event_id dedup over fail-closed. Rejected fail-closed because dropping views is a silent, unrecoverable undercount on a paid number; rejected naive fail-open because counting blind manufactures the inflation-then-purge drop. Cost: the per-(viewer,video) suppression leaks during a dedup outage (a re-view may count), bounded and corrected by the audit.
  • An independent out-of-band canary over relying on pipeline-derived SLIs. Rejected derived-only monitoring because every "count is wrong" signal computed from the pipeline goes green when the pipeline is confidently wrong. Cost: a separate prober plus its own liveness watchdog to run.
  • R=1 quorum reads on the counter over R=2. Rejected R=2 because a SUM over many cells is approximate-by-construction (not an atomic snapshot) and the count is eventual anyway, so per-cell read quorum buys nothing. Cost: a single lagging replica can briefly skew one cell's contribution — invisible against the seconds-of-staleness the display already tolerates.
Failure modes & mitigations
FailureDetection signalBlast radiusMitigation
Hot-key on a viral counter — stream/counterper-partition lag variance; per-shard CPU skewone viral video stalls; shard-mates suffersalt log keys; sharded cells; window pre-agg; edge absorbs reads
Consumer lag → counts hours stale — stream/aggregatorrecords-lag-max; staleness SLO breachevery count globally falls behindscale consumers to partitions; static membership; serve last-good
Double-count on retry — aggregator/counteringested ÷ unique-estimate > baseline; step-jumpsinflated counts → fraud-audit failsevent_id dedup; idempotent (video,region,window) UPSERT
Under-count on dropped events — ingest/streamedge-ingest vs log-accepted gapcounts read lowcommit after process; acks=all; DLQ residue, never silent drop
Dedup eviction (silent inflation) — dedupevicted_keys > 0 (page); mem > 80%fleet-wide inflation at peaknoeviction + page; separate dedup Redis; size for peak
Count goes backwards (checkpoint replay) — aggregatormonotonicity alarm on emitted countpublic "views went down"monotonic serve with verified-bit floor-reset (readsvc); event-time window_id; verifier writes a separate authoritative column + advances the watermark; externalized checkpoints
Silent gray failure (counts wrong, all 200s) — pipeline-wideindependent canary probe-delta miss + "canary-not-reporting"every displayed countout-of-band canary outside the pipeline's failure domain; reads ground truth direct
Bot-filter FP/FN — gw/verifierdrop-rate deviation; engagement-ratio anomalyFP: legit views dropped; FN: fraud inflatesGIVT inline + SIVT batch; shadow-score; manual-review queue
Poison / schema-bad event — stream/aggregatorconsumer restart loop on fixed offsetone partition stalls (one video)DLQ + skip (dlq); registry BACKWARD compat
Cache stampede on fresh-viral — countcache/counterorigin QPS spike on cache misscounter store browns outsingle-flight (readsvc); SWR at edge; 503 load-shed
Analytics compaction / replica lag — analyticssegment/compaction lag; replica seconds-behinddivergent dashboards (not the count)route to caught-up replicas; bound compaction backlog
IVT feed outage — ivtcircuit-breaker open; feed error/latencyverified count delayed / stale-signalcircuit breaker + cached verdicts; off the hot path
Open questions for human review
  • Single-region canvas, active-active posture. We bake active-active into the counter store's internals rather than drawing a second region's ingest+counter tier (the teaching focus is dedup/bot/aggregation, not cross-region staleness). Confirm that altitude vs a fully duplicated multi-region drawing.
  • Counter store = Cassandra rollup cells vs DynamoDB vs a custom counter (Netflix-style TimeSeries + EVCache). All three are real; we chose idempotent Cassandra cells. A reviewer who wants Netflix's exact event-store-+-rollup-+-EVCache shape would add a TimeSeries node.
  • Realtime OLAP ingest. We route analytics through the aggregator (pre-aggregated rollups). Pinot/Druid commonly ingest from Kafka directly; that would be a stream → analytics edge instead of (or in addition to) aggregator → analytics.
  • GIVT depth at the gateway. We keep inline filtering to signature/heuristic (no per-view external RPC). Confirm that's the right line vs a lightweight in-region bot-scoring sidecar for a subset of high-value beacons.
  • Verified-count UX. Showing a count that can drop is a product decision. Some platforms hide the drop (serve monotonic and silently correct internal/ad numbers); confirm whether the public number is allowed to decrease.
  • Dedup fail-open vs fail-closed. On a dedup-store outage we fail OPEN (the per-(viewer,video) suppression leaks) backstopped by the aggregator's in-window event_id dedup and the daily verifier; a team that treats any inflation as unacceptable might fail CLOSED (drop views during the outage) instead. Confirm the posture for this product.

Primary sources

  • Netflix Distributed Counter Abstraction (TechBlog, 2024)
  • Cloudflare HTTP analytics @6M req/s on ClickHouse
  • Reddit View Counting (HLL + Kafka + Cassandra)
  • Twitter TSAR / Summingbird + Manhattan
  • LinkedIn 'Who Viewed Your Profile' on Apache Pinot
  • Ably: Cassandra counter columns — nice in theory, hazardous in practice
  • Google Ad Manager GIVT/SIVT invalid-traffic methodology

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 View Count on a Video/Post yourself