Architecture One-Pager
Solution Architecture v1.0 · Google Cloud · Platform Architecture · 2026-10
Leaderboard & Counting Service · Solution Architecture v1.0 · Google Cloud · Platform Architecture · 2026-10
The event log is the system of record; every counter, window and ranking is a disposable, versioned projection of it — and the only freshness the platform spends is on the member's own contribution.
Almost every product people use daily puts a number in front of them: the view count under a video, the upvote tally on a post, the weekly league standing, the segment leaderboard, the end-of-year "top 2% of listeners" card. Each looks trivial and is not. It is a high-cardinality, high-write, skew-prone, user-visible aggregate, and it is wrong in a way people notice immediately. The convenient way to build it — a transactional counter store where the row is the truth — has no answer to the four things that arrive in the first year of operation: a partition that loses recent writes, a ranking build that overwrites a million rows with the wrong tie-break, a fraud campaign discovered six weeks late, and one key taking a thousand times the median write rate. Each is survivable alone; together they describe a platform that cannot be corrected. The second problem is smaller and sharper: the one reader who can tell that a count is stale is the member who just earned it, and a list that loads in 40 milliseconds without their last action reads as broken while a list twenty seconds behind that contains it reads as fine.
Make the append-only event log the only system of record, and every bucket aggregate, ranked view, histogram and cache entry a versioned projection of it that may be dropped and rebuilt at any time. Acknowledge a write when the event is durable, never when it has been counted, so a 4× burst is a capacity question rather than a correctness one. Deduplicate on a client-supplied idempotency key in the data model rather than trusting any transport. Attribute by event time with a declared lateness horizon and a visible late bucket. Let each counter declare whether it is exact, approximate or unique-counting, and show the price of that choice where it is made. Shard hot keys adaptively at ingestion and fan in on read, so the hottest key's cost stays with the hottest key. Answer rank-of-member from a value histogram outside the materialised head, because an exact rank at position 41,208 of 10 million is a query nobody will pay for. Then spend the platform's one piece of freshness on an overlay of the acting member's own deltas, declare the projection version and as-of on every response, retract by compensating event, and close a season explicitly so a reward can be paid against a standing that cannot silently change.
What it is, and what it is not
- An append-only log as the record, counters as projections — not A counter row that is the truth, with the log as an audit trail
- 202 means the event is durable — not 202 means the counter has moved
- Idempotency key in the data model — not Exactly-once delivery from the broker
- Event time with a declared lateness horizon — not Whenever the event happened to arrive
- Counter type declared per counter, with its price shown — not A platform-wide accuracy policy
- Read-your-writes for the actor, eventual for the audience — not Strong consistency for every reader
- Every response carries its version and as-of — not A stale answer that looks live
- Retraction as a compensating event, closure as a published boundary — not A decrement, and a standing that can be quietly rewritten
The decisions that are the architecture
- The log is the record; everything else is a projection (ADR-01) — Bucket aggregates, ranked views, histograms and caches are all derived, versioned and droppable. A corrupt partition, a bad ranking deploy, a wrong tie-break rule and a late-discovered fraud campaign therefore have one remedy — rebuild — instead of four disaster procedures, and rolling back a ranking is a pointer move.
- A write is acknowledged on durability, never on being counted (ADR-02) — The accept path's only synchronous dependency is the event log, so write latency is stable under pipeline restarts, projection rebuilds and cache loss. The weaker promise is the price, and the overlay rather than the documentation is what makes it acceptable to callers.
- Rebuild is routine, not a disaster procedure (ADR-03) — Every tenant's largest board is rebuilt monthly at ten times real time and diffed against live. That makes the 30-minute RTO a measured number, and makes a difference between rebuild and live the earliest possible warning that an accumulator has stopped being associative.
- Exactly-once effect comes from the data model (ADR-04) — A client-supplied idempotency key, deduplicated at admission and again in the pipeline over a 24-hour horizon. This covers the duplicate the broker never sees — the client that retried before the first request was acknowledged — which is the most common one in a consumer product.
- Event time, bounded, with the overflow made visible (ADR-05) — A ride that synced four hours late counts in the week it happened, up to a declared horizon; beyond it, a visible late bucket rather than a silent drop or a silent fold-in. Bounded lateness is what allows a window to be declared complete, which is the precondition for sealing a standing at all.
- Accuracy is declared per counter and carries a price (ADR-06) — Exact, approximate or unique-cardinality is a property of the versioned definition, enforced by the platform, with the unit cost per million events shown in the definition diff. That is what makes a distinct-viewer counter affordable at all, and what stops an estimate quietly ending up behind a number someone is paid against.
- The hottest key pays for itself (ADR-07) — Per-key write rate is detected at admission and the key's state is split across shards, fanned in on read, with shard count adapting as the key heats and cools. The alternative designs make either the median key or the 5-second RPO pay for the viral one.
- A rank outside the head is a percentile (ADR-08) — Exact ranks inside a capped materialised head; a value histogram beyond it, giving a constant-time percentile. Often the cheaper answer is also the better one — "top 3%" ages more gracefully than "#41,208" — and scoped boards exist to give an exact rank that means something.
- Freshness is spent where it is perceived (ADR-10) — The acting member's own accepted-but-unprojected deltas are applied at read time over a stale projection, so their own value and rank band are right within a second while everyone else's view of the ranking stays on the freshness budget. One extra keyed read on 30% of requests buys the only consistency a member can detect.
- Stale but honest, and stable within a session (ADR-11) — Every response carries its projection version and as-of, and a session is pinned to one version so a rank cannot oscillate between reads. Conditional reads make a polling leaderboard screen nearly free, which matters at 400,000 reads a second.
- Removal is a compensating event, and closure is published (ADR-12) — Counters are never mutated: a deletion, a ban, a disqualification or a fraud takedown is a compensating delta with cause, requester and authority, and affected rankings are recomputed rather than patched. After an explicit closure point, a retraction produces a correction record beside the original standing rather than a silent rewrite.
- Suspect contributions are held, not judged (ADR-14) — An inflation signal is probabilistic, so the contribution is quarantined — outside the counters, retained with original timestamps, reversible — and a freeze stops publication without stopping ingestion. A false positive then costs a delay rather than a member's effort, which is what buys the right to tune detection aggressively.
Why this should still be right in ten years
Brokers, stream processors, wide-column engines and sketch algorithms will all be replaced inside this platform's life. These are the properties that should outlast them.
- The boundary is not a technology. "The log is the record and everything else is a projection" is a statement about where authority lives, not about which broker carries it. Pub/Sub, Kafka, an object-store-native log or something not yet built can each satisfy it, and every decision downstream — versioned projections, rebuild as routine, retraction by compensation, rollback as a pointer move — survives the substitution unchanged.
- Skew is a property of human attention. One video, one post, one account taking three orders of magnitude more traffic than the median is a fact about how people pay attention, and algorithmic distribution has made it more extreme rather than less. No future storage engine removes the need to decide who pays for the hottest key; the engines that appear to handle it transparently are implementing some version of this decision internally.
- The exact-rank arithmetic does not improve. An exact global rank for an arbitrary member is a counting problem over a continuously changing set, and nothing about faster storage makes it cheap at 120,000 reads a second against 10 million members. If anything the pressure moves the other way, as products increasingly express standing as a percentile and the exact head becomes a smaller special case.
- Explicitly dated data gets more valuable. A response that states when it was true is correct; one that implies liveness is wrong. As more surfaces blend data of different freshness — and as generated summaries sit on top of live numbers — the value of carrying an as-of through the whole chain increases. This is a contract decision, so no implementation improvement makes it unnecessary.
- Append-only with compensation has centuries behind it. Recording corrections as entries rather than edits is the accounting profession's answer to the same problem, and it has lasted because a correction that leaves no trace is indistinguishable from a cover-up. Regulatory and cultural pressure on auditability has moved consistently in this direction and shows no sign of reversing.
- The costs are the ones that get cheaper. The price of this architecture is log storage, replay throughput, and one extra keyed read on a third of requests. All three are substrate properties that have improved by roughly an order of magnitude a decade. The price of the alternative — a user-visible number that cannot be corrected or defended — does not get cheaper with hardware.
Non-functional targets
Every figure below is a stated assumption for this design, chosen to be defensible and arguable rather than measured. A reviewer who changes one can follow it to the decision that depends on it.
| Quality | Target | How it is met | View |
|---|---|---|---|
| Ingestion availability | ≥ 99.99% monthly, measured as events durably logged and acknowledged | Stateless Cloud Run accept path across three zones; durability in a regionally redundant log before the 202 | 16 |
| Read availability | ≥ 99.95% monthly for top-N and rank | Cache tier plus fall-through to the ranked store; a cold cache is a latency event, not an availability one | 15 |
| Accept latency | p99 ≤ 25 ms single event, ≤ 80 ms for a 500-event batch | Validation, dedup lookup and publish only; no aggregation on the path | 13 |
| Read latency | Top-N p50 ≤ 10 ms, p99 ≤ 40 ms on cache hit, ≤ 150 ms on miss; rank p99 ≤ 120 ms | Version-keyed Memorystore snapshots; histogram percentile for ranks outside the head | 15 |
| Read-your-writes | Own value and rank band reflect an accepted event within 1 s | Query-time overlay of the member's own accepted-but-unprojected deltas | 04 |
| Global freshness | Staleness ≤ 5 s at p95, ≤ 30 s at p99; ≤ 120 s during a declared burst, flagged | Streaming aggregation plus incremental projection builds; as-of on every response | 13 |
| Throughput | 250,000 events/s steady, 1,000,000/s for 120 s; 400,000 reads/s | Log admits the burst, typed backpressure beyond it; read path scales independently of processing | 02 |
| Volume | 3.5 billion live counters, 120,000 active leaderboard scopes, 25 TB log per 90 days | Bigtable rows per (counter, member, bucket, shard); idle scopes dematerialised after 30 days | 11 |
| Retention | Log 90 d hot + 13 mo archive; buckets 400 d; standings and audit 5 y | Cloud Storage lifecycle classes; Bigtable cell TTLs; Spanner for sealed standings | 11 |
| Durability | RPO 5 s for accepted events; RPO 0 for standings and configuration | Replicated log before acknowledgement; multi-region Spanner for the exact state | 16 |
| Recovery | RTO 10 min read path in another region; RTO 30 min full projection rebuild | Warm standby over an async Bigtable replica; Dataflow replay at ≥ 10× real time | 16 |
| Accuracy | Exact counters zero error against replay; unique counters ≤ 2% at p95, ≤ 5% at p99 | 64-bit sums per shard; HyperLogLog++ sketches merged across shards and windows | 10 |
| Rank consistency | A reported rank agrees with its neighbours' reported values at the same version, 100% of the time | Ranking is a pure function of counter state at a version; sessions pinned to one version | 15 |
| Cost | ≤ $0.08 per million events ingested, aggregated and stored 90 d; ≤ $0.60 per million reads | Counter type priced at declaration; log tiering by age; idle-scope dematerialisation | 17 |
| Operability | A new counter or leaderboard live by configuration within one working day; config effective ≤ 60 s | Versioned definitions in Spanner, sandbox tenant for replay, per-field compatibility validation | 05 |
Scope
In scope
- A counting API with idempotency keys, batching, schema and bounds validation, quota enforcement and typed backpressure
- A durable, replayable event log as the system of record, with a tiered archive beyond the broker's retention
- Event-time aggregation with declared lateness, a visible late bucket, and exact, approximate and unique-cardinality counter types
- Adaptive hot-key write sharding with read fan-in, and a live shard map
- Versioned ranked views with deterministic tie-break, rank histograms, scoped boards and a version-keyed serving cache
- Rank-of-member, neighbourhood and percentile reads, with a read-your-writes overlay for the acting member
- Windows, seasons, non-blocking rollover and immutable closed standings
- Retraction by compensating event, bulk retraction by cause, correction records after closure, and reversible abuse quarantine with scope freeze
- Versioned counter and leaderboard configuration, per-tenant isolation and quotas, a disposable sandbox tenant
- Scheduled projection rebuild with diff against live, cost attribution per tenant and counter, and warehouse export
Explicitly out of scope
- What earns a point — the product decides and emits the event; the platform counts it
- The screens that render a number, and product notification logic beyond emitting a rank-change event
- Analytical querying of counting history; the warehouse receives an export and is never on the serving path
- Monetary ledgers and anything that must balance to the cent
- Real-time game state and matchmaking; the platform ranks outcomes, it does not run the game
- Member identity and display names, which are resolved from the product's own directory at read time
- A public third-party read API, active-active counting and progressive rollout of ranking algorithms, all named and deferred
What a four-week prototype should prove
The prototype's job is to falsify the two central claims — that a projection rebuilt from the log is identical to the live one, and that an overlay makes eventual consistency acceptable to the member who just earned the points — on one tenant and two boards, not to build a platform.
- One tenant, two leaderboards: one small board with exact ranks for everyone, one 10-million-member board with a materialised head and a histogram
- Drive 25,000 events/second with a deliberately skewed key distribution including one key at 1000× the median, and measure accept p99 and top-key read latency
- Build the scheduled rebuild and prove byte-identical output against the live projection, then break associativity deliberately and prove the diff catches it
- Implement the own-write overlay and measure, from the client's point of view, the time between a 202 and the member's own value being visible
- Seal a season while events continue to land in the next one, then retract a member post-closure and produce a correction record
- Quarantine a simulated inflation campaign, clear half the holds and replay them, and show the ranks end contiguous
- A client retries the same event five times with the same idempotency key: the counter moves once and the response echoes the key every time
- A batch of 500 events contains three with out-of-bounds deltas: 497 are accepted, three are rejected individually, and the batch does not fail
- A ride syncs four hours late: it lands in the week it happened; one that syncs eight days late lands in the visible late bucket and nowhere else
- The aggregation job is stopped for ten minutes: reads keep answering with an ageing as-of, the staleness flag flips, and no number is wrong
- A ranking build ships with an inverted tie-break: readers are pointed back at the previous version within seconds and no counter is touched
- The serving cache is flushed entirely at peak: reads fall through at the miss budget and no correctness signal moves
- A member's own counter is written by another member's token: the request is refused at admission on subject mismatch and logged against the attempt
Open risks, carried rather than hidden
| Risk | If it lands | Response |
|---|---|---|
| A product renders a stale answer as live, because the as-of field is easy to ignore | The platform's honest degradation becomes a wrong number on a screen, and the complaint lands on the platform | Make the as-of and the degradation flag part of the client SDK's rendered output rather than a field in a payload, publish the per-board freshness budget as a product-facing figure, and review new leaderboard surfaces for whether they show it. |
| The overlay and the projection disagree by more than the member's own delta | Two code paths producing one number, with no tie-break — the hardest class of bug to diagnose because both answers look plausible | Keep the overlay arithmetically trivial — add the member's own deltas, never recompute — and emit a continuous comparison between overlaid and projected values for sampled members, alerting on any divergence beyond the member's own outstanding deltas. |
| The scheduled rebuild is cut in a cost review, because it ships no feature | The rebuild path rots, and the first real use is also the first test, during an incident | Report the rebuild as a correctness control with a named owner and a measured RTO rather than as compute spend, and keep its diff output on the same dashboard as the user-visible freshness signal so its value is visible monthly. |
| A single write region makes a regional failure a counting outage, not just a ranking one | Member effort is lost unless every product client buffers and retries with stable idempotency keys — a contract the platform cannot enforce | Ship buffering and stable key generation in the client SDK rather than documenting it, test it by failing the write region in a game day, and report per-tenant retry success as a platform signal rather than assuming compliance. |
| Histogram bucket width is set once and silently governs the accuracy of every non-head rank | A board where "top 3%" is wrong by a percentage point for millions of members, with no signal that anything is off | Make bucket width a per-board setting with a published implied accuracy, validate it against the board's observed value distribution on each projection build, and alert when the distribution has moved far enough that the implied accuracy no longer holds. |
The reasoning behind every component and technology choice is in the Architecture Decision Record: 17 records across 6 areas, each with the alternatives that lost and what the choice costs.