Like Button at Scale — a worked solution
Eventual consistency, but the liker sees their own write. Counts are approximate by design; hot keys are the real enemy.
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 Like Button at Scale workspaceThe problem
Build the like button for a social product at scale: every post, video, story, and tweet carries a heart. A user taps it once to like, taps again to un-like; everyone else sees a count that updates in near-real-time. The system is read-dominated by 100–500×, hot-keyed on viral posts, must give the liker read-your-write on their own heart, and accepts approximate counts everywhere else.
The interview answer ("one Redis INCR, one Postgres row, done") survives ~10K writes/s and 100K reads/s on a single non-viral product. We're sizing for 2B likes/day with single-key viral peaks at 50K writes/s and read amplification 100–500×, and we have to defend it at 3am when Ronaldo posts.
The reference architecture
What each component is for
- Client (SDK)iOS / Android / Web SDK
Renders the post with its like count and hearted/un-hearted state, captures taps, sends POST /v1/likes with an Idempotency-Key, and applies an optimistic UI flip on tap so the user sees their own write before the server has even ack'd. Reconciles the local state on the next GET.
Why it exists. The SDK owns read-your-write for the liker. The backend is eventually consistent across viewers; the liker MUST see their heart flip immediately even if the cache / CDN / replica has not yet caught up. Client-side optimism + a server-confirm token is the only honest way to deliver that without an active session pin into every replica.
When it fails. Network drops mid-write: client retries the same Idempotency-Key — server dedups, no double-count. Optimistic state stuck if server never acks → SDK reconciles on next GET (state from server wins). A bug in the optimistic state machine = thousand-tap users seeing wrong heart state → caught by the count_flicker_per_user metric backend-side.
- CDN (edge cache)Cloudflare / Fastly
Caches the per-post {count, topReactions} JSON at the edge for 5–30 s (TTL is per-post heat — viral posts get 5 s, cold posts 30 s; heat signal comes from the gateway's hot-key sampler). Personalized hasLiked is bypassed — POSTs and any request bearing an Authorization header are pass-through to origin.
Why it exists. The dominant read shape is many viewers reading the same hot post's count. Without an edge cache, ~99% of cold-load reads hit origin and a viral post's 5M+ reads/s would saturate every read-svc pod globally. The CDN absorbs 80–95% of read QPS; origin is sized for the residual 5–20% plus the hasLiked traffic.
When it fails. Regional POP outage → traffic fails over to neighbor POP, colder cache, miss rate spikes from ~5% to ~20%, origin sees 4× more traffic; count-cache + read-svc must absorb. CDN serves stale count after post-delete / mass-unlike → purge-fanout failure is the primary cause; safety net is the 30 s natural TTL.
- API GatewayEnvoy + WAF + rate-limit + hot-key sampler
Terminates TLS, applies WAF, enforces per-user (100 writes/min) and per-IP (1000 writes/min) rate limits on POST /v1/likes, mTLS-binds to downstream services, samples 1% of traffic for hot-key detection via a count-min sketch, and routes by URL: POST /likes → write-svc, GET /posts/:id/likes → read-svc.
Why it exists. A single common edge that does the things every service would otherwise re-implement: TLS, authn, rate-limit, request tracing. Pulling it out lets read-svc / write-svc focus on the like-specific work. The hot-key sampler lives here because the gateway sees the un-salted, un-sharded request URLs — downstream services see post_ids already smeared by sub-counter shards and can't tell which key is hot.
When it fails. Pod crash: LB ejects within 5 s; cluster of 40 absorbs one loss with no user impact. Hot-key sampler misses a flash spike (sub-30 s viral curve): count-cache shard saturates for 30–60 s before key-split kicks in — mitigated by always-replicating count:* to 4 hot pools even before the sampler flags. WAF false-positive blocks a legit author: page on the waf_block_rate_per_authed_user SLO.
- Like-Read ServiceGo / gRPC
Stateless service that assembles the per-post like response: GETs the aggregate count:{post_id} from count-cache (or SUMs sub-counters / falls through to counter-store on miss), and for authenticated callers checks hasLiked(user, post) from a per-user Bloom filter, falling through to a per-(user, post) lookup in edge-store. Returns the merged payload in a single round-trip.
Why it exists. Splitting the read and write services lets each scale independently — reads run 100–500× writes (~18M/s peak vs ~93K/s peak, per the capacity anchors) and need wildly different connection-pool sizing, GC tuning, and JIT warmup than the write side. Co-locating them would mean a write deploy could ship a read regression on the busiest endpoint in the product.
When it fails. count-cache shard unavailable → fall through to counter-store; p99 jumps from 5 ms to 100 ms — the explicit graceful-degradation budget. counter-store unavailable too → serve {count: stale, lastFresh: ts} from the last good response (1 m TTL); UI degrades to 'about N' rounded. The read NEVER 5xx's — likes are decorative, not critical to the product loading.
- Like-Write ServiceGo / gRPC
Accepts POST /v1/likes with an Idempotency-Key. GETs the key from idem-store to dedup retries; if new, issues a Cassandra BATCH against the edge-store partition for the user that writes BOTH the edge row AND a row into the outbox table (partition-local atomicity — both succeed or both fail). After the batch acks, publishes to Kafka likes.committed and SETs the idem key with {state, intent_uuid, event_id}. Returns 202 once the publish acks. If the publish fails after the batch commit, the outbox row remains and a relay process drains it on the next pass — the like is durable and the event WILL appear, late.
Why it exists. Outbox separates 'durable enough for the user's 202' from 'fanned-out to all downstream'. Edge-store is the system-of-record for (user, post, state) — a viewer's hasLiked query reads from it. Kafka is the spine for everything else (counters, anti-abuse) — without a single event stream, each consumer polling edge-store would melt it.
When it fails. Edge-store quorum loss → write fails fast with 503 + Retry-After (better than hanging). Kafka publish times out after edge-store BATCH commit → outbox row survives; relay drains within 5 s; consumer dedups via intent_uuid (NOT event_id — two write-svc pods retrying the same intent generate DIFFERENT event_ids but the same intent_uuid, and the aggregator dedups on intent_uuid). Idem-store down → fail-CLOSED (refuse the write); accepting an un-dedupable write risks double-count, the worse outcome.
- Idempotency StoreRedis Cluster (24h TTL, AOF-persisted)
Key idem:{user_id}:{post_id}:{intent_uuid} → {state, ts, event_id}. write-svc GETs before any side-effect to dedup; SETs the key after Kafka publish + edge-store commit, with TTL=86400 s.
Why it exists. At-least-once retries are the default on every layer — mobile SDK retries on connection drop, the gateway retries on 503, internal RPC retries on transient errors. Without a dedicated dedup store, every retried POST becomes a double-increment. The dedup window has to outlive the worst-case retry pipeline (mobile app restored from background after a flight = measured 18 h tail), so 24 h is the floor.
When it fails. Shard unavailable → fail-CLOSED on writes that route to it (refuse rather than risk double-count); read-svc unaffected. Cluster restart → 5 s of writes may 503; mobile SDK retries succeed once cluster is back; no data loss because the post-commit SET hasn't happened yet so the retry simply re-runs the full path.
- Counter CacheRedis Cluster (mcrouter-style hot-key replication)
Stores the per-post aggregate count:{post_id} and the per-shard sub-counter set count:{post_id}:s0..s63. read-svc GETs the aggregate (5 s TTL); on miss it SUMs the 64 sub-counters in one MGET. Aggregator updates sub-counters on each Flink emit and re-SETs the aggregate with a monotonic generation_id (slow emitter's stale write LOSES the SET race).
Why it exists. A single Redis INCR on count:{post} for a viral post is a hot-key disaster — one shard saturates at ~100K ops/s while the rest sit idle. Per-shard sub-counters spread the write across N keys; the read pays N GETs but those are 1–2 µs each. The pre-aggregated count:{post} lets ≥95% of read-svc reads do a single GET; only on TTL expiry / cache miss does the SUM run.
When it fails. Shard primary failure: 15 s failover RTO (cluster mode, async repl), 5 s of recent INCRs may be lost — recovered by Flink replaying the last 5 s window from Kafka. Cluster restart cold: all SUM operations hit counter-store; counter-store sized to absorb 10× steady miss rate for ~5 min. Hot key on an under-replicated post: one read pool saturates → mitigated by hot-key detector → key-split bumps replication factor.
- Like-Edge StoreCassandra / ScyllaDB (RF=3, LOCAL_QUORUM)
Source-of-truth for (user_id, post_id, state, ts, reaction_type, void_flag) rows. Sharded by user_id (so hasLiked queries are single-partition). Write path UPSERTs with IF state != ? semantics; read path is point-lookup by (user_id, post_id) for hasLiked.
Why it exists. We need an authoritative source for who liked what, separately from the counter. Counters can be approximate; the like edge cannot (a user must see whether they liked a given post; we need to support 'undo' and GDPR audit). Sharding by user_id — not post_id — optimizes the hot read pattern: hasLiked-for-this-feed-of-50-posts becomes 50 same-shard lookups, not 50 cross-shard scatter-gathers.
When it fails. Single AZ loss: writes continue with W=2 from 2 surviving replicas; no user-visible impact. Quorum loss (region-level partition isolates 2 of 3 replicas): writes BREAK in that region — write-svc fails fast → SDK retries → eventual consistency once partition heals.
- Counter Store (CRDT)Cassandra/Scylla counter columns (PN-Counter per region × shard)
Stores per-(post_id, region, shard) counter cells. Each region writes only its own cells; the per-post aggregate is SUM(cell for cell in (post_id, *, *)). Aggregator's batched increments land here; read-svc SUMs on count-cache miss. Decrements from anti-abuse voids land here too (PN-Counter semantics).
Why it exists. Multi-region active-active writes on a single counter row would either fight (CAS retry storm) or lose increments (last-write-wins). A G-Counter / PN-Counter (per-actor partition + sum on read) is the only safe pattern — each region owns its slot, merging is commutative + associative, and a 5 min cross-region replication blip just means the global count is 5 min stale, never wrong.
When it fails. Single Cassandra node failure: writes continue at W=2; reads downgrade to LOCAL_ONE. Cross-region link severed: each region's count diverges by local_delta; on heal, CRDT merge converges automatically — no manual reconciliation. Hot post: per-shard cells balance the write load; cell-level hotspot mitigated by post_id key-split (post_id:0..post_id:7 sub-IDs at the aggregator).
- Likes Event SpineKafka (256 partitions, RF=3, acks=all)
Single durable event log. Topic likes.committed (every committed like) and likes.voided (anti-abuse retractions). Aggregator and abuse both consume from likes.committed independently, at their own pace.
Why it exists. Outbox + Kafka is the only way to fan one write out to N consumers without N×point-to-point coupling. If the aggregator and abuse scorer each polled edge-store, the edge-store reads would explode and a deploy of any consumer could break the others. Kafka decouples — durability is in the log; consumers can lag, replay, or fail without affecting write-svc latency.
When it fails. Single broker loss: ISR shrinks, controller re-elects partition leaders; ~5 s of producer blocking, no data loss. Cross-AZ partition: producers retry against surviving brokers; one-replica writes refused (good — don't sacrifice durability for tail latency). Consumer lag spike on likes.committed: counts go stale (visible after ~10 s), abuse voids land late — surface as consumer_lag_seconds SLO alert.
- Counter AggregatorApache Flink (KeyBy post_id, 1s tumbling window)
Consumes likes.committed AND likes.voided from Kafka, KeyBy (post_id, region), maintains a 1 s tumbling window in RocksDB state, and at each window-tick: (a) batches the deltas (positive for commits, negative for voids) into Cassandra counter-store updates, (b) writes the new aggregate to count-cache with versioned SET + issues a CDN Surrogate-Key purge. Dedups upstream via intent_uuid (NOT event_id) so two write-svc pod retries for the same user intent count once.
Why it exists. 1M individual INCRs into Cassandra / Redis would saturate both. Windowing collapses N+1 increments into one batched +N — at peak, a viral post sees 50K events/s collapsed into one +50K every second. Flink (vs custom code) for the at-least-once + checkpointed exactly-once-effect semantics and the operational tooling (Savepoints, RocksDB state, k8s operator).
When it fails. Flink job crash: checkpoint restore is ~2 min from S3; counts go stale during recovery — UI shows last-fresh count + a 'live' badge that hides. Slot OOM on a viral post: producer-side sub-keying spreads load across 8 slots. RocksDB compaction stall: Savepoint + restart on new pod; window state is bounded (1 s) so no replay storm. Poison event: 3 retries then publish to likes.dead-letter with the original event_id; per-partition unblocked.
- Anti-Abuse ScorerPython + gRPC ML sidecar, consumes likes.committed
Consumes likes.committed, scores each like with an ML model (per-user rate, IP-cluster entropy, device fingerprint, account age — a calibrated probability that the like is bot-driven). For likes above the threshold, publishes a compensating event to likes.voided; the aggregator decrements; the edge-store row is marked void=true.
Why it exists. Synchronous anti-abuse in the write path would push the write p99 above the budget (model inference is 20–50 ms). Async-after-the-fact lets the like commit immediately (good for the user) and retroactively void the bot ones (the counter briefly shows a higher count, then settles). This is the trade-off accepted by every major platform: visible counts are eventually-consistent with the abuse-cleaned ground truth.
When it fails. Model service down: defer scoring (consumer buffers in Redis with 24 h TTL); when model recovers, drain and void asynchronously — at most 24 h of bot likes leak through. Pod crash mid-publish: void event NOT lost (offset not yet committed); restart re-publishes; aggregator dedups. False-positive cascade (bad model deploy): manual rollback + cleanup via re-publish of likes.unvoided (we keep that capability for exactly this case). Storm of bot likes: rate-limiter degrades the user's account (1/min) and pages the abuse team.
- Config Storeetcd (Raft, 5-node cluster)
Distributes runtime control state: the gateway's count-min-sketch top-K hot post_ids (so count-cache and aggregator know which keys to sub-key / replicate), the per-post CDN heat hint (so cdn can drop its TTL on viral posts), and a small set of feature flags (canary %, kill-switches).
Why it exists. The hot-key story collapses without a control plane the dataplane can subscribe to. Pushing hot-key flags via Kafka would be slow (consumer lag = stale replication policy); pushing via direct RPC would couple gw to count-cache topology. etcd's Raft + watch semantics give us sub-second propagation with a sane consensus story for a small state size (<1 MB).
When it fails. etcd quorum loss: existing dataplane caches keep serving the last good config for 5 min, then expire — gradual degradation, not cliff. Stale hot-key flag: count-cache replicates a key that's no longer hot — wastes memory but doesn't break. Bad config push: a feature-flag panic mode is the same lever the SRE pulls.
- Schema RegistryConfluent Schema Registry (managed)
Stores versioned Avro/Protobuf schemas for every Kafka topic. write-svc validates outgoing events against the registered schema before publish; every consumer (aggregator, abuse) fetches the schema on first contact with a new schema_id and caches it.
Why it exists. Without a registry, the failure mode is: producer adds a new reaction field, a consumer does not recompile, that consumer silently drops or crashes — and you find out via the drift reconciler 24 h later. The registry's FORWARD compatibility mode enforces 'events written with the new schema must remain readable by consumers still on the old schema' at PUBLISH time, so the bad write is rejected at the source rather than 24 h later. (BACKWARD is the opposite guarantee — new consumers reading old data — and would not protect the stale consumer.)
When it fails. Registry 503 on publish: write-svc fails the publish — rather than ship an un-validated event we 503 the user; mobile SDK retries succeed once registry's back. Bad schema change (incompatible): registry rejects at the producer — the bad deploy doesn't even ship. Cache poisoning: consumer holds a bad schema for 60 min; fix is forced cache-flush via config-store kill-switch.
- Dead-Letter TopicKafka topic likes.dead-letter (32 partitions, 30d retention)
Receives events that aggregator / abuse couldn't process after N retries (typically: schema-incompatible payload, malformed UUID, downstream poison). Owned by the like-platform team; daily report of DLQ size goes to the team Slack. Runbook: inspect, fix root cause, optionally replay back into the source topic.
Why it exists. Without a DLQ, one poison message per partition stalls the entire partition forever — at 256 partitions of likes.committed, one bad event can stall 1/256th of all downstream processing indefinitely. A DLQ converts that into a quarantine the on-call can deal with on Monday morning rather than at 3am.
When it fails. DLQ topic itself unavailable: consumers fall back to logging the bad event + advancing offset (last-resort, alerted). DLQ grows unbounded: page when daily delta > 1000 events — usually means a producer / consumer version mismatch shipped without registry catching it.
- OTel CollectorOpenTelemetry Collector (gateway sampling) → Tempo
Receives OTLP spans from every service (gw, read-svc, write-svc, aggregator, abuse). Applies tail-based sampling (keep 100% of error traces + p99-latency outliers + 1% of healthy), buffers locally on disk for 30 s of upstream-down resilience, and ships to Tempo for query.
Why it exists. A 'like' touches ≥4 services on the sync path and ≥3 more on async; without distributed tracing, debugging a 'why was that 99.9th-percentile write 800 ms?' takes hours of log-grep. Tail-sampling at the collector keeps storage cost bounded while still capturing every error and outlier — which is what on-call actually opens during a page.
When it fails. Collector pod loss: spans buffer locally on app pods for ~30 s; longer outage → spans dropped (alerted on otlp_export_failures). Tempo backend down: collector buffers up to disk limit; sustained outage = we're blind to traces but the product keeps working. Bad sampling config: drop too much = blind; keep too much = storage cost explodes — guard rail on traces_kept_per_second.
- Metrics + AlertingPrometheus + Mimir (long-term) + Alertmanager
Scrapes every service for RED metrics (Rate, Errors, Duration) and every datastore for USE metrics (Utilization, Saturation, Errors). Mimir backs the historical store; Alertmanager evaluates SLO burn-rate rules and pages on-call via PagerDuty / Opsgenie.
Why it exists. Every SLI named in the observability section (redis_shard_qps_skew, consumer_lag_seconds, replication_lag_seconds, canary_drift_ratio, count_flicker_per_user, mm2_lag_seconds, traces_kept_per_second) is a metric that needs scraping, storing, and alerting. Without a metrics plane named in the diagram, those SLIs are imaginary.
When it fails. Prometheus pod loss: 15 s scrape gap; auto-replaced by k8s; Mimir backfills on next scrape — minimal data loss. Mimir storage corruption: short-term metrics survive in Prom local TSDB; long-term query degrades; recovery via S3 backup. Alertmanager misconfiguration: silent SLO breach — dead-man's-switch alert every 5 min as the canary.
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:
- Reaction types — just
like, or{like, love, haha, wow, sad, angry, care}? Affects the row size on the edge store, the aggregate payload, and the counter cardinality (7× per post). - Public vs private accounts — a private-account like is only visible to mutual followers; affects read ACLs but not the counter.
- Counts visible to the un-authenticated viewer? — yes typically (drives CDN cacheability).
- Un-like behavior — toggles state to
unliked, decrements counter. The double-tap UX is a real failure mode. - Anti-abuse approach — sync-block (rejects p99 budget) or async-void (commits then retroactively removes). Industry default = async.
- Multi-region posture — active-active (every region accepts writes) or active-standby (one writer)? Affects CRDT vs single-leader choice.
Assumptions to state:
- 2B likes/day → ~23K/s avg, ~93K/s diurnal peak, single-hot-key peak 50K/s (viral post).
- Read:Write 200:1 → ~4.6M/s avg, ~18M/s diurnal peak; CDN absorbs ≥80%, origin sized for ~3.6M/s.
- Reactions: 7 types (
{like, love, haha, wow, sad, angry, care}). - Multi-region active-active across 3 regions.
- 90-day raw event retention hot, then aggregates only.
- Anti-abuse async (commit then void).
02Functional reqs
What must this system actually do?
- Tap → like a post (or react with a specific emoji); idempotent on retry.
- Tap again → un-like / clear reaction; idempotent.
- Read a post's like count and top-3 reactions for the un-authenticated viewer.
- For an authenticated viewer, additionally return
hasLiked(user, post)so the heart renders correctly. - All counts approximate within ±0.5% steady-state; ±2% during cross-region partition.
03Non-functional
What must it promise about speed, uptime and correctness?
- Availability: 99.99% on the read path (reads are the product); 99.95% on the write path (write failure shows "tap again" — acceptable but tracked).
- Latency budgets: GET count p99 < 30 ms server-side (excluding CDN), p99.9 < 100 ms. POST /likes p99 < 150 ms. hasLiked p99 < 20 ms. Read-your-own-write end-to-end < 200 ms.
- Durability: A committed like (return 202) survives any single broker / replica loss. RPO 0 within-region, RPO ~5 s cross-region (during partition heal).
- Consistency: Read-your-write for the liker, eventually consistent for viewers. Counts approximate by design.
- Scalability: Horizontal on every tier. Single-key write rate up to 50K/s before key-split. Single-key read rate up to 5M/s with CDN absorption.
- Security: Rate-limit per-user + per-IP; bot-like scoring async; GDPR erasure within 30 days; no PII in the like-event payload beyond user_id (hashed).
04Capacity estimation
How much load and data does this have to hold?
We work the design from the published anchors, not from intuition.
Write QPS. 2B likes/day / 86,400 s ≈ 23K writes/s avg. Peak / avg = 4× → 93K writes/s diurnal peak. Viral hot-key peak (one post, observed on Twitter/Cristiano/MrBeast traffic shapes) is 50K writes/s on a single post_id for 10–20 min.
Read QPS. Read:write ratio = 200:1 (every viewer renders a count; only a small fraction click like). Avg reads = 23K × 200 = 4.6M/s avg; peak = 18M/s. CDN absorbs ≥80% (per Cloudflare / Fastly published hit rates for short-TTL cache-tag invalidatable objects), so origin sees ~3.6M/s peak spread across CDN-miss + hasLiked traffic.
Kafka spine sizing. Each committed like = 1 event. Fanout factor 2× (aggregator + abuse) means consumers process ~186K events/s peak at the read fanout. Producer ingest = peak writes = ~93K/s. 256 partitions × per-partition ceiling ~8K events/s = 2M events/s headroom (comfortable for hot-key skew).
Counter cache sizing. Hot working set = top 1% of posts (1T total posts × 1% = 10B keys). Aggregate count:{post} = ~32 B; sub-counters set ~2 KB. Top 1% in cache ≈ 10B × 2 KB ≈ 20 TB — too large. Reduce to top 0.1% (1B keys × 2 KB = 2 TB) across 64 shards × ~32 GB/shard, replicated 3×.
Counter store sizing. Each post × 3 regions × 64 shards × 7 reaction types = up to 1,344 cells (192 per reaction) × 48 B ≈ 64 KB per post fully populated. Counter cells materialize on first write, so the footprint is bounded by likes, not posts: 2B likes/day creates at most ~0.7T new cells/year ≈ 35 TB/year, and only viral posts ever fill all 1,344 — fits in Cassandra/Scylla at scale (Discord runs trillions of message rows on ScyllaDB).
Edge-store sizing. 2B likes/day × 80 B = 160 GB/day; 90-day hot retention = ~14 TB. After 90 days, drop raw rows and keep aggregate (user × month → liked_posts_count) summary. Sharded by user_id over 256 partitions.
Bandwidth. Peak read response = 18M/s × ~120 B = 2.2 GB/s; CDN absorbs ~80% → origin egress ~440 MB/s.
Architectural flip points (where a layer falls over).
| Component | Ceiling | Next tier |
|---|---|---|
| Single Redis INCR per post | ~100K ops/s | Per-shard sub-counters (this design uses 64) |
| Cassandra counter cell per post | ~5K ops/s | Per-(region, shard, reaction) cells (this design uses 192 per reaction) |
| Kafka partition (single) | ~8K events/s | Increase partition count (this design uses 256) |
| CDN edge cache (per POP) | ~50K rps/POP | Multi-POP + per-POP replication (Cloudflare default) |
05API design
What does the outside world call, and what comes back?
Anonymous variant — no Bearer. Shared-cacheable; never carries per-user fields.
GET /v1/posts/:post_id/likes
200 OK
Cache-Control: public, s-maxage=15, stale-while-revalidate=60
Cache-Tag: post:{post_id}:count
ETag: "gen-{generation_id}"
Vary: Authorization
{
"post_id": "p_abc",
"count": 18234,
"reactions": { "like": 14000, "love": 3200, "haha": 1000, "wow": 34 },
"generation_id": 184293,
"stale": false
}
Authenticated variant — Bearer present. Adds per-user hasLiked/myReaction; MUST NOT be shared-cached.
GET /v1/posts/:post_id/likes
Authorization: Bearer <token>
200 OK
Cache-Control: private, no-store
Vary: Authorization
{
"post_id": "p_abc",
"count": 18234,
"reactions": { "like": 14000, "love": 3200, "haha": 1000, "wow": 34 },
"hasLiked": true,
"myReaction": "love",
"generation_id": 184293,
"stale": false
}
The edge-bypass trigger keys on the Authorization header, not on ?u=: the anonymous variant — the only one a shared cache ever stores — carries Vary: Authorization, and the authed variant is private, no-store, so a CDN/proxy can neither store one viewer's hasLiked/myReaction nor replay the stored anonymous payload to a request bearing a token (public/s-maxage would otherwise permit exactly that reuse under RFC 9111). The anonymous variant is the only thing the CDN fans out.
POST /v1/likes
Authorization: Bearer <token>
Idempotency-Key: <user_id>:<post_id>:<intent_uuid>
{
"post_id": "p_abc",
"reaction": "like" # or "love"|"haha"|...|"clear"
}
202 Accepted
{
"event_id": "ev_xz192",
"state": "liked",
"ts": "2026-05-17T12:34:56.789Z"
}
409 Conflict # idempotency-key replayed with different body
Why 202, not 201: the like is durably committed to edge-store + Kafka, but the counter aggregation and anti-abuse scoring fan-out is async. 201 would imply "fully reflected everywhere" which is a lie.
06Data model
What gets stored, and what is it looked up by?
Like-edge store (Cassandra/Scylla, RF=3, sharded by user_id over 256 partitions):
| field | type | notes |
|---|---|---|
| user_id | uuid | partition key |
| post_id | uuid | clustering key |
| state | enum | liked / loved / ... / clear |
| reaction | tinyint | reaction type ordinal |
| ts | timestamp | event time |
| void | bool | set by anti-abuse (excluded from counts) |
| origin_region | tinyint | which region accepted the write |
| event_id | uuid | matches Kafka event_id for trace |
Counter store (Cassandra counter columns, partitioned by post_id):
| field | type | notes |
|---|---|---|
| post_id | uuid | partition key |
| region | tinyint | clustering key (0/1/2 for 3 regions) |
| shard | tinyint | clustering key (0..63) |
| reaction | tinyint | clustering key (0..6 — per reaction type) |
| count | counter | PN-Counter cell — region+shard owns its slot |
Read SELECT SUM(count) FROM counters WHERE post_id = ? for the global aggregate; group by reaction for the breakdown.
Idempotency store (Redis Cluster, 24 h TTL, sharded by user_id):
| key | value |
|---|---|
idem:{user_id}:{post_id}:{intent_uuid} | {state, event_id, ts} |
Count cache (Redis Cluster, sharded by post_id):
| key | value |
|---|---|
count:{post_id} | {count, reactions, generation_id, ts} (aggregate, 5 s TTL) |
count:{post_id}:s{0..63} | per-shard sub-counter ints |
count:{post_id}:hot | flag set by gateway sampler — read-replicates to 4 pools |
Kafka topic likes.committed:
| field | type |
|---|---|
| event_id | uuid |
| user_id | uuid |
| post_id | uuid |
| state | enum |
| reaction | tinyint |
| ts | timestamp |
| origin_region | tinyint |
| author_id | uuid |
07High-level design
Which components handle a request, and in what order?
The diagram (in the frontmatter) lays out:
- Read path (top lane, y≈120): client → CDN → gateway → read-svc → count-cache → fallback to counter-store (cache-miss branch).
- Write path (middle lane, y≈320): client → gateway → write-svc → {idem-store dedup} + {edge-store outbox write} + {Kafka outbox publish}.
- Async fanout (bottom lane, y≈480+): Kafka → aggregator → counter-store + count-cache invalidate; Kafka → abuse → publish likes.voided (compensating saga).
The four load-bearing patterns:
- Real outbox at the write — write-svc issues a Cassandra BATCH inserting the edge row AND an
outboxrow into the SAME user-partition (partition-local atomicity is the only ACID Cassandra offers and we lean on it deliberately). After commit, it publishes to Kafka; if that fails, a relay sweeper inside write-svc drains pending outbox rows within 5 s. The user's 202 returns on first successful publish; durability is guaranteed by the BATCH alone. Txn grouplike-commit, modeoutbox. - Per-shard counter cells — counter-store holds (post, region, shard) cells, summed on read. Spreads write load and survives cross-region partition via CRDT G-Counter merge.
- mcrouter-style hot-key replication on the cache — gateway publishes top-K hot post_ids to config-store; count-cache replicates those keys to 4 pools to absorb single-shard read peaks. For posts above the hot threshold, write-svc additionally sub-keys the Kafka partition (post_id + write_svc_pod_id % 8) so the upstream partition itself doesn't saturate.
- Saga compensation for anti-abuse — abuse-scorer publishes
likes.voidedAND writes the void flag back to edge-store. A SEPARATE txn grouplike-commit-compensate(notlike-commit) so the compensating saga is distinguishable from the original outbox commit — both reference the originalevent_idfor trace correlation.
08Deep dives
Which part breaks first, and what do you do about it?
1. Hot-key cascade — what actually breaks at 50K writes/s on one post.
Without sharding, every increment hits one counter cell and one cache key:
- Cassandra counter cells max out at ~5K ops/s due to internal coordination (per-row LWT-ish path); 50K/s would queue.
- Redis INCR on one key tops out at ~100K ops/s per shard; viable for writes but reads on the same key from a viral post (5M/s) crush the shard's CPU.
The fix is split + replicate:
- Writes: distribute across 64 sub-counter cells per (post, region). Each cell sees ~800 ops/s peak — well within ceiling.
- Reads: gateway samples QPS via a 30 s count-min sketch; top-K keys are pushed to the count-cache config which replicates them across 4 read pools. Reads pick a pool at random (mcrouter pattern).
Reference: Facebook's "Scaling Memcache at Facebook" (NSDI 2013) is the canonical write-up.
2. See-your-own-write across regions.
User likes from EU, switches to US WiFi 90 s later, opens the app. Three things have to work:
- The SDK has the optimistic state from the original tap (local SQLite).
- The first GET hits CDN — CDN doesn't know hasLiked, so it bypasses to origin for authed requests.
- Origin's read-svc looks up edge-store; cross-region replication is async ~1–5 s, so the EU-origin write may or may not be on the US replica.
- If hasLiked returns false but client SQLite says true → SDK trusts SQLite for 60 s (read-your-write grace window) and silently retries the lookup. After 60 s, server wins.
Cross-region edge-store replication uses Cassandra's DC-aware repl; the SDK's read-your-write token is a (user_id, ts) cookie that the read-svc uses to reject stale-replica reads when the lookup ts < cookie ts.
3. Double-tap toggle bug, and the subtler write-svc dedup contract.
User taps like, network is slow, taps again to un-like. SDK sends like → POST 1; before ack arrives, sends un-like → POST 2. Server processes in some order.
Each tap gets a fresh intent_uuid → different Idempotency-Key → server treats them as distinct intents. Order is determined by server arrival; the edge-store row's state is overwritten by the later write (last-writer-wins by ts). UX is correct: the final state matches the final tap. Counter sees +1 then -1 (PN-Counter handles this naturally).
If the user retries the SAME intent (network retry on a single tap), the SDK reuses the SAME Idempotency-Key → server dedups. No double-count.
The subtle bug — write-svc pod retry across DIFFERENT pods. Two retries of the same intent can land on two different write-svc pods if the gateway retries on 503. Both pods read empty from idem-store, both write to edge-store (UPSERT is idempotent so the row is fine), both publish to Kafka — with DIFFERENT event_ids but the SAME intent_uuid. If the aggregator dedups on event_id, the counter increments TWICE for one user intent. The contract is: aggregator dedups on intent_uuid (the field carried in every published event), not on event_id; the dedup set is a 24 h RocksDB-backed Bloom + exact-set with intent_uuid as key. This is documented in aggregator.keyChoices. Without it, idempotency leaks. The gateway-side fix is also live: e4 retryPolicy is now none (the SDK is the only retry layer for writes) — this removes the double-dispatch path entirely, with the consumer-side intent_uuid dedup as defense-in-depth.
4. Count drift and reconciliation.
A nightly job re-computes counts from likes.committed (replayed from Kafka or the S3 tiered archive) and compares to counter-store. Drift > 0.5% pages the on-call. Root causes seen in practice:
- CRDT cross-region merge bug — a Cassandra version mismatch dropped one region's cells silently. Fix: full repair, re-publish lost increments.
- At-least-once consumer without idempotency — counter aggregator restarted, replayed 30 s of events, double-counted. Fix: aggregator carries
last_seen_event_id_per_partitionin RocksDB state; dedupes on replay. - Queue drop — Kafka partition leader election dropped 2 in-flight increments. Fix: acks=all + min.ISR=2 (already in design); on suspicion, replay the affected partition window.
- Missed anti-abuse void — abuse scored a like as bot but the void publish 503'd and didn't retry. Fix: voids carry their own Idempotency-Key and retry until acked.
Reconciliation never SETs the counter — that would cause a visible jump. Instead it computes the delta and feeds it through the aggregator as a synthetic event with correction=true, which fades in over the next 60 s.
09Trade-offs
What did this design cost, and what breaks at 10×?
Accepted trade-offs, with the alternatives we rejected and why:
- Real outbox via Cassandra BATCH instead of full distributed-2PC OR naive dual-write. Rejected: (a) 2PC across Cassandra + Kafka (no production-grade coordinator survives node loss); (b) bare dual-write (Cassandra success then Kafka publish, lose the event if publish fails). Why: the BATCH writes both rows in the user-partition atomically (the only ACID Cassandra offers); the publish is post-commit relay; on publish failure the outbox row persists for the relay sweeper. This is the textbook outbox pattern adapted to Cassandra. A daily reconciler compares edge-row inserts to
likes.committedKafka offsets and re-publishes any missed events as defense-in-depth — explicitly designed to catch the case where the outbox relay itself fails.
- CRDT (eventual) counters instead of Spanner / single-leader linearizable. Rejected: Spanner. Why: linearizable cross-region writes mean every write pays cross-region latency (~80 ms for cross-Atlantic) — that's our entire p99 write budget. PN-Counter cells trade exact serial consistency for write-availability; counts are approximate by product spec, so the lost guarantee was free.
- Async anti-abuse instead of synchronous block. Rejected: sync ML scoring on the write path. Why: model inference is 20–50 ms; synchronous scoring would blow the 150 ms write p99. The cost is brief over-counting from bot likes — voids settle within a minute. Industry default.
- No write-through cache — async invalidate instead. Rejected: write-svc writes count-cache synchronously on each like. Why: writes are 23K/s; touching one cache key per write means 23K SETs/s per hot post — same hot-key problem we solved at the counter. Async invalidation via the aggregator coalesces.
- Outbox pattern instead of 2PC. Rejected: 2PC between edge-store and Kafka. Why: 2PC across Cassandra + Kafka has no production-grade coordinator that survives node loss. Outbox = local atomicity (edge-store row write) + post-commit relay (publish). At-least-once + idempotent consumers handle duplicates.
- Per-shard sub-counters instead of CRDT-only. Rejected: "just trust the CRDT, no sub-counters." Why: CRDT solves cross-region merge but not within-region hot-cell contention. Sub-counters spread the write load within a region; CRDT merges across regions. Both are needed.
- Edge-store sharded by user_id, not post_id. Rejected: post_id sharding (would be natural for counter joins). Why: the hot read pattern is hasLiked-for-this-feed-of-50-posts. User_id sharding makes that 50 same-shard lookups instead of 50 cross-shard scatters. Counts live in counter-store with their own sharding — separation of concerns.
- No write-side ML throttle (decided async-only). Rejected: per-user adaptive concurrency limit at write-svc. Why: adds another point of failure for legit users when an attacker has compromised an account; abuse is better handled async by the scorer, and the gateway rate-limit (100/min/user) is the only blunt instrument we keep on the sync path.
- Author push notifications and the "liked posts" tab are deliberately scoped out. Why: both are derived products, not the counting spine — each layers back as one more consumer group on
likes.committed(a notification dispatcher digesting per author-tier; a reverse-index builder materializing a per-user liked-posts view) without touching the write path, the outbox, or the counter. Keeping one derived consumer (the anti-abuse scorer) proves the fan-out pattern; the rest are the same shape.
Primary sources
- TAO: Facebook's Distributed Data Store for the Social Graph (USENIX ATC 2013)
- Scaling Memcache at Facebook (NSDI 2013)
- How Discord Stores Trillions of Messages (Discord Eng blog, 2023)
- Twitter — The Infrastructure Behind Twitter: Scale (Twitter Eng blog)
- PSY's Gangnam Style 32-bit overflow incident (2014)
- Pinterest Flink Counter Framework (Current 2025)
- Google SRE Workbook ch.5 — Alerting on SLOs (burn-rate alerts)
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 Like Button at Scale 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`.
- Reddit / Hacker NewsVote-driven ranking with hot/top/new at scale.
- View Count on a Video/PostDedup, bot-filter, batched aggregation.
- Trending TopicsSliding windows + Count-Min Sketch + top-K.