Architecture Decision Record
Solution Architecture v1.0 · Google Cloud · Platform Architecture · 2026-10
Leaderboard & Counting Service · Solution Architecture v1.0 · Google Cloud · Platform Architecture · 2026-10
The argument these decisions serve is summarised in the Architecture One-Pager.
Seventeen decisions make up this architecture. Everything else across the twenty-one views is either a consequence of one of them or a detail that could be decided differently next quarter without anybody having to redraw the set. Each record carries more than the classic context / decision / consequences triple: the forcing question, how the decision is actually realised on Google Cloud, the alternatives including those that are right for a different organisation, the conditions that would flip the choice, why the choice should still hold as scale and technology change, and the transferable lesson.
Status of this document. This is a design, not a report on a running system. Every rate, latency, ratio, threshold and retention figure is a stated assumption chosen to be defensible and arguable rather than measured. The operating context assumed throughout is a consumer platform of 80 million monthly active members across 25 product tenants, accepting 250,000 counting events a second at steady state with a 4× burst for 120 seconds, serving 400,000 leaderboard reads a second split 70/30 between top-N and rank-of-member, over 3.5 billion distinct live counters in 120,000 active leaderboard scopes, with an event log of 25 TB per 90 days after compression. A reviewer who disagrees with a number can follow it to the decision that depends on it; that is what the numbers are for.
How to read a record
- Question: The forcing question: why a decision was needed at all.
- Context: The requirement, the scale and the constraint that make it hard.
- Decision: What this architecture does, stated so it can be checked.
- How it is realised on AWS: The concrete mechanism: which service or package, configured how, in which subscription.
- Options weighed: Chosen, rejected, deferred, or right elsewhere, with the reason for each.
- Consequences: What the choice buys and what it costs, both kept visible.
- Choose differently when: The conditions that would flip the decision for your system.
- Why it holds up over time: What keeps the decision right as scale, staff and technology change.
- Lesson: The principle that transfers beyond this platform.
Decision map
Truth and rebuildability: The decisions that make every number in the platform recoverable rather than precious.
- ADR-01 · The event log is the system of record; every counter and ranking is a disposable, versioned projection
- ADR-02 · A write is acknowledged when the event is durable, never when it has been counted
- ADR-03 · Rebuild is a scheduled routine at ten times real time, not a disaster procedure
Counting correctness: The decisions that make a total right under duplicate delivery, late arrival and a declared error budget.
- ADR-04 · Exactly-once effect comes from an idempotency key in the data model, not from the transport
- ADR-05 · Events are attributed by event time, with a declared lateness horizon and a visible late bucket
- ADR-06 · Counter type — exact, approximate or unique-cardinality — is declared per counter and carries its price
Skew and the cost of a rank: The decisions that stop the hottest key and the deepest leaderboard setting the price of everything else.
- ADR-07 · Hot keys are sharded adaptively at ingestion, with read fan-in
- ADR-08 · A rank outside the materialised head comes from a value histogram, not from a count-less-than query
- ADR-09 · Windows are pre-aggregated into the finest declared bucket and composed by addition
Freshness and honesty: The decisions about where freshness is spent, and how the platform admits what it does not know yet.
- ADR-10 · Read-your-writes is a query-time overlay of the member's own recent deltas
- ADR-11 · Every response declares its projection version and as-of time, and a session is pinned to one version
Results, retraction and integrity: The decisions that make a published standing trustworthy and a correction auditable.
- ADR-12 · Retraction is a compensating event, never a mutation of a stored counter
- ADR-13 · Closure is explicit, and a retraction after it produces a correction record rather than a rewrite
- ADR-14 · Suspect contributions are quarantined — held and reversible — rather than admitted or discarded
State placement and reach: The decisions about which store holds what, how far it reaches, and how a definition changes.
- ADR-15 · Counters live in a wide-column store; definitions, seasons and standings live in a relational one
- ADR-16 · One write region, with a read standby whose cache is rebuilt rather than replicated
- ADR-17 · Counter definitions are versioned, and a breaking change creates a new version rather than mutating the old one
Technology by capability
Google Cloud was chosen for this exercise deliberately, and for two reasons that point the same way. The first is rotation: across this repository's previous use cases Azure and self-hosted open source dominate, with Amazon Web Services next, and Google Cloud is the least-exercised of the three hyperscalers — so building here teaches something the practice has not already learned. The second is fit, which matters more. This architecture needs exactly two kinds of storage: a wide-column store where a row key can carry a shard component and a window range scans contiguously, and a strongly consistent relational store for the small exact part. Bigtable and Spanner are an unusually clean match for that pair, and Pub/Sub plus Dataflow give an event log and an event-time streaming model without assembling them. Every requirement in ask.md is written vendor-neutrally; the table below is where the architecture commits.
| Capability | Choice | Origin | Credible alternative | Why this one | Record |
|---|---|---|---|---|---|
| Durable event log | Pub/Sub topic with 31-day retention, one subscription per aggregation job | Google Cloud | Kafka on GKE or Confluent Cloud; Kinesis Data Streams; Event Hubs | Managed, regionally redundant, and 31-day retention covers streaming replay without the archive being on the hot path | ADR-01 |
| Replayable record beyond the broker | Avro on Cloud Storage, standard class for 90 days then archive class to 13 months | Google Cloud | Kafka tiered storage; S3 with lifecycle rules; a lakehouse table format | Retention of the record must be set by the architecture, not by the broker's maximum | ADR-01 |
| Event-time aggregation | Dataflow streaming with per-counter allowed lateness and keyed dedup state | Google Cloud + Apache Beam | Flink on GKE; Spark Structured Streaming; a hand-written consumer with RocksDB state | Event-time windowing with bounded lateness is the native model rather than a layer on top of it | ADR-05 |
| Bucket aggregate store | Bigtable, row key (counter, member, window bucket, shard) | Google Cloud | Cassandra or ScyllaDB; DynamoDB with a sharded sort key; HBase | A shard component in the row key makes hot-key splitting a data operation, and a window composes as a contiguous range scan | ADR-07 |
| Ranked views and rank histogram | Bigtable under a version key, built by a Dataflow batch job | Google Cloud | Redis sorted sets as the source of truth; a search index; a columnar store | Versioned ordered rows make publish and rollback pointer moves, and keep the ranked view a projection rather than a record | ADR-08 |
| Serving cache | Memorystore for Redis, HA tier, version-keyed snapshots | Google Cloud | Self-managed Redis or Valkey; a CDN cache for public boards | Top-N at p99 ≤ 40 ms needs an in-memory tier, and a cache with no RPO is capacity rather than state | ADR-11 |
| Control plane, seasons and closed standings | Spanner, multi-region, RPO 0 | Google Cloud | Cloud SQL or AlloyDB with regional replicas; CockroachDB; DynamoDB with transactions | Closure and retraction span several entities transactionally, and this is the only state that is restored rather than rebuilt | ADR-15 |
| Request services | Cloud Run for the counting, query, config and operations APIs | Google Cloud | GKE Autopilot; Cloud Functions; EKS or AKS equivalents | Stateless request services that must scale with a 4× burst and hold no shared mutable state in the accept path | ADR-02 |
| Edge protection and rate limiting | Global external load balancer with Cloud Armor rules per tenant | Google Cloud | Cloudflare or Fastly in front; an API gateway with per-key quotas | Per-member and per-tenant limits are cheapest to enforce before a request reaches the accept path | ADR-14 |
| Member authentication | Identity Platform tokens, verified at the edge, subject bound to the member reference | Google Cloud | The product's own OIDC provider; Auth0 or Okta | The one check that matters is that the token's subject equals the member being written, and it belongs before admission | ADR-17 |
| Season close and rebuild orchestration | Workflows driven by Cloud Scheduler, launching Dataflow batch jobs | Google Cloud | Airflow on Composer; Argo Workflows; Step Functions | Closure and scheduled rebuild are short, auditable, low-frequency orchestrations, not data pipelines in their own right | ADR-03 |
| Observability | Cloud Monitoring, Logging and Trace, with projection lag per leaderboard as a first-class metric | Google Cloud | Prometheus and Grafana on GKE; a commercial APM | The signal that matters most is a derived per-board freshness number, which needs a custom metric wherever it is hosted | ADR-11 |
| Analytics export | BigQuery, loaded from the event archive and the closed standings on a schedule | Google Cloud | Snowflake or Databricks over the same archive | Analysts must never be on the serving path, and the archive is already the right shape to load | ADR-01 |
The decisions, and the alternatives that lost
Truth and rebuildability
The decisions that make every number in the platform recoverable rather than precious.
ADR-01 · The event log is the system of record; every counter and ranking is a disposable, versioned projection
Status: Accepted · Shown on views: 02, 11, 18
When a counter's value and the log disagree, which one is wrong?
Context. The obvious way to build a counting service is a transactional counter store: the value lives in a row, every event increments it, and the row is the truth. That design is simpler to reason about and gives immediately consistent reads. It also has no answer to four questions that arrive within the first year of operation. A partition loses recent writes — what is the correct value now? A ranking build ships with a wrong tie-break rule and overwrites a million rows — what do we restore from? A fraud campaign is discovered six weeks late — how are its contributions removed from counters that have been incremented in place? And the hottest key in the estate takes a thousand times the median write rate, so the row that is the truth is also the row with the deepest contention. Each of those is survivable individually; together they describe a platform that cannot be corrected.
Decision. The durable append-only event log is the only system of record. Every bucket aggregate, every ranked view, every histogram and every cache entry is a projection of it, carries a version, and may be dropped and rebuilt at any time without loss. A projection and the log never disagree: if they differ, the projection is wrong by definition and is rebuilt.
How it is realised on AWS. Accepted events are published to a Pub/Sub topic with 31-day retention and continuously archived as Avro to Cloud Storage, hot for 90 days and archive class to 13 months. A streaming Dataflow job materialises bucket aggregates into Bigtable; a batch job builds ranked views and rank histograms into Bigtable under a version key; Memorystore holds published top-N snapshots keyed by the same version. Nothing but the log and the archive is backed up — the rest is rebuilt from them by a Dataflow replay job, which runs on a schedule and not only after an incident.
| Option | Verdict | Reasoning |
|---|---|---|
| Append-only log as the record, counters and rankings as versioned projections | Chosen | Makes every downstream store correctable, and turns four separate disaster procedures into one routine operation |
| Transactional counter store as the record, log as an audit trail | Rejected | Immediately consistent and simpler to read; has no rebuild story, and couples the write path to the hottest key's contention |
| Counter store as the record with periodic snapshots for recovery | Rejected | Snapshots recover a point in time, not a correction — a fraud campaign discovered later is still unpickable from an incremented row |
| Log as the record with no retained archive, relying on broker retention | Right elsewhere | Right where the rebuild window is hours, not months; here the broker's retention would set the record's retention |
What it buys
- A corrupt partition, a bad ranking deploy, a wrong tie-break rule and a late-discovered fraud campaign all have the same remedy: rebuild
- Rolling back a ranking becomes pointing readers at the previous projection version, which is a seconds-long operation with no data movement
- The write path owes the caller only durability, which is what makes a 4× burst a capacity question rather than a correctness one
What it costs
- Counter values are eventually consistent by construction, so read-your-writes has to be engineered deliberately (ADR-10) rather than inherited
- The log must be retained long enough and replayable fast enough to rebuild the largest leaderboard inside the RTO — a throughput commitment, not a backup
- Two copies of the record (live topic and archive) double the ingest storage floor and add a boundary a replay must cross correctly
Choose differently when. If every counter in the estate were small, bounded, never disputed and never ranked — a feature-usage tally read by nobody but a dashboard — the rebuild machinery would be dead weight and a transactional counter store would be the right answer. The decision is justified by counters that are user-visible, occasionally disputed, and sometimes paid against.
Why it holds up over time. "The log is the record" is a statement about where authority lives, not about which broker carries it. Pub/Sub, Kafka, a successor nobody has shipped yet, or an object-store-native log can each satisfy it, and every decision downstream — versioned projections, rebuild as routine, retraction by compensation — survives the substitution unchanged. The cost of the rule is storage and replay throughput, both of which have fallen by roughly an order of magnitude a decade; the cost of abandoning it is an uncorrectable number, which does not fall at all.
Lesson. A number people see will eventually be wrong. Decide, before you build, whether being wrong is a repair or an incident — and make the repair the cheap path.
ADR-02 · A write is acknowledged when the event is durable, never when it has been counted
Status: Accepted · Shown on views: 02, 08, 13
What has the platform promised when it returns 202?
Context. A product backend emitting 250,000 events a second needs a hand-off that is fast and, more importantly, one whose latency does not depend on anything downstream. If the acknowledgement waits for the aggregation pipeline, then a pipeline slowdown becomes a product-wide write latency event, a pipeline restart becomes an outage, and a rebuild — which deliberately replays the log at ten times real time — competes with live traffic for the same acknowledgement path. The temptation to wait is real: a caller that gets a 202 and then cannot see its own effect looks broken, and the natural fix is to make the acknowledgement mean more.
Decision. The Counting API acknowledges a batch once every accepted event in it is durably committed to the replicated log, and never waits on aggregation, projection or cache. The 202 means "this event will be counted", not "this event has been counted". The gap that creates for the caller's own read is closed by the overlay in ADR-10, not by making the write slower.
How it is realised on AWS. Admission validates schema, bounds and clock, checks the idempotency key, applies quota and the inflation signal, then publishes to Pub/Sub and waits for the publish acknowledgement before responding. Per-event acceptance is reported inside a batch response so one rejected event does not fail the other 499. When the log cannot admit — a publish latency breach or a quota ceiling — the API returns a typed rejection with a retry hint rather than queueing in memory.
| Option | Verdict | Reasoning |
|---|---|---|
| Acknowledge on log durability, aggregate asynchronously | Chosen | Decouples the product's write latency from every piece of processing the platform does |
| Acknowledge after the counter is updated | Rejected | Gives the caller a stronger promise and makes every pipeline event a product incident |
| Acknowledge on receipt, before durability | Rejected | Fastest, and silently converts an instance crash into lost member effort — the one loss members notice |
| Synchronous counter update with asynchronous ranking | Right elsewhere | Reasonable where counters are few and cold; at 3.5 billion keys with a 1000× hot key it reintroduces contention on the write path |
What it buys
- Write latency is a function of the log alone, so the p99 accept budget of 25 ms is defensible and stable under rebuilds
- A pipeline restart, a projection rebuild or a cache loss are all invisible to writers
- Backpressure becomes an explicit, typed product contract instead of an unbounded in-memory queue
What it costs
- The acknowledgement is weaker than callers expect, and documentation will not fix that — the overlay has to
- A member's effect is visible to that member before it is visible globally, which is a behaviour the product UI must be designed around
- An event that is durable but unprocessable still counts as accepted, so the poison sink becomes a monitored queue rather than an edge case
Choose differently when. If the product surface genuinely required a globally consistent count at the moment of the write — a live auction, a seat allocation, an inventory reservation — this decision would be wrong, and the right design would be a transactional reservation service rather than a counting platform. The split here is that a leaderboard is a view of accumulated effort, not a contended resource.
Why it holds up over time. The property that outlasts the technology is that a system's synchronous promise should be the smallest one its callers can build on. Every later generation of this platform will have faster aggregation and will be tempted to fold it back into the acknowledgement; the reason not to is that the promise then inherits the availability of everything behind it, and that arithmetic does not improve with hardware.
Lesson. Promise the least you can and still be useful. A strong acknowledgement is borrowed from the availability of everything it waits for.
ADR-03 · Rebuild is a scheduled routine at ten times real time, not a disaster procedure
Status: Accepted · Shown on views: 10, 11, 18
How does the platform know its rebuild works before it needs it?
Context. A rebuild path that exists only in a runbook does not exist. It depends on an archive format nobody has read in months, a replay job whose dependencies have moved, a Bigtable write path sized for live traffic rather than for a ten-times burst, and an accumulator whose associativity nobody has checked since the last change. Every one of those fails quietly until the day a projection is corrupt, which is exactly the day there is no time to discover them. The second reason to run it routinely is subtler: a rebuild of a healthy projection should produce an identical result, so a difference is the earliest possible signal that an accumulator has stopped being commutative or that a late-event rule has drifted.
Decision. Projection rebuild from the retained log is a first-class, scheduled operation. Every tenant's largest leaderboard is rebuilt at least monthly from the archive at a replay rate of at least ten times real time, the output is compared against the live projection, and a difference is a correctness alert rather than an expected variance. The same mechanism serves recovery, so the recovery path is the exercised path.
How it is realised on AWS. A Dataflow batch job reads the Avro archive for a tenant, counter and time range, runs the same transform library as the streaming job, and writes to a new projection version rather than over the live one. A Workflows pipeline orchestrates the schedule, the comparison and the promotion-or-alert decision. Replay throughput is provisioned as a declared burst against Cloud Storage and Bigtable, with the rebuild job running at a lower Bigtable app-profile priority than the live write path.
| Option | Verdict | Reasoning |
|---|---|---|
| Scheduled rebuild with diff against live, same code path as recovery | Chosen | Proves the recovery path monthly and turns accumulator drift into an alert instead of a mystery |
| Rebuild only on incident | Rejected | Cheapest in compute; the first real use is also the first test, under time pressure |
| Continuous shadow pipeline running permanently alongside live | Rejected | Strongest correctness signal and roughly doubles steady-state processing cost for a property a monthly diff establishes |
| Periodic snapshot-and-restore instead of replay | Right elsewhere | Right when the transform is expensive and stable; here the transform is cheap and the point is to re-derive, not to restore |
What it buys
- The 30-minute RTO for the largest leaderboard is a measured number rather than an estimate
- Accumulator drift, archive format breakage and replay throughput regressions are all caught by the same monthly job
- A disputed standing can be answered by re-deriving it, which is a materially stronger answer than reading the current value
What it costs
- A recurring compute and read cost with no feature attached, which is the first line cut in a cost review and must be defended as insurance
- Replay at ten times real time is a large bursty read against the same storage the ingest path writes, so it needs its own capacity envelope
- Comparison needs a tolerance model for approximate counters, or every sketch-backed board reports a false difference
Choose differently when. If the log were small enough to replay in minutes on demand, the schedule would be unnecessary — any failure could simply be rebuilt when it happened. The schedule is justified by a 25 TB hot archive where a cold rebuild path would take days to debug.
Why it holds up over time. Exercising a recovery path on a schedule is one of the few operational practices that has never been made obsolete by better infrastructure. Managed services reduce the number of things that can break; they do not make an unexercised path trustworthy, and they add a new failure mode — a platform change that silently alters replay semantics — that only a routine run detects.
Lesson. A recovery path you do not run is a hypothesis. Run it often enough that it is boring, and use its output as a correctness signal rather than just a drill.
Counting correctness
The decisions that make a total right under duplicate delivery, late arrival and a declared error budget.
ADR-04 · Exactly-once effect comes from an idempotency key in the data model, not from the transport
Status: Accepted · Shown on views: 08, 13, 14
When the same event arrives twice, what stops the counter moving twice?
Context. Every delivery path in this architecture is at-least-once. A mobile client retries on a timeout it cannot distinguish from a failure; a product backend retries a batch whose response was lost; the log redelivers on consumer restart; a pipeline worker is replaced mid-bundle. Transport-level exactly-once exists in some brokers but only within the broker's own boundary, and it does nothing about the client that retried before the first request was ever acknowledged — which is the most common duplicate in a consumer product and the one a member can cause by tapping twice.
Decision. Every counting event carries a client-supplied idempotency key, and deduplication is performed on (tenant, counter, idempotency key) over a declared horizon of at least 24 hours. Deduplication happens twice, at different boundaries: at admission, so a client retry costs a lookup rather than a pipeline reprocess, and inside the pipeline, so a log redelivery cannot double-count. The transport's own guarantees are treated as a performance property, not a correctness one.
How it is realised on AWS. Admission checks a Bigtable dedup table keyed by (tenant, counter, idempotency key) with a 24-hour cell TTL and responds with an idempotency echo so a client can retry safely. The Dataflow job holds keyed state on the same key for the same horizon, which bounds the pipeline's largest state by the horizon rather than by the counter count. A duplicate is counted as a signal — dedup hit rate per tenant — because a rising rate is a client bug worth telling a team about.
| Option | Verdict | Reasoning |
|---|---|---|
| Client-supplied idempotency key, deduplicated at admission and in the pipeline | Chosen | Covers client retries, broker redelivery and worker replacement with one mechanism and one key |
| Rely on the broker's exactly-once delivery | Rejected | Cheapest to implement and silent about the duplicate that happens before the broker is involved |
| Server-generated event id with content hashing for duplicate detection | Rejected | Needs no client cooperation; cannot distinguish a retry from two genuine identical actions a second apart |
| Idempotent counters by construction — set-to-value rather than increment | Right elsewhere | Right where the client knows the absolute total, such as a synced device counter; most products only know the delta |
What it buys
- A retry is free and safe, which lets every client retry aggressively and removes lost-event anxiety from the product teams' designs
- Deduplication is testable and observable: the hit rate is a number per tenant rather than a property of a broker configuration
- The pipeline's state size is bounded by the dedup horizon, which makes its restart time predictable
What it costs
- Keyed dedup state at 250,000 events/second over 24 hours is the pipeline's largest state and its slowest restart
- A duplicate arriving after the horizon is counted, so the horizon is a stated correctness boundary rather than a guarantee
- Clients must generate stable keys, and a client that regenerates the key on retry defeats the mechanism invisibly
Choose differently when. If every event carried an absolute value rather than a delta, deduplication would be unnecessary because applying the same value twice is a no-op. The decision is forced by delta semantics, which the products need because they do not hold the running total.
Why it holds up over time. Putting idempotency in the data model rather than in the transport has survived every generation of messaging technology, because the duplicate it protects against is generated outside the transport. As clients become more aggressive about retrying — and offline-first mobile behaviour has only pushed that further — the fraction of duplicates the broker never sees grows, which strengthens the decision over time.
Lesson. Put the deduplication key where the retry is generated, not where the message is delivered. The transport cannot see the duplicate that happened before it was involved.
ADR-05 · Events are attributed by event time, with a declared lateness horizon and a visible late bucket
Status: Accepted · Shown on views: 10, 12, 14
A member's ride synced four hours after it finished — which week does it count in?
Context. Mobile clients buffer offline. A lesson completed on a plane, a ride recorded in a valley, a round played on a train all arrive hours after they happened, and sometimes after the week they belong to has been published. Attributing by arrival time is trivial to implement and produces a visibly wrong answer: effort appears in the wrong week, and a member who was offline is penalised for it. Attributing purely by event time is correct and unbounded — a window can never be declared complete, so no standing can ever be closed. Both failure modes are user-visible, in opposite directions.
Decision. Every event is attributed to the window its event timestamp falls in, up to a lateness horizon declared per counter. Events later than the horizon are attributed to a visible late bucket rather than being dropped or silently folded into the current window, and the server's own ingestion time is retained alongside the client's event time so clock skew is detectable and auditable.
How it is realised on AWS. The Dataflow job windows on event time with the counter's declared allowed lateness, and emits beyond-horizon events to a separate late column in the bucket row. Admission rejects timestamps more than 7 days old or more than 60 seconds in the future outright. The late-event ratio per tenant is an emitted signal, because a rising ratio means either a client problem or a horizon set too tight for how the product is actually used.
| Option | Verdict | Reasoning |
|---|---|---|
| Event time with a declared, per-counter lateness horizon and a visible late bucket | Chosen | Correct for the member, bounded for the platform, and honest about what fell outside |
| Arrival time only | Rejected | Trivial and robust; puts an offline member's effort in the wrong week, which is the complaint that arrives first |
| Unbounded event time with windows reopened on late arrival | Rejected | Most correct in principle; no window can ever be closed, so no standing can be sealed or paid against |
| Event time with late events silently discarded | Right elsewhere | Defensible for a pure engagement metric nobody competes on; indefensible when the number decides a reward |
What it buys
- A member who was offline is not penalised, which removes a class of support contact that is impossible to answer well
- A window can be declared complete at a known time, which is the precondition for sealing a standing at all
- The late bucket makes the trade-off measurable instead of a silent policy
What it costs
- The pipeline must hold window state for the full horizon, which sets its memory footprint and its restart cost
- Late events arriving after a window is published change a number a member may already have seen, so the as-of becomes load-bearing
- A visible late bucket needs a product answer — what a member is told about effort that arrived too late — that the platform cannot supply
Choose differently when. If every client were online and server-authoritative — a web-only product where the backend emits at the moment of the action — event time and arrival time would be within milliseconds and the whole mechanism would be redundant. The decision is justified by mobile clients that buffer for hours.
Why it holds up over time. Event time versus processing time is one of the few genuinely settled results in stream processing, and it is settled in favour of event time with bounded lateness. Connectivity improving does not remove the problem: as products move further into offline-first and background sync, the gap between when something happened and when the platform hears about it gets wider, not narrower.
Lesson. Attribute by when it happened, bound how long you will wait, and make what fell outside the bound visible. Two of the three are not enough.
ADR-06 · Counter type — exact, approximate or unique-cardinality — is declared per counter and carries its price
Status: Accepted · Shown on views: 05, 10, 12
Which of these numbers is allowed to be slightly wrong, and who gets to decide?
Context. Counting distinct viewers of a video is an entirely different problem from counting its views. The view count is a sum: cheap, exact, one number per key. The distinct-viewer count needs to know who has already been seen, which means either storing the member set — at 3.5 billion counters, the dominant cost in the platform by an order of magnitude — or accepting a bounded error from a sketch. Platform teams tend to resolve this by picking one answer for everyone: exact everywhere, which prices unique counting out of existence, or approximate everywhere, which eventually puts an estimate behind a number someone is paid against.
Decision. Each counter declares its type — exact, approximate, or unique-cardinality — along with its bounds, monotonicity and tolerated error, and the declaration is part of the versioned definition. The platform enforces the declaration and refuses writes that contradict it, and it publishes the unit cost per million events for each type so the choice is made against a price rather than a preference.
How it is realised on AWS. Exact counters accumulate as 64-bit sums per (key, window bucket, shard) in Bigtable. Unique-cardinality counters accumulate HyperLogLog++ sketches, mergeable across shards and windows without storing member identity, with a declared relative error of ≤ 2% at p95 and ≤ 5% at p99. The counter definition in Spanner carries type, bounds, lateness horizon, client-writability and the unit cost figure that the tenant console shows in the definition diff.
| Option | Verdict | Reasoning |
|---|---|---|
| Type declared per counter, enforced, with unit cost shown at declaration time | Chosen | Puts the cost-versus-accuracy choice in front of the person who knows what the number is for |
| Exact everywhere | Rejected | Simplest contract and simplest support answer; makes distinct-counting the dominant storage line in the platform |
| Approximate everywhere above a size threshold | Rejected | Cheap and uniform; the threshold is a platform decision silently applied to a number that may decide a reward |
| Exact with periodic sketch-based estimation for dashboards only | Right elsewhere | Right where approximate numbers are never member-facing; here uniques are exactly what the member-facing card shows |
What it buys
- A distinct-viewer counter becomes affordable at a stated error, which makes a whole category of product surface possible at all
- Members' and auditors' questions about accuracy have a documented answer per counter rather than a platform-wide hedge
- Cost attribution becomes meaningful, because the expensive choice is visible at the moment it is made
What it costs
- Three accumulator families to implement, test and keep associative, rather than one
- A sketch is mergeable but not retractable: removing a member from a unique counter requires rebuilding that window's sketch from the log
- Tenants will occasionally choose exact where approximate would do, because exact is the safe-sounding word
Choose differently when. If uniqueness were never counted — only sums and maxima — the whole approximate branch would be unnecessary. It is justified by product surfaces that show distinct participants, distinct viewers and distinct listeners, which members read as headline numbers.
Why it holds up over time. Sketch algorithms improve and will be replaced; what endures is that accuracy is a product requirement with a price, and that the price belongs in front of the person choosing. Storage getting cheaper narrows the gap slowly, but the ratio between a sum and a member set grows with cardinality, so the choice does not disappear as the estate grows — it gets more consequential.
Lesson. When accuracy costs money, make the price visible at the point of choice and enforce the choice in the data model. A platform-wide accuracy policy is a decision taken by whoever is least informed about the number.
Skew and the cost of a rank
The decisions that stop the hottest key and the deepest leaderboard setting the price of everything else.
ADR-07 · Hot keys are sharded adaptively at ingestion, with read fan-in
Status: Accepted · Shown on views: 08, 12, 17
One key is taking a thousand times the median write rate — who pays for that?
Context. Skew is not an edge case in a counting platform; it is the shape of the traffic. One video, one post, one segment or one celebrity account routinely takes three orders of magnitude more writes than the median key, and in a wide-column store a single row is a single contention point on a single tablet. The three available answers price differently. A fixed shard count per counter makes every key pay the hottest key's read amplification. Local pre-aggregation in the ingest process collapses a thousand writes into one, which is the cheapest answer and introduces a flush interval during which counts live only in memory. Adaptive splitting costs a control loop and a shard map.
Decision. The admission tier detects per-key write rate, assigns a shard number from a live shard map, and the aggregation job writes per-shard state that is fanned in on read. Shard count per key adapts to observed rate and reduces again when the key cools, so the median key does not pay the hottest key's read amplification. Counts are never held only in process memory for the sake of absorbing a burst.
How it is realised on AWS. Bigtable row keys include a shard component, so splitting a key is a change to the shard map rather than a schema migration. The shard map lives beside the aggregates and records observed rate, current shard count and split and merge times. The projection builder fans in shards when materialising a ranked view; a direct read of a hot counter fans in at read time. Top-key write share is an emitted signal, and the splitter's reaction time is held to the burst envelope the platform claims to absorb.
| Option | Verdict | Reasoning |
|---|---|---|
| Adaptive per-key sharding at ingestion with read fan-in | Chosen | The hottest key's cost stays with the hottest key, and no count ever lives only in a worker's memory |
| Fixed shard count per counter | Rejected | No control loop to build; every one of 3.5 billion keys pays the fan-in cost of the worst one |
| Local pre-aggregation in the ingest process with periodic flush | Rejected | Cheapest by a wide margin; a worker crash loses a flush interval of member effort, which the 5-second RPO does not allow |
| A dedicated path for a hand-maintained list of hot keys | Right elsewhere | Workable where hotness is predictable, such as scheduled live events; a viral post arrives without warning |
What it buys
- A viral key degrades its own read cost and nothing else, which keeps the 40 ms top-N budget credible across the estate
- Splitting and merging are data operations on the shard map, with no schema change and no migration
- Hot-key behaviour becomes observable — shard count per key is a number an engineer can look at during an incident
What it costs
- A control loop that must react faster than a launch arrives, and whose oscillation would itself be a cost problem
- Read fan-in means a hot key's read is several row reads, so the hottest key has the worst read latency in the system
- The shard map is a new piece of state on the hot path of admission, cached per instance and therefore eventually consistent
Choose differently when. If the write distribution were close to uniform — a per-member counter with a hard rate limit, for instance — sharding would be pure overhead and a single row per key would be correct. The decision is justified by a measured expectation of 1000× skew, which is normal in any content or social product.
Why it holds up over time. Skew is a property of human attention, not of infrastructure, and it has become more extreme as distribution has become more algorithmic. No future storage engine removes the need to decide who pays for the hottest key; the engines that handle it transparently do so by implementing some version of this decision internally, which is an argument for understanding it rather than for delegating it blindly.
Lesson. Size the write path for the hottest key and the read path for the median one. A design that averages the two is wrong at both ends.
ADR-08 · A rank outside the materialised head comes from a value histogram, not from a count-less-than query
Status: Accepted · Shown on views: 04, 12, 15
What is the rank of the member at position 41,208 of 10 million, and what is it worth?
Context. "Where am I" is 30% of reads and the single most expensive query in the system. For a member inside the materialised top N it is a lookup. For everyone else it means counting how many members have a value greater than theirs, which is a full scan, an index maintained on a value that changes constantly, or a stored rank refreshed on every projection build. At 120,000 rank reads a second against boards of up to 10 million members, the exact answer is affordable for none of them. The product question hiding inside the technical one is whether the member needs the integer 41,208 or needs to know they are in the top 3%.
Decision. Exact ranks are served only for members inside the materialised head, capped per leaderboard scope. Every other member's position is answered from a value histogram with bucket counts, giving a percentile in constant time, with the neighbourhood of members immediately above and below served from a ranked-store slice. A leaderboard may declare dense or sparse rank for its head, and the response always states which kind of answer it is.
How it is realised on AWS. The projection builder emits, alongside the ranked view, a histogram of value-bucket counts for the scope. A rank query resolves the member's value from the aggregate store, locates its bucket, and derives a percentile from the cumulative counts; the member's own value is always exact. Histogram bucket width is a per-board setting, because it sets the accuracy of every non-head rank on that board.
| Option | Verdict | Reasoning |
|---|---|---|
| Exact rank in the materialised head, histogram percentile beyond it | Chosen | Constant-time, affordable at 120,000 reads/second, and honest about which kind of answer it gave |
| Indexed count-less-than over the aggregate store | Rejected | Exact and fresh for every member; the index is maintained on the one column that changes on every write |
| Per-member stored rank refreshed on each projection build | Rejected | Exact and a simple read; materialises 10 million rows per scope per version, for a number most members never look at |
| Exact rank for everyone, computed on demand | Right elsewhere | Right for boards of thousands, such as a league division or a friend leaderboard — which is why scoped boards exist |
What it buys
- Rank-of-member cost is independent of leaderboard depth, which is what makes a 10-million-member board affordable at all
- The product gains a cheaper and often better answer — "top 3%" ages more gracefully than "#41,208"
- Scoped boards (league, cohort, friends) become the designed way to give a member an exact rank that means something
What it costs
- A percentile where a member expected an integer is a product regression even when it is architecturally correct, and needs designed presentation
- Histogram bucket width is a single global-feeling parameter that silently sets accuracy for every non-head member
- Two code paths for one question, with a visible discontinuity at the edge of the materialised head
Choose differently when. If every leaderboard were small — tens of thousands at most, because the product only ever ranks within a division — the histogram would be unnecessary complexity and exact ranks for all would be correct. The decision is justified by at least one board with 10 million ranked members.
Why it holds up over time. The arithmetic here is structural: exact global rank for an arbitrary member is a counting problem over a continuously changing set, and no storage technology makes that cheap at this read rate. What may change is the product's appetite — if "top 3%" becomes the normal way products express standing, the exact head becomes a smaller and smaller special case, which moves in this decision's favour.
Lesson. Before paying for an exact answer, ask what the reader does with the precision. A cheaper answer that is also more meaningful is the best outcome available in a ranking system.
ADR-09 · Windows are pre-aggregated into the finest declared bucket and composed by addition
Status: Accepted · Shown on views: 10, 12, 14
A tenant wants "this week" and "last 30 days" over the same counter — are those two aggregates or one?
Context. Tenants declare several windows over the same counter: all-time, today, this week, this month, this season, last 24 hours. Maintaining each declared window as its own running aggregate makes every read a single lookup and every new window definition a backfill, which turns a product configuration change into a platform operation. Composing windows from a fine pre-aggregated bucket makes a new window free to declare and makes the read a fan-in — 168 hourly buckets for a 7-day window, 720 for 30 days. Recomputing from the log on demand is correct, trivially flexible and far too slow for a screen.
Decision. Events are pre-aggregated into the finest bucket any declared window requires — one hour — and coarser windows are composed by bucket addition rather than by re-reading the log. Day-level roll-ups are maintained for windows longer than seven days so that the fan-in stays bounded, and rolling-window state expires automatically once it can contribute to no declared window.
How it is realised on AWS. The bucket row key is (counter, member, window bucket, shard), so composing a window is a range scan over contiguous rows in Bigtable. The projection builder reads the range for the board's declared window and emits the ranked view. Hourly buckets are retained for 400 days to serve year-over-year comparisons; day roll-ups are derived from them and are themselves a projection.
| Option | Verdict | Reasoning |
|---|---|---|
| Finest-bucket pre-aggregation with day roll-ups, composed by addition | Chosen | A new window is a configuration change, and the fan-in stays bounded by the roll-up tier |
| One running aggregate per declared window | Rejected | Cheapest read in the design; every new window becomes a backfill and a platform ticket |
| Pure hourly buckets with no roll-ups | Rejected | Simplest model; a 30-day window is a 720-row fan-in on every projection build |
| Recompute windows from the log on demand | Right elsewhere | Right for an analytical surface with minutes of latency budget; impossible inside a 40 ms read |
What it buys
- A tenant can declare a new window and have it serve historical data immediately, with no backfill
- Window expiry is mechanical: a bucket that no declared window can reach is deletable without a bespoke cleanup job
- The same buckets serve every window, so storage is proportional to active keys rather than to the number of windows
What it costs
- Composition is a fan-in, so the read cost of a long window is a function of the roll-up tiers, which is one more thing to size
- Roll-ups are projections of projections, so a correctness bug in the roll-up step is one level removed from the log
- Daylight-saving transitions have to be resolved consistently between the bucket boundary and the tenant's declared timezone
Choose differently when. If tenants only ever needed one window per counter and never added one, a single running aggregate would be the right design and the whole composition machinery unnecessary. It is justified by five or more declared windows per counter and a product appetite for adding them.
Why it holds up over time. Pre-aggregating to a fine grain and composing upwards is the same decision OLAP systems settled on decades ago, for the same reason: the set of questions is not known in advance. Nothing about faster storage changes the asymmetry between a backfill and an addition.
Lesson. Pre-aggregate at the grain of the question you cannot anticipate, not at the grain of the question you have. Addition is cheap; backfill is a project.
Freshness and honesty
The decisions about where freshness is spent, and how the platform admits what it does not know yet.
ADR-10 · Read-your-writes is a query-time overlay of the member's own recent deltas
Status: Accepted · Shown on views: 04, 08, 15
How does the member who just earned the points see them, when nobody else needs to?
Context. ADR-01 and ADR-02 together mean counter values are eventually consistent. For almost every reader that is fine: a top-10 that is five seconds behind is indistinguishable from a live one. For exactly one reader it is not fine, and that reader is the member who just did the thing. A list that loads in 40 milliseconds and does not contain their last action reads as broken, while a list that is twenty seconds stale and does contain it reads as correct. The naive fixes are both bad: making every read strongly consistent prices the read path out, and making the write synchronous reintroduces the coupling the architecture exists to remove.
Decision. The Query API applies an overlay at read time: the acting member's own accepted-but-not-yet-projected deltas are fetched and applied on top of the stale projection, so their own value, rank band and window totals reflect their action within one second. Everyone else's view of the ranking is served from the projection unchanged, at the declared freshness budget.
How it is realised on AWS. Admission writes a short-lived per-member delta record alongside the log publish, keyed by (member, counter) with a TTL slightly longer than the worst-case projection lag. The Query API reads it only when the request's authenticated subject is the member being asked about, adds the delta to the projected value, and recomputes the member's rank band from the histogram. A session carries its pinned projection version, so the overlay never double-counts a delta the projection has already absorbed.
| Option | Verdict | Reasoning |
|---|---|---|
| Query-time overlay of the member's own recent deltas | Chosen | Spends freshness only where it is perceived, and bounds the cost to one extra read on 30% of requests |
| Strongly consistent reads for everything | Rejected | No special case to build; the read path then inherits the latency and availability of the counter store |
| Session token carrying last-write version, with the read waiting for the projection to catch up | Rejected | Elegant and needs no second store; converts a freshness problem into a latency problem the member feels directly |
| Route the member's reads to the aggregation shard that owns their key | Right elsewhere | Right where the pipeline is a stateful service with addressable partitions; Dataflow's model does not expose that, and it couples reads to processing |
What it buys
- The one thing a member notices is correct, which is what keeps the rest of the eventual-consistency design acceptable to them
- The cost is bounded and predictable: one extra keyed read, on the requests where the subject is the member
- Global freshness can be a budget rather than a guarantee, which is what makes a 4× burst survivable
What it costs
- Two code paths produce the same number, and if they ever disagree by more than the member's own delta the platform has two answers and no tie-break
- A second short-lived store on the read path, with a TTL that must track the worst-case projection lag or the overlay will double-count or go blind
- Version pinning per session means a long-lived session can hold a projection version the platform would rather retire
Choose differently when. If the product never showed a member their own position — a pure spectator leaderboard, or a board of teams rather than individuals — the overlay would be unnecessary and the whole read path would be one lookup. It is justified by the fact that the member is almost always in the list they are looking at.
Why it holds up over time. The underlying observation is about perception, not technology: consistency matters where a reader has independent knowledge of the truth, and nowhere else. That is why this decision will keep paying off as the platform's aggregation gets faster — the overlay gets cheaper and smaller, and never becomes wrong.
Lesson. Spend consistency where someone can tell. For everyone else, be stale and say so.
ADR-11 · Every response declares its projection version and as-of time, and a session is pinned to one version
Status: Accepted · Shown on views: 04, 13, 15
What does the platform owe a reader it cannot give a live answer to?
Context. An eventually consistent ranking served without qualification is a system that lies by omission: the response looks exactly like a live one, so the product renders it as live, and the first time a member compares two screens the platform is wrong rather than late. The second problem is subtler and generates more complaints: two reads in one session can land on different replicas at different projection versions, so a member watches their rank move up and down with no action taken. Both are solved by making the version explicit, and neither is solved by making the pipeline faster.
Decision. Every read response carries the projection version it was served from and the as-of timestamp at which that projection is accurate, and a reader's session is pinned to a single version so a rank cannot oscillate within it. When freshness is outside its budget — during a declared burst, or a pipeline stall — the response says so rather than silently ageing.
How it is realised on AWS. The projection builder publishes a monotonic version per (board, scope) and warms the Memorystore top-N for that version before it becomes readable. A conditional read returns nothing when the client's version is current, so a polling client refreshing a visible leaderboard transfers no payload. The staleness tagger is a single component, so the contract cannot be forgotten by one endpoint.
| Option | Verdict | Reasoning |
|---|---|---|
| Version and as-of on every response, version pinned per session | Chosen | Removes both the implied-liveness lie and the in-session oscillation, and makes conditional reads free |
| Serve the freshest available with no version, document eventual consistency | Rejected | Simplest API; the product cannot distinguish stale from current, so it renders stale as current |
| A staleness header, but no session pinning | Rejected | Fixes the honesty problem and leaves the oscillation, which is the one members actually report |
| Hide staleness behind a hard freshness SLO and page on breach | Right elsewhere | Right where the budget can be met essentially always; here a declared 4× burst widens staleness by design |
What it buys
- The product can render "as of" and be honest at no engineering cost beyond displaying a field it already receives
- Rank oscillation inside a session disappears, which removes a complaint class that no latency metric would have surfaced
- Conditional reads make a polling leaderboard screen nearly free, which matters at 400,000 reads a second
What it costs
- A product that chooses not to render the as-of converts a declared degradation back into a wrong number, and the platform cannot enforce that
- Pinned versions must be retained until sessions drain, which delays retirement of superseded projections
- Publishing a version is now a two-step operation — warm then expose — and getting the order wrong turns a publish into a latency cliff
Choose differently when. If the platform could guarantee sub-second global freshness at this read rate and cost, staleness would not need declaring, because there would be none worth declaring. It is justified by a design that deliberately trades global freshness for burst tolerance.
Why it holds up over time. Explicitly dated data outlives every caching and replication technology, because the honesty problem is in the contract rather than in the implementation. As more product surfaces blend data of different freshness — and AI-generated summaries over live data make this sharply worse — a response that states when it was true becomes more valuable, not less.
Lesson. A stale answer that states its age is correct. One that looks live is wrong, and no amount of engineering on the pipeline fixes the contract.
Results, retraction and integrity
The decisions that make a published standing trustworthy and a correction auditable.
ADR-12 · Retraction is a compensating event, never a mutation of a stored counter
Status: Accepted · Shown on views: 06, 12, 21
A post is deleted, an account is banned, a campaign is proved fraudulent — how do those contributions leave the counters?
Context. Removal is a routine requirement, not an exception: posts are deleted, accounts are banned, participants are disqualified, fraud campaigns are discovered weeks later, and members exercise deletion rights. The direct implementation is to decrement the counter, or to recompute it from a filtered set and overwrite it. Both destroy the ability to answer the question that follows every removal — what exactly was taken out, by whom, on what authority — and both make the number unverifiable from the log, which breaks the guarantee ADR-01 exists to provide.
Decision. A retraction is expressed as a compensating event referencing the original idempotency key, published to the log like any other event. Bulk retraction by named cause is one auditable operation across every counter and window the affected events touched. Counters are never mutated in place, and any ranking whose membership changes is recomputed rather than adjusted, so ranks stay contiguous.
How it is realised on AWS. The Operations API writes a retraction record to Spanner with cause, requester, authority and the event set, then publishes compensating deltas to the log. The aggregation pipeline applies them like any other event, which keeps the log and the counters consistent by construction. Affected ranked views are rebuilt at a new version rather than patched, and the retraction record is retained for 5 years.
| Option | Verdict | Reasoning |
|---|---|---|
| Compensating events through the log, with rankings recomputed | Chosen | Keeps the log authoritative, makes every removal auditable, and keeps ranks contiguous |
| Decrement the counter directly | Rejected | Fastest and simplest; the counter no longer matches a replay of the log, which breaks the rebuild guarantee |
| Filter at read time against a retracted-event set | Rejected | Preserves the log untouched and puts a growing exclusion set on the hot path of every read |
| Delete the original events from the log | Right elsewhere | Required in some jurisdictions for personal data; here it would destroy the evidence that justifies the removal, so it is confined to the deletion path with its own record |
What it buys
- Every removal has an actor, an authority and a record, which is what makes a disputed standing answerable rather than arguable
- The log stays authoritative: a rebuild after a retraction produces the corrected value with no special case
- Bulk retraction by cause lets a fraud campaign be removed as one operation instead of as thousands of edits
What it costs
- Retraction is the most powerful operation in the platform, so its authorisation and its reversibility matter more than its speed
- A ranking recompute after a bulk retraction is a full projection build, which is expensive and visible as a version jump
- A sketch-backed unique counter cannot be compensated arithmetically and must have its window rebuilt from the log
Choose differently when. If removal were never required — a counter of anonymous, unattributable, never-disputed events — compensation would be unnecessary complexity. It is justified by deletion rights, bans, disqualifications and fraud, all of which arrive within the first year.
Why it holds up over time. Append-only with compensation is the accounting profession's answer to the same problem, and it has lasted centuries for a reason: a correction that leaves no trace is indistinguishable from a cover-up. Regulatory pressure on explainability and auditability has moved consistently in this direction and shows no sign of reversing.
Lesson. Record corrections as events. A number you can only correct by overwriting is a number you cannot defend.
ADR-13 · Closure is explicit, and a retraction after it produces a correction record rather than a rewrite
Status: Accepted · Shown on views: 06, 11, 15
The season closed on Sunday and a cheater was proved on Wednesday — what happens to the published standings?
Context. A standing that can be silently rewritten was never a result. The moment a reward is granted against a published standing, the platform has two obligations that conflict: the standing must reflect what the members were told, and it must reflect what is true. Rewriting it satisfies the second and quietly breaks the first, including for the member who was paid. Refusing all post-closure correction satisfies the first and leaves a known-wrong published result standing. The resolution is not technical but it has to be expressed in the data model.
Decision. Each window or season has an explicit closure point after which its standings become immutable, snapshotted independently of the live projections. A retraction affecting a closed season produces a documented correction record published alongside the original standing — never a silent edit — and the rollover itself happens without a write freeze, so the new season counts while the previous one is being closed.
How it is realised on AWS. A Workflows pipeline driven by Cloud Scheduler seals each season's standings into Spanner at its declared closure time, in the tenant's declared timezone, while new events continue to land in the next season's buckets. Closed standings are read by season id and never by projection version, which makes them immune to any later projection rebuild. The reward service reads only closed standings.
| Option | Verdict | Reasoning |
|---|---|---|
| Explicit closure, immutable snapshot, correction records after it | Chosen | Keeps both obligations visible: what members were told, and what is now known |
| Retract everywhere and rebuild, including closed seasons | Rejected | Always truthful about the present; changes a published result a reward was paid against, with no trace |
| Refuse all post-closure correction, handle it out of band | Rejected | Simplest contract; leaves a known-wrong standing published and pushes the problem to support |
| No closure — standings always live | Right elsewhere | Fine where nothing is ever paid or awarded; impossible once a standing has consequences |
What it buys
- A reward can be granted against a standing that provably cannot change afterwards
- Post-closure corrections are visible to members and auditors instead of being invisible to both
- Rollover without a freeze means the busiest moment in a weekly league is not also a write outage
What it costs
- Two sources for "the standing": the live projection and the sealed snapshot, with a product obligation to show the right one
- A member whose rank improves because someone above them was removed is owed an explanation the platform does not generate today
- Closure needs a timezone-correct boundary per tenant, which is a persistent source of off-by-one-hour bugs
Choose differently when. If no leaderboard ever drove a reward, a badge or a public claim, closure could be a presentation detail and live standings would be enough. It is justified by the reward service being in the context diagram.
Why it holds up over time. The separation between a published result and a current estimate is older than computing and will outlast this platform. What changes over time is how quickly disputes arrive and how publicly they are argued, which increases the value of having the original standing and its correction both on the record.
Lesson. Decide when a number stops being an estimate and becomes a result. After that point, correct it in public or not at all.
ADR-14 · Suspect contributions are quarantined — held and reversible — rather than admitted or discarded
Status: Accepted · Shown on views: 06, 19, 20
The platform thinks this contribution is fake but is not certain. Count it, or drop it?
Context. A leaderboard with a visible reward is a target, and inflation signals are probabilistic: a value rising faster than any legitimate interaction path allows, many members rising in lockstep, a burst from one device or network. Both confident answers are wrong in different ways. Admitting the contribution and correcting later means the exploit is visibly winning for as long as detection takes, which is the window in which members lose faith. Discarding it means a false positive silently destroys legitimate member effort, and the member has no path to recover it — the worst possible outcome for a product whose whole proposition is that effort counts.
Decision. A contribution matching an inflation signal is quarantined: held outside the counters, retained with its original timestamps, and reversible. A cleared quarantine replays its held events in order so nothing is lost; a confirmed one becomes a retraction. A tenant may additionally freeze a leaderboard scope, which stops publication without stopping ingestion.
How it is realised on AWS. Admission writes held events to a quarantine table in Bigtable with the matched signal and state, and does not publish them to the log. The trust & safety console lists holds by signal and tenant; clearing a hold publishes the events with their original event times, so event-time windowing attributes them correctly even hours later. Freeze is a read-side flag on the board scope, so the Query API stops serving a ranking while aggregation continues.
| Option | Verdict | Reasoning |
|---|---|---|
| Reversible quarantine, with freeze as a separate read-side control | Chosen | A false positive costs a delay rather than a member's effort, and an exploit stops being visible immediately |
| Admit and correct later | Rejected | No false-positive cost; the exploit is publicly winning for the whole detection window |
| Discard on signal match | Rejected | Stops the exploit hardest and makes every false positive an unrecoverable loss of legitimate effort |
| Shadow-count suspect contributions in a parallel ranking | Right elsewhere | Useful for tuning detectors offline; as a live mechanism it doubles the ranking surface and the questions members can ask |
What it buys
- A false positive is recoverable, which is what makes it safe to tune the detector aggressively
- Freeze separates "stop showing this" from "stop counting this", so an investigation does not create a data gap
- Quarantine volume becomes a signal — a rising hold rate is either an attack or a detector regression, and both need attention
What it costs
- A held contribution is invisible to the member, who sees their action not counted and has no way to learn why
- Quarantine is a second write path with its own state, retention and authorisation model
- Clearing a hold replays events late, which pushes them against the lateness horizon and sometimes past it
Choose differently when. If detection were effectively certain — a hard server-side rate limit on an action that cannot legitimately exceed it — hold-and-review would be pure overhead and a straight refusal would be correct. The decision is justified by probabilistic signals with real false-positive rates.
Why it holds up over time. Hold-pending-review is the standard answer wherever automated judgement meets irreversible consequence, from payments to content moderation. As detection moves further towards statistical models, the false-positive rate becomes more central rather than less, which strengthens the case for reversibility.
Lesson. When a decision is probabilistic, choose the option whose mistakes are recoverable. Reversibility buys the right to be aggressive.
State placement and reach
The decisions about which store holds what, how far it reaches, and how a definition changes.
ADR-15 · Counters live in a wide-column store; definitions, seasons and standings live in a relational one
Status: Accepted · Shown on views: 11, 12, 16
Can one store hold both 3.5 billion counters and the handful of rows that must be exact?
Context. This platform holds two kinds of state with almost nothing in common. The counters are 3.5 billion keys across 120,000 leaderboard scopes, written 250,000 times a second, read in ordered ranges, and entirely rebuildable. The control plane is a few million rows — definitions, tenancy, quotas, seasons, closed standings, retraction records, audit — that must be strongly consistent, transactional and never lost. A single store that is excellent at both of those at this ratio does not exist, and choosing one means either paying strong-consistency cost on every counter write or holding a reward-bearing standing in an eventually consistent store.
Decision. Two stores, with a hard rule: the high-volume rebuildable counters never share a store with the exact transactional state. Bigtable holds bucket aggregates, ranked views, histograms, the shard map and the quarantine. Spanner holds counter and leaderboard definitions, tenancy and quotas, seasons, closed standings, retraction records and the audit log. Memorystore holds the serving cache and carries no RPO at all.
How it is realised on AWS. Bigtable row keys are designed for three access patterns at once — bucket upsert, window range scan, ranked slice — with a shard component so a hot key splits without a schema change. Spanner is multi-region with RPO 0, holds the only state that is backed up rather than rebuilt, and is read on the admission path through a per-instance cache for counter definitions. The split is visible in the storage-zone view, because it is the first thing a reviewer should be able to check.
| Option | Verdict | Reasoning |
|---|---|---|
| Wide-column for counters, relational for the control plane, cache with no RPO | Chosen | Each store is used for what it is good at, and the recovery story differs per zone by design |
| One relational store for everything | Rejected | One operational surface and one consistency model; pays transactional cost on 250,000 writes a second |
| One wide-column store for everything | Rejected | Cheap and uniform; puts a reward-bearing closed standing in an eventually consistent store |
| Document store for the control plane | Right elsewhere | Reasonable at smaller scale and fewer tenants; the cross-entity transactions around season closure are what argue for relational here |
What it buys
- The counter path is cheap and splittable, and the control-plane path is exact — neither compromises for the other
- Three recovery stories (backed up, rebuilt, lost) map one-to-one onto three storage zones, which makes DR explainable
- A reward-bearing standing cannot be written by a projection build, because it is not in a store the builder touches
What it costs
- Two data stores to operate, monitor, capacity-plan and reason about, with a 9-person team
- Counter definitions sit in Spanner but are read on the hot path of admission, so they are cached and therefore eventually consistent by up to the cache TTL
- Any invariant that spans both stores — a quarantine hold in Bigtable against a retraction record in Spanner — cannot be transactional and needs reconciliation
Choose differently when. At a tenth of this scale — 300 million counters rather than 3.5 billion — a single relational store would hold both comfortably and the operational simplicity would outweigh everything in this record. The split is justified by the ratio between the two workloads, not by either one alone.
Why it holds up over time. The specific products will change; the rule will not. Wherever a system holds one dataset that is enormous and rebuildable and another that is small and must be exact, separating them by store is the decision that keeps both affordable. Convergent databases narrow the gap in the middle of the range and do not change the arithmetic at a 1000:1 ratio.
Lesson. Let the recovery story choose the store. State you would rebuild and state you would restore do not belong in the same place.
ADR-16 · One write region, with a read standby whose cache is rebuilt rather than replicated
Status: Accepted · Shown on views: 15, 16, 21
Should counting continue in a second region when the first one is gone?
Context. Multi-region counting sounds like a resilience requirement and is really a merge-semantics question. Two regions accepting deltas for the same key must either coordinate on the hot path, which spends the cross-region latency on every write, or merge independently accumulated state, which is tractable for commutative counters and intractable for the ranked views and seasons built on top of them. Meanwhile the actual product requirement is modest: members must be able to read standings during a regional failure, and must not lose effort — which a client retry with a stable idempotency key already provides.
Decision. One write region. A second region runs a warm read path over an asynchronously replicated copy of the ranked store, serving at declared staleness, and rebuilds its serving cache on cutover rather than replicating it. The strongly consistent control-plane store is multi-region with RPO 0, because it holds the only state that cannot be rebuilt.
How it is realised on AWS. Cloud Run services are deployed in both regions with the standby scaled to minimum; Bigtable replicates asynchronously to the standby; Memorystore in the standby is cold and warms from the ranked store after cutover. Pub/Sub and Cloud Storage are already multi-region, so the log and archive survive a regional loss without intervention. RTO is 10 minutes for the read path and RPO 5 seconds for accepted events.
| Option | Verdict | Reasoning |
|---|---|---|
| Single write region, warm read standby, multi-region log and control plane | Chosen | No merge semantics to invent, and the product requirement — readable standings, no lost effort — is still met |
| Active-active counting with CRDT-style merge | Rejected | Elegant for commutative counters; the ranked views, seasons and closures built on them have no comparable merge rule |
| Active-active with cross-region coordination on write | Rejected | Globally consistent; spends cross-region latency on all 250,000 writes a second to buy a case that happens rarely |
| Single region with no standby | Right elsewhere | Right where an hour of unavailability is acceptable; a leaderboard that is blank during a regional event is a visible product failure |
What it buys
- No merge rule to invent for rankings, seasons or standings, which removes the hardest correctness problem in a multi-region design
- Members keep reading standings through a regional failure, at a staleness the response declares
- The log and the control plane — the only state that matters — are multi-region without any application logic
What it costs
- A regional failure stops counting, not just ranking, so clients must buffer and retry with stable idempotency keys — a contract on the product teams
- A cold cache on cutover means the first minutes after failover run at ranked-store latency rather than cache latency
- The standby is a permanent cost line for capacity that is idle almost always
Choose differently when. If the product were global and competitive in real time — one worldwide tournament where members in two regions must see the same live ranking within a second — a single write region would be untenable and the merge problem would have to be solved. The decision is justified by leaderboards that are mostly regional or cohort-scoped.
Why it holds up over time. The reason not to go active-active is not infrastructure cost; it is that derived, ordered state has no natural merge. That will remain true however cheap cross-region replication becomes, which is why this decision should survive even as the economics shift.
Lesson. Before paying for active-active, write down the merge rule for your derived state. If you cannot, you are buying availability you will not be able to use correctly.
ADR-17 · Counter definitions are versioned, and a breaking change creates a new version rather than mutating the old one
Status: Accepted · Shown on views: 05, 09, 12
A tenant wants to change a counter's tie-break rule. What happens to the events already counted under the old one?
Context. Definitions change: a tenant wants a different tie-break, tighter bounds, an extra window, a switch from exact to approximate. Applying the change in place is one update statement and silently reinterprets every event already counted — the stored value now means something the events never meant, and a rebuild from the log produces a different number than the live projection, which breaks the guarantee in ADR-01. The failure is invisible at the moment of the change and surfaces weeks later as a drift nobody can explain.
Decision. Counter and leaderboard definitions are versioned configuration. Every event records the counter version it was counted under. A change that would reinterpret already-counted events — type, bounds, tie-break, window set — is refused as an in-place edit and creates a new counter version instead. Compatible changes (a new window, a raised cap, a cost note) apply to the existing version. Validation happens before the change takes effect, not after.
How it is realised on AWS. Definitions live in Spanner with a version key; the Config API classifies each field by compatibility rule and refuses an incompatible edit with the reason and the new-version path. Admission reads the current version — cached per instance, effective within 60 seconds — and stamps it onto the event. A disposable sandbox tenant lets a team replay sample events against a candidate definition before promoting it.
| Option | Verdict | Reasoning |
|---|---|---|
| Versioned definitions, per-field compatibility rules, new version for breaking changes | Chosen | Keeps every stored value meaning what its events meant, and keeps rebuild deterministic |
| Mutable definitions | Rejected | One update and no migration; silently reinterprets history and breaks rebuild equivalence |
| Immutable definitions — every change is a new counter | Rejected | Strongest guarantee; makes adding a window a migration and guarantees tenants work around the platform |
| Mutable with a changelog and no enforcement | Right elsewhere | Adequate where no number is disputed and nothing is rebuilt; here it would make the drift traceable but not preventable |
What it buys
- A rebuild from the log always reproduces the live value, because both apply the version each event was counted under
- A tenant can evolve a counter without a platform ticket, and cannot evolve it into incoherence
- The refusal is a teaching moment: the error names which field is incompatible and why
What it costs
- Versions accumulate, and without a retirement path the definition store becomes an archaeology site
- A counter with several live versions needs a rule for which one a leaderboard reads, which is one more concept for tenants
- Per-field compatibility classification is platform logic that must be maintained as fields are added
Choose differently when. If definitions were owned centrally by the platform team and changed a handful of times a year under review, a changelog and discipline would be enough. Versioning is justified by self-service: 25 teams changing their own definitions without a platform engineer in the loop.
Why it holds up over time. Schema and semantic versioning of configuration is a settled practice wherever stored data outlives the code that wrote it, and counting data outlives it by years. The pressure towards self-service configuration has only increased, and self-service without enforcement is how the reinterpretation bug gets shipped.
Lesson. Configuration that changes the meaning of stored data is schema. Version it, classify its fields, and refuse the edit that rewrites the past.
Every package used, in one table
These terms are used precisely in this package. Several are used loosely in the wider literature on counting and ranking systems, and the difference matters when reading the decision records.
| Package | What it is | What it does here | Considered instead |
|---|---|---|---|
| Counting event | One signed delta against one counter for one member, with an event time and a client-supplied idempotency key. | The unit of admission, deduplication, retraction and audit. | A 'hit' or 'increment', which imply a delta of exactly one and no retraction path. |
| Counter | A named, versioned definition of something being counted, with a type, bounds and a set of windows. | What a tenant declares; what an event names; what a leaderboard ranks over. | A 'metric', which in common usage means a monitoring signal with no member dimension. |
| Bucket aggregate | The accumulated value for one (counter, member, window bucket, shard). | The finest-grained projection; every window is composed from these by addition. | A 'counter value', which hides the fact that the value is a composition rather than a stored total. |
| Projection | Any stored value derived from the log — a bucket aggregate, a ranked view, a histogram, a cache entry. | Droppable and rebuildable by definition. If a projection and the log disagree, the projection is wrong. | A 'materialised view', which is accurate but implies a database feature rather than an architectural rule. |
| Ranked view | A versioned, ordered list of members for one leaderboard scope, with a deterministic tie-break. | What a top-N read is served from, and what is rebuilt rather than patched after a retraction. | A 'leaderboard', which in this package means the declared definition, not the materialised list. |
| Scope | One slice of a leaderboard — global, a region, a cohort, a league division, a friend graph. | The unit of materialisation, freezing and dematerialisation. One member may appear in several. | A 'segment', which is used for marketing audiences elsewhere in the estate. |
| Projection version | A monotonic identifier for one build of a ranked view, published after its cache is warm. | What a response declares, what a session pins, and what a rollback points back to. | A 'snapshot', which is reserved here for the immutable sealed standing. |
| As-of | The timestamp at which a served projection is accurate. | The platform's statement of what it does not yet know; the field a product renders as "as of". | A 'last updated' time, which usually means when the record was written rather than when it was true. |
| Lateness horizon | The declared period after an event's event time during which it is still attributed to its own window. | Bounds the pipeline's window state and makes closure possible at a known time. | A 'grace period', which suggests a courtesy rather than a correctness boundary. |
| Late bucket | The visible accumulation of events that arrived beyond the lateness horizon. | Makes the cost of the horizon measurable instead of a silent policy. | A 'dead letter', which here means an unparseable event rather than a late one. |
| Compensating delta | An event that reverses a previously counted one, referencing its idempotency key. | The only mechanism by which a counter decreases on removal; what makes retraction auditable. | A 'decrement', which does not carry the cause, the authority or the link to the original event. |
| Closure | The declared point after which a window or season's standings are immutable. | The precondition for paying a reward against a standing; the boundary after which correction is published rather than applied. | A 'reset', which describes the counting side of a rollover and says nothing about the result. |
| Correction record | A published statement that a closed standing is now known to be wrong, with the prior and revised ranks. | What a post-closure retraction produces instead of an edit. | An 'amendment', used in some products for a member-initiated change rather than a platform one. |
| Quarantine | A hold on a suspect contribution: outside the counters, retained with original timestamps, reversible. | The answer to a probabilistic inflation signal; replayed on clearing, retracted on confirmation. | A 'block', which implies a permanent refusal and no path back for a false positive. |
| Shard map | The live record of how many physical shards a hot key's aggregate state is split across. | What makes hot-key splitting a data operation rather than a schema migration. | A 'partition map', which in this package refers to the storage engine's own tablet layout. |
| Own-write overlay | The read-time application of a member's own accepted-but-unprojected deltas to a stale projection. | How read-your-writes is delivered without making any read strongly consistent. | A 'write-through cache', which would make the projection itself fresh for everyone and cost far more. |