Reddit / Hacker News — a worked solution
Vote-driven ranking with hot/top/new at scale.
Try it yourself first.
You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.
Open the Reddit / Hacker News workspaceThe problem
Build a Reddit / Hacker News-shaped product: users submit links, comment in threaded discussions, and vote up/down on submissions and comments. Each subreddit (and the front page) renders hot / top / new / best listings ranked by a vote-driven formula. The product is read-tail-dominated: a typical session loads ~5 listings before producing a single vote. Voting is the highest-rate write workload; ranking is the largest async pipeline; comment-trees are the trickiest read path.
The reference architecture
What each component is for
- Client
Renders ranked listings (front page, /r/:sub/hot, /comments/:id), opens long-running websockets for live comments, and POSTs votes/submits with an Idempotency-Key header. Browsers receive HTML on cold loads + JSON on subsequent paginations; mobile is JSON-only.
Why it exists. Reddit-shape products are read-tail-dominated — a session typically loads 5+ listings before a single vote — so the client is where the SLO budget lives. We picked HTML+JSON over GraphQL precisely because edge-cacheable URLs (one URL = one ranked listing) are the load-bearing primitive that lets Fastly absorb 95% of read QPS.
When it fails. If the client retries non-idempotent votes after a 408/5xx — and they will, on every cellular handoff — vote-double-counting visibly skews ranking. Detection is the vote-audit job that compares raw vote-row count to score deltas; mitigation is the server-side idempotency dedup keyed by (account, target, vote_value) with a 24h TTL.
- CDNFastly
Fronts every listing fetch (/r/:sub/hot.json, /comments/:id.json) and signed media URLs. Caches listing JSON for 10s with surrogate-key purge keyed on (sub, sort), so the ranker can soft-purge a listing in <1s after each ZADD. Signed URLs gate any non-public sub.
Why it exists. At ~32K listing QPS peak with 40% landing on the front page (one logical key), origin would melt without short-TTL edge caching. We considered cache-asiding only at the gateway and skipping the CDN, but origin's listing-read tier still gets crushed by the front-page hot key — only the edge gives geographic spread that absorbs viral traffic. Fastly's surrogate-key model gives us the freshness invariant a time-based TTL alone cannot.
When it fails. Edge cache poisoning: a stale 'hot' object pinned at one POP for 60s while the actual hot listing has rotated away. Detection is the synthetic prober (hosted on the monitor) that asserts Age on /hot is <30s; mitigation is the surrogate-key purge fired on every ranker emit, plus a hard 30s ceiling on edge-ttl. Loss of CDN entirely would push origin to ~10× current load — the runbook is to flip to a secondary CDN provider via the LB's backup origin pool, accepting a 90s propagation gap.
- Load Balancer · Anycast EdgeAWS NLB + BGP anycast
Terminates anycast traffic from the CDN miss path, runs L7 health checks against the gateway tier (200ms timeout, 2-of-3 threshold), and round-robins requests across healthy gateway pods cross-AZ. On region failover, anycast withdraws the unhealthy region's prefix and traffic re-routes within ~30s — strictly faster than DNS-based failover by avoiding TTL propagation.
Why it exists. Region failover via DNS depends on resolver TTL respect — production resolvers cache aggressively past TTL, so the 'half users on dead region' pattern (failure mode 19 in the SRE Workbook) hits exactly here. We considered a 30s DNS TTL on regional records but rejected it because (a) DNS amplification spikes resolver load, and (b) anycast withdrawal is sub-30s without any client-side caching games. ECMP at the edge spreads traffic uniformly; gateway health-checks remove unhealthy pods within one check interval.
When it fails. Gray failure: health checks pass but real traffic fails (gateway returns 200 on /health but 503 on /vote). The LB happily routes to broken pods. Detection: the monitor's synthetic prober exercises a full /vote round-trip every 10s and alerts on >1% failure rate. Mitigation: deep health check that touches downstreams (rate-limit-redis ping + JWKS fetch) so /health fails when the gateway is not actually serviceable.
- API GatewayEnvoy + ratelimit sidecar
Terminates TLS, runs WAF rules, enforces per-IP/per-account token-bucket rate limits (60 votes/min/account, 10 submits/hour/account, 600 listing fetches/min/IP) by INCR-ing buckets in rate-limit-redis, validates session cookies or OAuth bearer (JWT verified locally against a JWKS cache, refreshed every 5min), then dispatches: /hot|/top|/new → listing-read; /comments/:id → comment-tree; POST /vote|/submit → vote-svc. Strips Idempotency-Key into a request attribute.
Why it exists. Submission-spam waves and vote-storms are abuse problems, not capacity problems — they need to be shed AT THE EDGE before they consume any backend capacity. The gateway is also where karma-farm shadow-banning shunts traffic to a low-priority queue (the ranker still sees the votes, but tags them suspect). Doing rate-limiting in each backend service would mean every backend grows its own Redis client and rules diverge.
When it fails. If rate-limit-redis is partitioned, gateway must fail-open or fail-closed. Fail-closed = real users see 429s during a Redis blip; fail-open = a spam wave punches through. We fail-OPEN with a circuit-breaker that re-engages on Redis recovery, AND the fraud-scrubber re-evaluates the burst window after the fact — the spam still gets shadow-removed, just not blocked at the door. Detection is the ratelimit-redis p99 + fail-open-rate gauge. Second failure: JWKS-cache staleness — if the IdP rotates and JWKS cache misses propagate, every authed request 401s for up to 5min. Detection: 401 rate spike post-JWKS-rotation; mitigation: pre-warm JWKS cache on rotation event (the IdP webhooks the gateway on every rotation).
- Cache · Rate Limit BucketsRedis 7 (single shard + 1 follower)
Holds the per-(account, ip) token-bucket state for the gateway's rate-limiter. Lookup key: rl:{account|ip}:{endpoint}. Each gateway request issues a LUA-scripted INCR-and-check + EXPIRE, returns the remaining tokens; the gateway decides allow/deny. Hot keys: any IP behind a NAT (corporate / school VPNs) drives a single key to thousands of QPS.
Why it exists. Storing rate-limit state in the gateway's process memory is cheap until you scale to 12 gateway pods — a per-pod bucket lets a single bot burst ×12 by hitting different pods. We considered hashing client IP to a sticky pod (so each pod owns a deterministic IP slice) and rejected it because asymmetric load and pod failures would drop limit state. Centralized Redis is the correct trade — sub-ms p99 limit checks, shared state across all gateway pods.
When it fails. If rate-limit-redis is unreachable, gateway fails OPEN (per gw.failureMode), so a partition is non-fatal. The harder failure is bucket-key contention on a NAT'd IP — a corporate proxy with 50K users behind one IP saturates one bucket key. Detection: per-key INCR rate >10K/s alarms on the monitor. Mitigation: hash subkey on user_agent for IPs flagged as NAT'd; the fraud-scrubber labels them.
- Service · Listing ReadGo (Baseplate.go)
Resolves /r/:sub/hot.json by ZRANGEing the precomputed sorted set in listing-cache, multi_getting per-thing data (title, score, thumbnail) from render-cache (with a thing-db replica fallback for cold rows), and returning the page. Front-page (/r/all, /r/popular) is a fan-in over O(50K) sub listings precomputed by the ranker into one global candidate ZSET.
Why it exists. Computing hot at read time would mean SELECT…ORDER BY hot(ups,downs,date) LIMIT 25 across millions of posts on every request — at 32K QPS peak, that is a query-planner death spiral. Precomputed Redis ZSETs reduce read-time work to O(log N) per page. We rejected Postgres materialized views because cross-sub fan-in for /r/all needs to merge 50K row groups in <50ms, and Postgres simply cannot do that. We rejected client-side ranking because half our DAU is on phones that cannot afford the CPU.
When it fails. If listing-cache misses (Redis evicts under memory pressure or a shard fails over), every concurrent request stampedes the ranker queue and thing-db replicas. Detection is listing_miss_rate >5% paged within 60s by the monitor; mitigation is single-flight (request-coalesced ranker recompute, others block on a promise) plus a 30s stale-while-revalidate served from the prior cached page. The edges to listing-cache and listing-read use circuit-breaker (not exp-backoff) for exactly this reason — at 32K QPS, exp-backoff produces a 27× retry amplification under partial Redis failure. Pi Day 2023 was this shape: bring the cache back BEFORE the front door reopens.
- Service · Comment TreeGo (Baseplate.go)
Fetches /comments/:post_id by reading the precomputed preorder traversal of the comment tree from thing-db (a flat list of comment_ids with depth/parent), multi_getting per-comment HTML chunks from render-cache, and decorating each comment with the requesting user's vote state (read from vote-store via Bloom-filter negative-lookup). Trees deeper than 250 children at any level surface a MoreChildren stub; the client paginates into them via /morechildren.
Why it exists. The largest threads (e.g. /r/AskReddit megathreads) reach 50K+ comments. Building this tree from a recursive Postgres CTE on every fetch would push p99 to seconds and DOS the DB. Pre-rendered HTML fragments + a flat preorder array means a thread fetch is one tree-skeleton query plus one memcached multi_get plus one Cassandra Bloom probe per displayed comment.
When it fails. When the #1-on-front-page thread's render-cache expires, every concurrent fetcher races to recompute. With 10K concurrent viewers, that is 10K parallel CPU-bound renders. Detection: render_p99 >5s paged; mitigation is single-flight per (post_id, subtree_root) plus chunked subtree caching (so a vote on one child does not bust the whole tree). Discord's Scylla migration retro covers this exact fanout.
- Service · Vote / SubmitGo (Baseplate.go)
Handles POST /vote and POST /submit. For votes: validates Idempotency-Key, runs Cassandra LWT (Paxos compare-and-set) on (user_id, target_id, idempotency_key), atomically writes (a) the vote row + (b) an outbox event into the same Cassandra LOGGED batch keyed on user_id, returns 200 with the new score read from the fast post-counter row in thing-db. Best-effort first-attempt Kafka publish runs in-line; the outbox-tailer sidecar handles catch-up. For submits: post_id = deterministic hash(idempotency_key, author_id) so retries hit the same row, INSERT into thing-db with ON CONFLICT DO NOTHING, then emit a 'new submission' outbox event.
Why it exists. Votes MUST be idempotent — clients retry on every network blip and a duplicated vote skews ranking visibly. We picked Cassandra LWT over Postgres unique-constraints because the vote-history fan-out (a viral post takes 1K votes/sec on a single key) saturates a Postgres shard. We rejected Redis-only counters because we lose the per-(user, target) audit row that the fraud-scrubber depends on. We chose the outbox pattern over 2PC across Cassandra + Kafka because 2PC over Kafka is operationally a tar pit.
When it fails. Three signals on this tier, all on the monitor: (1) LWT_failures rate spike — Cassandra cannot form a write quorum (two AZs partitioned), 503 to client, retry with same Idempotency-Key. (2) Outbox-row-age — now() - min(outbox.created_at WHERE published_at IS NULL) >30s pages on a dead outbox-tailer (the silent-killer; both db_score and ranker_score stop moving together so the score-divergence reconciliation cannot detect it). (3) Score-divergence — the reconciliation job alerts on |db_score - ranker_score| > 5%, and on alert we drain the outbox before re-rank.
- Worker · RankerGo + Kafka consumer
Consumes vote events off Kafka in a per-(post_id) keyed partition. Recomputes hot = log10(max(|score|,1)) + sign(score)·t_submit/45000 for posts, where t_submit is the submission time in epoch-seconds — newer submissions earn a larger term, so recency is rewarded; Wilson-score lower-bound at 95% confidence for comments (the 'best' sort); HN-style (P-1)/(T+2)^1.8 wired in for /r/all gravity decay. Reads a 5s-TTL per-post tally cache (in-process) before fanning out to the 16-way counter-shard-split tally read on viral posts. ZADDs the new score into listing-cache for every (sub, sort) the post is in, plus the global candidate set for /r/all. Cadence: per-vote on viral posts (top 1K), per-batch every 5s on cold posts to amortize Redis ops.
Why it exists. Read-time ranking does not scale beyond ~500 listing QPS — Reddit pre-2010 melted exactly here. Async ranker keeps the vote-write path at 120ms p99 regardless of how many listings the post participates in. The fan-out (one vote → up to 4 ZADDs) happens off the user's request. We considered Kafka Streams stateful processors but rejected them — the score function is dumb arithmetic; the value is in batched coalescing, not in stream-SQL.
When it fails. Consumer lag is the silent killer — if the ranker falls behind 5+ minutes, the front page is yesterday and users disengage with no error firing. Detection: the monitor scrapes the ranker's freshness heartbeat (now() - max(listing.generated_at) >5min) and pages on staleness, not on errors — symptom-based alerting per the SRE Workbook. Mitigation: load-shed low-value events (anonymous-view signals before votes), autoscale on lag not CPU, and a kill-switch to drop fraud-suspect votes from the rank input when scrubber is also behind.
- Worker · Fraud ScrubberGo + Kafka consumer + graph clustering
Tails the same vote-stream on its own consumer group. Maintains a rolling fraud-score per (account, IP, device, vote-time-correlation) using sliding-window graph clustering (1h windows, edges weighted by co-vote count). On suspicion, writes a 'scrub' record into vote-store and emits a 'recompute(post_id)' control event to the ranker — the surfaced votes remain visible to the offending voter (shadow-ban semantics, the load-bearing trick) but are excluded from rank input. Also shadow-marks accounts in thing-db so subsequent submits go to a low-priority ranker queue.
Why it exists. Vote manipulation is existential for a ranked product — without an offline scrubber, a $100/day karma farm pollutes the front page and the product dies. We chose async over inline because the fraud signals (graph clustering on co-voting, cross-account IP overlap, temporal pattern matching) need a window of votes to evaluate; inline would push vote-ack p99 past the 120ms budget. We chose shadow-removal (not blocking) because if attackers see their votes get blocked, they iterate the attack — invisibility is the moat.
When it fails. Scrubber lag pollutes the product visibly. SLO: scrubber lag <1h. If breached, runbook: (a) drop low-confidence votes from rank input until lag drains, (b) re-rank affected posts post-scrub. Detection is consumer-lag-on-fraud-topic page from the monitor; mitigation is the inline lightweight check at vote-time (gateway-level, IP velocity only) as a load-bearing fallback.
- Worker · Outbox TailerGo sidecar (1-per vote-svc replica)
A sidecar (one per vote-svc replica) that reads vote-store rows where published_at IS NULL AND created_at > now() - 24h, publishes each event to vote-stream (Kafka), and updates published_at = now() on success. Runs every 100ms. Vote-svc still publishes synchronously on the happy path (e14); the tailer is the catch-up path only fires when the synchronous publish failed (Kafka blip, partition leader election, broker GC pause).
Why it exists. The synchronous publish from vote-svc is best-effort — if Kafka is briefly unavailable, the vote is durable in vote-store but absent from the ranker's input. Without the tailer, votes stuck in the outbox stay invisible to the ranker forever, the front page silently freezes for those posts. We considered a single global tailer (one process tailing all outbox rows) but rejected it as a SPOF; a sidecar-per-replica gives natural sharding by user_id partition assignment.
When it fails. If a tailer dies and no replacement spawns, votes pile up in the outbox; rank visibly freezes for those users' votes (depending on partition coverage). Detection: the monitor scrapes now() - min(outbox.created_at WHERE published_at IS NULL) >30s — the canonical's load-bearing outbox-lag signal. Mitigation: monitor's heartbeat-absence alarm plus k8s liveness probe on the sidecar. Re-publish is at-least-once; the ranker dedupes via idempotency_key.
- Cache · Listing ZSETsRedis 7 (cluster mode)
Holds the precomputed hot/top/new listings as Redis sorted sets keyed listing:{sub}:{sort}. Each ZSET holds 1000 thing_ids ranked by score. Reads are ZRANGE; writes are ZADD/ZREMRANGEBYSCORE on every vote (top-of-listing posts) or ranker-batched (cold posts). The front-page ZSET is replicated across all 16 shards as listing:_all:hot:0..15 and reads round-robin across them — this is the load-bearing mitigation for the 13K front-page-read QPS hot key.
Why it exists. Listings ARE the read product — at 32K QPS peak with 40% on front-page (one logical key), Postgres replicas would saturate before reaching 5K QPS. Redis ZSET gives O(log N) updates and O(K) range reads in-memory, which is what an 80ms p99 cached budget demands. We rejected DynamoDB DAX because the per-key write rate (1K vote-ZADDs/sec on the #1 post) exceeds DAX node throughput. We rejected an in-process LRU on each listing-read replica because cross-replica score divergence is a worse failure than network hops.
When it fails. Cold restart is the page that shaped the May 2023 Reddit outage. If the cluster restarts cold, hit-rate drops below 50% and Cassandra-thing-db gets crushed. Detection: monitor scrapes Redis hit-rate; <90% pages. Mitigation: refuse cold-start traffic — admission control (at the LB) caps origin QPS to 10% of normal until hit-rate recovers, ranker pre-populates from a snapshot, the front door re-opens only after hit-rate >90%. Cross-region cold-rebuild RTO for the full ZSET set (~50K subs × 1K thing_ids) is ~10 min — this is the real RTO, not the 30s edge cutover. Hot-shard meltdown on the front-page key is mitigated by the 16-way fan-out described in keyChoices.
- Cache · Render FragmentsMemcached
Holds rendered HTML fragments — per-comment HTML, per-listing-card HTML, and the user-flag overlay for vote-arrow state. Reads are mget over hundreds of keys per page; writes are issued by listing-read and comment-tree services on cold renders, and invalidated by the ranker on score changes greater than the 'visible delta' threshold (avoids invalidating on every fuzz).
Why it exists. Pre-rendering once and serving N times is what makes /r/AskReddit megathreads feasible. Rendering 50K comments per request would saturate the comment-tree CPU at <50 QPS. We chose Memcached over Redis here because Memcached's slab allocator handles variable-sized HTML fragments better than Redis's hash table, and we already use Redis for the sorted-set workload — keeping caches workload-specialized lets us tune memory and eviction independently.
When it fails. Two pages on this tier. (1) Connection-pool exhaustion: a megathread fanout drives every comment-tree replica through its per-replica Memcached connection limit (~256 conns/replica) before any infrastructure error fires. Detection: mget p99 spike with no Memcached-side error. Mitigation: connection-pool autosize on QPS, single-flight per-key. (2) Topology change: a node added/removed effectively invalidates ~1/N keys, triggering a render storm on whichever threads land on the new shard. Detection: render p99 spike. Mitigation: ketama hashing + rolling-add-then-fill (warm the new node from thing-db before exposing it to reads).
- Stream · Vote EventsKafka 3.7
Single source of truth for every vote and submission event. Producer: vote-svc (synchronous best-effort) + outbox-tailer (catch-up). Consumers: ranker-worker (per-post partition), fraud-worker (separate consumer group), and a search-indexer consumer (out-of-scope for the main flow). Per-message body: (user_id, target_id, target_kind, dir, ts, idempotency_key, fraud_signals, ranker_version_pin).
Why it exists. The outbox pattern (write the vote row + the event in the same Cassandra LOGGED batch) decouples vote-write durability from ranker availability — the ranker can be down for 30 minutes and votes still durably commit; the rank just lags. We rejected synchronous gRPC fan-out from vote-svc to ranker because every additional sync hop multiplies vote-ack p99, and ranker recomputation on a viral post is itself bursty. We rejected RabbitMQ (Reddit's historical choice) because partition-keyed ordering is load-bearing and Kafka's per-partition leader election handles vote-storms more gracefully than RabbitMQ mirrored queues.
When it fails. Two leading indicators on the monitor. (1) under-replicated-partitions >0 — with min.insync.replicas=2 of RF=3, one broker loss puts every partition at the edge of write availability, and ISR shrinkage is 30+ seconds ahead of the producer 503 symptom. (2) Consumer lag on votes.events >30s p99 — the dominant pipeline failure shape. The poison-pill scenario (one bad message blocks one partition) is mitigated by max-retries-then-DLQ; DLQ runaway is bounded by per-partition consumer isolation so 1/256 of votes never blocks the rest. Loss of in-sync replicas below min.insync.replicas → producer 503s (we deliberately fail closed; back-pressure to vote-svc is preferable to acking a non-durable vote).
- Monitoring · Freshness + SLOPrometheus + Alertmanager + synthetic prober
Scrapes per-service metrics every 15s — listing-cache hit-rate, ranker freshness heartbeat, vote-svc outbox-lag, vote-store coordinator queue depth, vote-stream under-replicated-partitions, render-cache mget p99, gateway 401-rate (JWKS), rate-limit-redis INCR rate per key. Hosts the synthetic prober (Go binary, every 10s) that exercises /hot, /comments, and /vote round-trips end-to-end from each region. Pages on SLO breach via PagerDuty integration.
Why it exists. Gray-failure detection — silent stale ranker, dead outbox-tailer, JWKS-rotation cascade, edge cache poisoning — has no error signal in any single service's logs. Symptom-based alerting requires a CENTRAL place to express 'now() - max(listing.generated_at) >5min' across all ranker replicas. We considered piping these signals through a log-shipper to ELK and triggering on log queries, but rejected it because ELK alert latency is minutes — the freshness probe must page within 60s.
When it fails. If the monitor itself is degraded, every other system loses its detection signal — the worst possible failure (you don't know your systems are down). Mitigation: cross-region peer health-checks (each monitor instance watches the others' heartbeats); detection: 'absence of expected heartbeats' from any monitor instance pages the on-call's pager via a parallel cross-region path.
- SQL DB · Thing DBPostgres 16 + Citus
Canonical store for posts (thing_link), comments (thing_comment), accounts (thing_account), and the EAV data_<type> tables. Holds the comment-tree skeleton (ordered preorder traversal with depth/parent), post metadata (title, url, author, created_at, deleted, spam, ranker_version_emitted), and the (fuzzed) cumulative ups/downs counters used for read-time score reconciliation. The submission_idempotency table dedups POST /submit by (Idempotency-Key, author_id) with a 24h TTL.
Why it exists. Reddit's classic ThingDB-on-Postgres survives because (a) submissions and comments are write-rare (~400 writes/sec total) so single-leader-per-shard is fine, and (b) Postgres's relational guarantees on author + thread_id integrity matter for product invariants (comments cannot reference a deleted post). We rejected Cassandra here because the comment-tree skeleton query (ORDER BY preorder_idx WITHIN post_id) is exactly the range scan Cassandra does poorly without a hand-tuned partition key, AND we need transactional INSERT-and-bump-counter on submissions.
When it fails. Two pages on this tier. (1) Compaction/autovacuum storm at peak hour — write p99 spikes from 10ms to 2s and bleeds into vote-svc (which carries both submission and comment writes — e12/e30). Detection: monitor scrapes pg_stat_progress_vacuum + write p99. Mitigation: per-shard autovacuum cost-limit caps tied to peak QPS (back off at >80% peak), manual VACUUM scheduled off-peak, per-shard write-rate cap at vote-svc so back-pressure flows to the user (429) rather than to the ranker queue. (2) Sync replication backpressure — a slow follower stalls the primary's commit, write p99 climbs across the shard. Detection: pg_stat_replication.lag >50ms; mitigation: degrade the offending follower to async, alert on quorum loss.
- NoSQL DB · Vote StoreCassandra 4.1 (RF=3, quorum)
The high-write side of the system. Stores (user_id, target_id, target_kind, dir, ts, idempotency_key, scrubbed_flag) — partition key user_id (a user's vote-history is one partition), clustered by target_id. A second materialized view is partition-keyed by target_id for the ranker's per-post tally. Bloom filters give sub-ms 'did user X vote on target Y?' negative answers used on every comment-tree decoration. The outbox table (outbox) is co-partitioned with the vote table on user_id so a single LOGGED batch atomically writes the vote row + the outbox event.
Why it exists. Votes are the heaviest write workload (~15K QPS peak, ~500M/day, 73TB indexed at 5yr) AND the heaviest negative-lookup read (the comment-tree decoration: 'has this user voted on each of these 100 comments?'). Cassandra was specifically chosen by Reddit in 2010-2011 over Postgres because Bloom-filter negative lookups are sub-ms at this scale. Postgres unique-index lookups, by contrast, force a B-tree probe per check — orders of magnitude slower at 5K decoration QPS × 100 keys. We considered DynamoDB; rejected because the per-partition write ceiling (~3K WCU) hits the viral-post hot key.
When it fails. Three pages on this tier. (1) Hot-partition write meltdown (#1 viral post takes 1K votes/sec): leading indicator is coordinator_queue_depth / pending_native_transport_requests >100 — paged 30+ seconds ahead of the trailing write_latency p99 >100ms. Mitigation: counter shard-splitting at the application layer — vote-svc writes a vote into one of N=16 sub-partitions chosen by hash(user_id), and the ranker sums on read. Cost: 16-way fan-out on the ranker's tally read = 16K Cassandra reads/sec for the #1 post at 1× scale, 640K at 10× — the ranker's 5s per-post tally cache (see ranker.keyChoices) is the load-bearing mitigation for this read-amp. (2) MV lag drift: monitor scrapes db_score - mv_score per-post; alert at >5% drift triggers outbox drain + re-rank. (3) LWT contention storm under the same hot-partition load: external retries on LWT multiply Paxos contention — the e13 edge uses retryPolicy=none for exactly this reason.
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Typical clarifications to surface:
- Scope of "ranking"? Just hot/top/new/best, or rising/controversial too?
hot= log10(score) + sign·t_submit/45000 with t_submit = submission time in epoch-seconds (newer posts score higher);best= Wilson lower bound;top= score over a time window. - Sub-communities + front page? Front page is a fan-in over many sub listings; that drives a 50K-job/min ranker workload separately from per-sub ranking.
- Comment depth? Long-tail: 99% of threads are <100 comments; the top 0.1% are >10K. The renderer must handle both.
- Authenticated vs anonymous? Anonymous reads dominate. Votes are auth-gated and per-user dedup'd.
- Retention? Votes forever (rebuild listings, abuse detection). Listings rebuilt continuously. Comment-tree HTML 10m TTL.
- Anti-fraud surface? Vote-cabal detection is existential; without it, a $100/day farm pollutes the front page.
- SLOs? p99 listing fetch <80ms cached / <400ms cold; vote-ack <120ms; rank-visibility <5s; comment-tree <2s.
Assumptions to state:
- 73M DAU, ~10 PV/DAU, 3 sessions/DAU, 5 listing fetches/session.
- 500M votes/day, 3M submissions/day, 30M comments/day.
- Front page ≈ 40% of listing QPS (one logical key — the load-bearing skew).
02Functional reqs
What must this system actually do?
- Submit a link or text post to a sub.
- Vote up/down on a submission or comment (idempotent per user × target).
- Render a hot/top/new/best listing for a sub or for the front page.
- Render a comment thread with depth, the requesting user's vote-state, and MoreChildren stubs at >250 children.
- Reply to a comment.
- (Optional) search submissions/comments — out of scope for this canonical's main flow.
03Non-functional
What must it promise about speed, uptime and correctness?
- Availability: 99.95% on the read path; 99.9% on the write path; rank-visibility SLO 99% within 5s.
- Latency:
- Listing cached: p50 20ms / p99 80ms / p99.9 200ms.
- Listing cold (cache miss): p50 100ms / p99 400ms / p99.9 1.2s.
- Vote-ack: p50 30ms / p99 120ms / p99.9 300ms.
- Comment-tree: p50 150ms / p99 800ms / p99.9 2s.
- Durability: Votes must not be lost (the rank can be rebuilt from votes; the inverse isn't true). Submissions must not be lost.
- Consistency: Read-your-write for the voter (their own vote arrow reflects state immediately); eventual for ranking (the global hot list takes ~5s to update). Comments are RYW for the author.
- Scalability: Horizontal on every tier. The viral post (1K votes/sec on a single key) must not melt anything.
- Anti-abuse: Vote-cabal scrubbing is in-band durability for ranking quality. Shadow-removal — not blocking — is the canonical pattern.
04Capacity estimation
How much load and data does this have to hold?
Inputs: 73M DAU, 3 sessions/DAU, 5 listings/session, 500M votes/day, 3M submissions/day, 10 comments/submission, peak 2.5×.
- Listing QPS =
73M × 3 × 5 / 86400≈ 12.7K avg, 32K peak. Front page ≈ 40% = 13K QPS on one logical key at peak — mitigated by listing-cache's 16-way front-page fan-out (listing:_all:hot:0..15). - Vote QPS =
500M / 86400≈ 5.8K avg, 15K peak. Comment votes ≈ 85% of total. Each vote fans out to 1 sync DB write + 1 sync Redis ZADD + ~3 async ZADDs (ranker emit). Effective write amp ~5–7× → up to 75K Redis ops/s on the listing-cache. - Submission QPS =
3M / 86400≈ 35 avg, 90 peak. - Comment QPS =
35 × 10≈ 350 avg, 900 peak. - Vote storage (5 yr, indexed) =
500M × 365 × 5 × 80B / 1e12≈ 73 TB. Heaviest table by far. - Comment storage (1 yr) =
30M × 365 × 1KB≈ 11 TB/yr. - Ranker compute =
50K active subs / 30s cadence≈ 1.7K jobs/s for the cold tail, plus per-vote rank for the top 1K posts. ~25 cores at peak — trivial CPU, I/O bound. - Egress (listings) =
32K × 50KB≈ 12.8 Gbps peak. - Hot-key write skew (Zipf α≈1) — top-100 posts capture ~35% of votes. The #1 post alone takes ~1K QPS on one Cassandra partition — exceeds Cassandra's per-coordinator queue ceiling at QUORUM (
pending_native_transport_requestsbudget), hence the 16-way counter-shard-split (hash(user_id) % 16). - Counter-split read amp (the hidden cost) — every ranker tally on the #1 post reads 16 sub-partitions. At per-vote rank cadence on the top-1K posts:
1K vote/sec × 16 sub-partitions = 16K Cassandra reads/sec for #1 post aloneat 1× scale;640K reads/sec at 10× scale. The ranker's 5s per-post tally cache drops this to ~200 reads/sec by amortizing across votes. - Ranker tally read-amp budget at 1× scale: 1K hot posts × 16 fan-out / 5s cache window = 3.2K Cassandra reads/sec total ranker tally load — well within the 48-node ring's 50K reads/sec/node aggregate.
05API design
What does the outside world call, and what comes back?
POST /api/v1/vote
Authorization: Bearer <token>
Idempotency-Key: <uuid>
Content-Type: application/json
{ "target_id": "t3_abc123", "target_kind": "post", "dir": 1 }
200 OK
{ "score": 423, "userVote": 1, "rankerVersionPin": "v17" }
# Conflict on idempotency replay (same body):
200 OK (returns the prior result — no double-count)
# Conflict on logical conflict (same idempotency_key, different body):
409 { "error": "idempotency_conflict" }
GET /api/v1/r/:sub/hot.json?after=<cursor>&limit=25
Cache-Control: public, max-age=10 # anonymous; edge-cached, keyed on (sub, sort, cursor)
200 OK
{ "after": "...", "items": [ { "id": "t3_abc", "title": "...", "score": 1234 }, ... ] }
# No per-user fields in the cacheable body — a shared CDN object is served to every user.
# The signed-in arrow state is fetched separately and merged client-side (see GET /api/v1/votes).
POST /api/v1/submit
Idempotency-Key: <uuid>
{ "sub": "programming", "title": "...", "url": "..." }
201 Created { "id": "t3_xyz", "permalink": "/r/programming/comments/xyz/..." }
# Idempotent on retry: post_id = hash(idempotency_key, author_id) so retries hit the same row.
POST /api/v1/comments/:post_id
Authorization: Bearer <token>
Idempotency-Key: <uuid>
{ "parent_id": "t1_def456", "body": "..." } # parent_id null/omitted for a top-level comment
201 Created { "id": "t1_ghi789", "permalink": "/r/programming/comments/xyz/.../ghi789" }
# Comment-write path (vote-svc dispatch): LWT-deduped on (author_id, idempotency_key), then
# INSERT thing_comment into thing-db with preorder_idx assignment + parent child-count bump,
# and a 'new comment' outbox event for the ranker (best-sort recompute). 24h idempotency TTL.
# Backed by hldEdge e30 (vote-svc → thing-db: thing_comment INSERT), sibling of e12 (submission INSERT).
GET /api/v1/comments/:post_id?sort=best&limit=200
Cache-Control: public, max-age=10 # anonymous, per-thread edge cache
200 OK
{ "post": { ... }, "comments": [ { "id": "t1_...", "depth": 0, "parent": null, "html": "...", "score": 12 }, ... ], "moreChildren": [ "t1_more_42", ... ] }
# Again no per-user fields in the cached body; arrows come from GET /api/v1/votes below.
GET /api/v1/morechildren?post_id=t3_abc123&children=t1_more_42,t1_more_57&sort=best&limit=200
Cache-Control: public, max-age=10 # anonymous, per (post_id, children-batch) edge cache
200 OK
{ "comments": [ { "id": "t1_...", "depth": 4, "parent": "t1_...", "html": "...", "score": 3 }, ... ], "moreChildren": [ "t1_more_88", ... ] }
# Resolves the MoreChildren stubs emitted by GET /comments: pass the post_id (link_id) plus the
# comma-separated stub ids, get back the next flattened batch (same comment objects) with deeper
# stubs of its own, so a 35K-comment tree paginates lazily. Same comment-tree read path as /comments
# (gw → comment-tree → render-cache mget + thing-db preorder skeleton); arrows via GET /api/v1/votes.
# Per-user vote-arrow overlay — authenticated, NEVER cached. The client passes the ids it
# just rendered from the (edge-cached) listing/comments and merges arrow state client-side.
# Backed by the same per-user Cassandra Bloom probe the comment-tree path already uses.
GET /api/v1/votes?targets=t3_abc,t1_more_42,...
Authorization: Bearer <token>
Cache-Control: private, no-store
200 OK
{ "votes": { "t3_abc": 1, "t1_xx": 0 } } # -1 | 0 | +1 per target06Data model
What gets stored, and what is it looked up by?
thing_link / data_link (Postgres, sharded by hash(post_id)):
id(post_id),sub,author_id,title,url,created_at,deleted,spam,ups,downs,score_fuzzed,ranker_version_emitted.- Index: PK on id; secondary on
(sub, created_at DESC)for new-listing fallback; partial onexpires_at WHERE deleted=false.
thing_comment / data_comment (Postgres, collocated with parent post on hash(post_id)):
id,post_id,parent_id,depth,preorder_idx,author_id,body,created_at,ups,downs,score_fuzzed.- Index: PK on id; covering on
(post_id, preorder_idx)for tree skeleton fetch; partial on deleted.
submission_idempotency (Postgres, sharded by hash(idempotency_key)):
idempotency_key,author_id,post_id,created_at. TTL 24h.
user_votes (Cassandra, partition key user_id, cluster key (target_id, target_kind)):
- Columns:
dir(-1, 0, +1),ts,idempotency_key,scrubbed_flag,fraud_score. - Bloom filter tuned to fp=0.01 — sub-ms negative-lookup on every comment-tree decoration.
- Materialized view partitioned by
target_idfor the ranker's per-target vote tally. MV write is async to source; bounded staleness ~10ms; reconciliation alerts at >5% drift.
outbox (Cassandra, partition key user_id — co-partitioned with user_votes for single-partition LOGGED batch atomicity):
- Columns:
event_id,event_body,created_at,published_at. Tailed by outbox-tailer.
listing:{sub}:{sort} (Redis ZSET, 16 hash slots):
- Members:
thing_id. Scores: hot/top/best score from ranker. ZRANGE 0,-1 returns the page. - Front-page exception:
listing:_all:hot:0..15is the same ZSET written to all 16 shards by the ranker; readers pick a shard at random.
comment_html:{post_id}:{comment_id} (Memcached, ketama hashed on post_id):
- Pre-rendered HTML chunk. TTL 600s, recomputed on score-delta-cross-fuzz-threshold.
votes.events (Kafka, 256 partitions, key=post_id, RF=3):
- Event body:
(user_id, target_id, target_kind, dir, ts, idempotency_key, fraud_signals, ranker_version_pin). - Retention 7d (long enough to replay through a multi-day ranker outage).
07High-level design
Which components handle a request, and in what order?
Architecture summary:
- Edge — Fastly CDN. 10s TTL on listing JSON, surrogate-key purge on every ranker emit. Absorbs ~95% of read QPS.
- Anycast Load Balancer. L7 health-checked, terminates anycast traffic to healthy gateway pods cross-AZ. Replaces DNS-based failover.
- API Gateway — Envoy. WAF, per-account/per-IP token-bucket rate limits via rate-limit-redis, JWT auth (local + JWKS cache), dispatch by URL.
- Three serving lanes (split by intent):
- Listing Read — ZRANGE listing-cache, multi_get render-cache, hydrate cold from Postgres. Go for tail-latency control.
- Comment Tree — preorder skeleton from Postgres, HTML chunks from render-cache, user-vote overlay via Cassandra Bloom probe.
- Vote / Submit / Comment — Cassandra LWT + LOGGED-batch outbox (vote-row + outbox row co-partitioned on user_id), Postgres counter bump + INSERT (for submits, idempotent via
submission_idempotency), Kafka publish best-effort. Comment-create (POST /api/v1/comments/:post_id) also dispatches here: an LWT-dedupedINSERT thing_commentwithpreorder_idxassignment + parent child-count bump, then a 'new comment' outbox event. This comment-write is backed bye30(vote-svc → thing-db:thing_commentINSERT + preorder_idx + child-count bump), the sibling ofe12(vote-svc → thing-db: submission INSERT + counter bump).
- Listing Cache — Redis cluster. Sorted sets keyed
listing:{sub}:{sort}. 16 shards, 1 follower per shard, async replication. Front-page replicated across all 16 shards for hot-key fan-out. - Render Cache — Memcached cluster. Per-comment HTML, per-listing-card HTML. 32 nodes, ketama hash on
post_id. - Vote Stream — Kafka. 256 partitions on
votes.events, RF=3, retention 7d. Per-partition ordering onpost_idso the ranker batch-coalesces a hot post's votes through one consumer slot. - Outbox Tailer. Sidecar per vote-svc replica, publishes outbox rows where
published_at IS NULL. The catch-up path; vote-svc publishes synchronously on the happy path. - Ranker Worker. Consumes Kafka, recomputes hot/best/top, ZADDs to listing-cache. Per-vote on top 1K posts (via 5s tally cache to bound counter-split read amp); per-batch every 5s on cold posts.
- Fraud Scrubber Worker. Separate Kafka consumer group, async cabal detection, shadow-removal.
- Thing DB — Postgres + Citus. EAV schema (thing_<type> + data_<type>), 16 hash shards on
post_id, sync followers cross-AZ for RPO=0 on submits. - Vote Store — Cassandra. RF=3, W=R=LOCAL_QUORUM, 48 physical nodes × 256 vnodes, partitioned by
user_idwith target_id MV. Bloom-filter negative-lookup for per-comment "did I vote?" decoration. - Monitor — Prometheus + synthetic prober. Hosts the freshness probe, the outbox-lag probe, the synthetic /vote/listing/comment round-trip prober. Alertmanager pages on SLO breach.
- Rate-limit-redis. Single-shard Redis + 1 follower for token-bucket state.
Read flow (cached): Client → CDN (hit) → 200 / cached JSON. Origin sees nothing.
Read flow (cold): Client → CDN (miss) → LB → Gateway → Listing Read → Redis ZRANGE → Memcached mget → respond. If any of those misses, fall through to Postgres replicas.
Comment tree flow: Client → CDN (per-thread, 10s TTL) → LB → Gateway → Comment Tree → Memcached mget for HTML → Postgres skeleton query (cold) → Cassandra Bloom probe for user vote-state → respond.
Write flow (vote): Client → LB → Gateway → Vote-svc → Cassandra LWT (vote-row + outbox in one LOGGED batch) → return 200. Vote-svc best-effort publish to Kafka. Outbox-tailer catches anything that didn't make it. Ranker consumes, recomputes (with 5s tally cache), ZADDs. Fraud scrubber consumes in parallel.
Submit flow: Client → LB → Gateway → Vote-svc → Postgres INSERT (with submission_idempotency ON CONFLICT DO NOTHING) → emit outbox event → return 201. Ranker picks it up; new posts seed the new ZSET immediately and graduate to hot once they accumulate score.
08Deep dives
Which part breaks first, and what do you do about it?
1. Why precomputed listings and not read-time ranking. The naive SELECT … ORDER BY hot(ups, downs, date) LIMIT 25 works at <500 listing QPS. Above that, the planner cost dominates — and front-page fan-in across 50K subs is impossible with replica reads. Reddit's pre-2010 architecture melted exactly here. Precomputed Redis ZSETs reduce read-time work to O(log N) and let the ranker batch-coalesce 1K vote-ZADDs into 1 update on viral posts.
2. The vote-storage decision (Postgres → Cassandra). Reddit's 2010-2011 migration was driven by the per-comment "has user X voted on Y?" lookup, which fires hundreds of times per thread fetch. Cassandra's Bloom filters give sub-ms negative answers at 5K decoration QPS × 100 keys. Postgres's B-tree probe at the same rate would saturate a shard. The downside: no relational integrity on the vote table — but votes are append-only, and the audit row is the source of truth, so this trade flies.
3. Comment-tree precompute and the MoreChildren stub. A 50K-comment thread cannot be rendered in <2s on every fetch. The canonical trick: store the comment tree as a flat preorder traversal (an ordered list of comment_ids with depth/parent), pre-render each comment's HTML into Memcached, and surface a MoreChildren stub at any level with >250 children. The thread fetch becomes one tree-skeleton query + one mget — a constant-time read regardless of thread size.
4. Outbox over 2PC, and the MV-atomicity caveat. The vote-row write and the Kafka publish must atomically commit or both not — otherwise we get a vote with no ranker-emit, or a Kafka event with no audit row. We do NOT use 2PC across Cassandra and Kafka because cross-system 2PC is operationally a tar pit. Instead, we co-partition the user_votes table and the outbox table on user_id and write both in one Cassandra LOGGED batch — single-partition, atomic for those two writes. The ranker is idempotent on (post_id, idempotency_key) so duplicate publishes are harmless.
The caveat: the materialized view (partition-keyed by target_id) is NOT in the LOGGED batch — Cassandra MV writes are async to the source. The ranker reads the MV via e17 with bounded staleness (~10ms typical). vote-store.failureMode authoritatively documents this and the >5% reconciliation drift trigger.
5. Idempotency on votes AND submits. Every POST /vote ships an Idempotency-Key. Vote-svc runs Cassandra LWT (Paxos compare-and-set) on (user_id, target_id, idempotency_key). If the key already exists with the same body, return the prior result. If the key exists with a different body, reject with 409. TTL 24h.
POST /submit is also idempotent: post_id = hash(idempotency_key, author_id), so a retry generates the same post_id, and the INSERT uses ON CONFLICT DO NOTHING against the submission_idempotency table.
The edges e12 and e13 are both retryPolicy: "none" because (a) the endpoint is server-side idempotent, the retry should come from the client with the SAME Idempotency-Key, and (b) external retries on LWT multiply Paxos contention on the hot partition.
6. Hot-key counter-split for the viral post. The #1 post takes 1K votes/sec on one Cassandra partition. Even with QUORUM=2/3 spread cross-AZ, the coordinator queue on the partition owner saturates. Mitigation: vote-svc writes a vote into one of N=16 sub-partitions chosen by hash(user_id), and the ranker sums on read.
The cost of the mitigation — the ranker tally read on the #1 post becomes a 16-way fan-out: at 1K vote/sec per-vote cadence, that's 16K Cassandra reads/sec just for the #1 post's score recompute. The ranker's 5s per-post tally cache (in-process, see ranker.keyChoices) bounds this to ~200 reads/sec by amortizing across votes. Without that cache, the counter-split fix would just move the meltdown from the write path to the read path.
7. Read-your-write for the voter. After POST /vote, the voter's next listing/thread fetch must show their own vote-arrow reflecting their action. Cassandra LOCAL_QUORUM=2/3 reads give RYW because W+R > N. The trickier case is the score: the user's vote is durable, but the ranker hasn't run yet and the score they see is the pre-vote score + an optimistic +1 client-side. We don't promise the GLOBAL re-rank is RYW — that's eventual within ~5s.
8. Gray failure: silent stale ranker AND dead outbox-tailer. The scariest failures have no error.
- Stale ranker — the freshness probe on the monitor (
now() - max(listing.generated_at) >5min) pages on staleness, not on errors. - Dead outbox-tailer — the score-divergence reconciliation alone CANNOT catch this because both
db_scoreandranker_scorestop moving together. The load-bearing signal isnow() - min(outbox.created_at WHERE published_at IS NULL) >30s, scraped by the monitor (e27).
Per the SRE Workbook, symptom-based alerts (the user's experience) catch gray failures that infrastructure metrics miss.
9. Fraud scrubber and shadow-removal. The detector runs async on the same vote stream (separate consumer group). On suspicion, it writes a scrub-mark in Cassandra and tells the ranker to recompute. The voter still sees their vote registered locally (shadow-ban semantics) so attackers don't iterate. Score-fuzzing on display (random ±N noise on ups/downs) is the same idea: if attackers see exact deltas, they reverse-engineer the scrubber.
10. Retry math at peak (why circuit-breakers, not exp-backoff, on e3 and e6). Naive exp-backoff on the read-path under partial Redis failure: 3 client retries × 3 cdn retries × 3 gateway retries = 27 attempts at the listing-read tier per user click. At 32K listing QPS peak with 50% Redis-shard partition, that's 16K × 27 = 432K attempted connections at listing-read inside the failover window — the listing-read replicas' connection pools fill in <1s. Circuit-breakers on e3 (gw → listing-read) and e6 (listing-read → listing-cache) shed load fast under degradation, matching the listing-cache.failureMode runbook of "admission control caps origin QPS to 10% until hit-rate recovers."
09Trade-offs
What did this design cost, and what breaks at 10×?
What breaks at 10× scale (~150K vote QPS, 320K listing QPS):
- Listing-cache front-page hot key. A single Redis shard can do ~150K ops/s; the front page alone produces ~40K ZADD-equivalents/s. Mitigation already designed in: front-page replicated across all 16 shards (
listing:_all:hot:0..15); at 10× we extend to 64 shards. - Cassandra hot-partition (#1 post = 10K votes/sec). Counter-shard-split factor needs to grow from 16 to 64; ranker's 5s tally cache extends to 10s to keep the read-amp manageable (
10K × 64 / 10s = 64K reads/sec for #1 post alone). - Render-cache megathread. A 500K-comment thread (next /r/AskReddit megathread) breaks the per-shard memory budget. Mitigation: chunked subtree caching with per-subtree expiry, plus a dedicated 'megathread' Memcached pool.
- Kafka 256 partitions become the rebalance-storm risk; 1024 partitions or topic-per-shard.
- Ranker becomes the bottleneck — at 150K vote QPS with 7× fan-out, that's ~1M Redis ops/s on listing-cache from the ranker alone. Counter-coalescing window grows from 1s to 5s; rank-visibility SLO loosens to 10s.
Per-component failure stories:
- CDN edge cache poisoning — surrogate-key purge fixes; hard 30s TTL ceiling caps blast radius.
- Listing-cache cold restart (May 2023 Reddit shape) — admission control at the LB caps origin QPS to 10% until hit-rate recovers; never reopen the door before the cache is warm.
- Render-cache topology change — ketama + rolling-warm-then-add, never expose a cold shard to reads.
- Vote-stream consumer lag — load-shed low-value events first, autoscale on lag not CPU, kill-switch on fraud-suspect votes.
- Cassandra hot partition — counter shard-split + ranker 5s tally cache; runbook documents the post_id-to-sub-partition fan-out and the read-amp ceiling.
- Postgres compaction storm — autovacuum cost-limits tied to peak QPS; manual VACUUM off-peak.
- Ranker version bump — atomic feature-flag flip with version pin; cursor-cache invalidates on version change so v1 and v2 scores never mix in one user session.
- Kafka poison message — max-retries-then-DLQ; per-partition isolation so 1/256 of votes never blocks the rest.
- Region failover — anycast at the LB instead of DNS to avoid the "DNS hasn't propagated" botch.
Trade-offs we accepted
- Async rank visibility (~5s) instead of synchronous re-rank on every vote. Synchronous would push vote-ack p99 above 120ms and saturate the listing-cache during vote storms. Voters see their own arrow update immediately (RYW); the global hot list catches up within 5s.
- Cassandra MV is async to the source write. The (target_id partition) MV used by the ranker for per-target tally is NOT in the same LOGGED batch as the (user_id partition) source write. Bounded staleness ~10ms typical; the >5% reconciliation drift trigger catches anything pathological. We accept this rather than the 2-partition LOGGED batch cost or routing the ranker through Kafka tally events.
- Memcached single-replica. Render-cache is recomputable from thing-db; we trade a 5-minute warm-up on cluster-loss for half the memory cost.
- Single-region active-standby with PER-TIER RTO (NOT a flat 30s).
- LB anycast withdraw: ~30s region cutover.
- thing-db: cross-region promotion ~10min from PITR or async-replica catch-up. The intra-region
failover.rtoSeconds=30documents the AZ failover; cross-region is bigger. - vote-store: with NetworkTopologyStrategy + LOCAL_QUORUM in standby DC, RTO ~2min, RPO=0 within standby; cross-DC writes during partition return 503 (CP under partition).
- listing-cache: cold-rebuild from Cassandra ~10min for full ZSET set (50K subs × 1K thing_ids); admission-control degraded service during rebuild.
- The "30s RTO" headline applies only to the edge cutover. The product is degraded for ~10min during cross-region failover. Building active-active would require a CRDT-based vote-merge layer; we do not yet need it.
- Score-fuzzing on display. Users see a ±N-noised score, not the true score, so cabal-scrubber output is invisible to attackers. This is a UX cost — power users notice — but it's the load-bearing anti-iteration moat.
- Counter-shard-split read-amp absorbed by 5s tally cache, not by a real counter table. A real counter-table (CRDT counters or write-through aggregate row) would eliminate the read-amp but introduce a second consistency surface (counter row vs vote rows). The 5s tally cache is simpler and the staleness is below the rank-visibility SLO. Acceptable.
Primary sources
- Reddit Architecture Overview (reddit-archive wiki)
- Reddit on Postgres → Cassandra (InfoQ 2010)
- Salihefendic — How Reddit ranking algorithms work
- Wilson-score lower-bound for comment ranking
- Fastly — How we built r/place
- Pi Day 2023 outage retro
- SRE Workbook — Monitoring Distributed Systems
- SRE Workbook — Addressing Cascading Failures
Now defend it
Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.
Work Reddit / Hacker News yourselfMore in Feeds, Timelines, Counters & Ranking
What to show and in what order: fanout on write versus read, hot/top/new scoring, approximate counters, trending, and recommendation.
- Twitter / X TimelinePush or pull? Both. The canonical fanout problem.
- Instagram News FeedRanked feed with cursor pagination. No `OFFSET`.
- Like Button at ScaleEventual consistency, but the liker sees their own write. Counts are approximate by design; hot keys are the real enemy.
- View Count on a Video/PostDedup, bot-filter, batched aggregation.
- Trending TopicsSliding windows + Count-Min Sketch + top-K.