The L5 rubric & the opener script
What the round is actually scoring, and the 3-minute script that front-loads the graded part.
1. What the L5 bar looks like, signal by signal
| Signal | What clears the bar |
|---|---|
| Requirements | Pins down QPS, latency SLO, consistency, region geometry, internal vs external before any box. |
| Architecture | Every box defended: what breaks without it, what it costs. |
| Storage | Names the store and the access pattern that picked it (point lookup vs range scan vs append-heavy). |
| Failure | Per-dependency: what happens when it dies, fail-open vs fail-closed, and which is right here. |
| Scale | Hot-key and hot-shard analysis with a concrete mitigation. |
The single most-cited rejection reason in the reports is jumping to architecture without clarifying. Not depth. Not novelty. That one.
2. The opener โ say this out loud in every mock
Five questions, ~3 minutes
1. Who calls this โ internal services or end users? (decides auth model, SLA, abuse surface) 2. Scale? Peak QPS, total entities, growth rate. (decides sharding and whether one box works) 3. Latency SLO at p99, and is it on the critical path? (decides sync vs async, cache vs no cache) 4. Consistency: can a reader see stale data, and for how long? (the single biggest architecture fork) 5. One region or global? Where are the writers? (decides replication topology and conflict handling)
Then: "I'll assume X, Y, Z unless you want to steer me." Stating assumptions out loud is what converts silence into a grade.
Then the shape of the hour
- ~10 min clarify + capacity math out loud.
- ~10 min high-level boxes only. Resist diving.
- ~25 min deep dive โ the interviewer picks the box. Have the atomic sequence ready for the one that has a race in it.
- ~10 min failure & trade-off. Volunteer this; don't wait to be asked.
- ~5 min evolution: MVP โ 10ร, and where you'd re-architect.
Verbal-round adaptation: when there's no whiteboard, narrate structure explicitly โ "three tiers: edge, decision, storage; let me walk the write path first, then the read path." Without a diagram the interviewer can only grade the words, so the words have to carry the structure.
3. Reusable moves that score in almost every prompt
Say these unprompted
- Idempotency key on every retryable write. Names the double-write race before it's asked.
- Fail-open vs fail-closed, tied to the business consequence, not to taste.
- Hot key mitigation: shard the key (
key:{0..N}), or local pre-aggregate then flush. - Backpressure: what the caller sees when the queue is full.
429withRetry-Afterbeats silent drop. - Watermarks / bounded state for anything streaming, so state doesn't grow forever.
Capacity math you should have memorized
| Thing | Number |
|---|---|
| Redis point op, same DC | ~0.5โ1 ms p99, ~100k ops/s/core |
| SSD random read | ~100 ยตs; ~50โ100k IOPS/device |
| Cross-region RTT (USโEU) | ~80โ100 ms |
| Spanner cross-region commit | ~50โ150 ms |
| Kafka partition | ~10 MB/s sustained |
| 1 day at 100k QPS | ~8.6 B events |
| Object storage | ~$20/TB/mo hot, ~$4/TB/mo cold |
Distributed rate limiter full
L6 Infra/Platform onsite, SD round 1, 2026-01. "Design a large-scale distributed rate limiter service used internally by multiple teams." Interviewer pushed hard on trade-offs rather than novelty.
1. Requirements in one line
Functional
check(tenant, key, cost) โ ALLOW | DENYcalled on the request path by any internal service.- Rules configurable per tenant, per API, per caller identity; hierarchical (org โ team โ key).
- Multiple algorithms: fixed window, sliding window, token bucket with burst.
- Observability: per-rule allow/deny counts, near-real-time.
Out of scope: authn (assume identity is already resolved), billing, WAF/bot detection.
Non-functional
- 10โตโ10โถ decisions/sec aggregate, multi-tenant.
- p99 < 5 ms added latency โ it's on every request's critical path.
- Multi-region; a region partition must not take down callers.
- Availability > accuracy: 99.99%+.
Core tension: a limiter that is perfectly accurate must be strongly consistent and therefore globally coordinated, which is exactly the thing you cannot afford on the critical path. Every decision below is buying accuracy back cheaply.
2. Core entities & schema
Rule (config store, read-mostly, versioned)
rule_id, tenant_id, scope -- "api:/v1/search", "org:ads"
algo -- TOKEN_BUCKET | SLIDING_WINDOW
capacity, refill_per_sec, window_ms
action_on_breach -- DENY | SHADOW | DEGRADE
version, updated_at
Counter (hot path, Redis / in-memory, TTL'd)
key = {tenant}:{scope}:{subject}:{bucket_ts}
value = tokens (float) | count (int)
aux = last_refill_ts
TTL = 2 ร window -- self-cleaning, no GC job
Decision log (async, sampled, โ Kafka โ BigQuery-ish)
ts, rule_id, subject, allowed, remaining, node_id
Why Redis for counters, not the SQL config store
The counter is a read-modify-write at request rate with a TTL and no need for durability โ losing a counter costs one over-admitted window, not money. That is exactly Redis's shape: single-threaded per shard so INCR/Lua are atomic without locks, and TTL eviction means no cleanup job. Putting it in SQL costs a disk fsync per request and a row-lock hot spot on popular tenants. Conversely rules are low-volume, must survive a full cache flush, and need audit history โ so they live in a replicated SQL store and are pushed to limiter nodes, never read on the request path.
3. API interfaces
Hot path (called by every service)
POST /v1/check
{ tenant, scope, subject, cost=1 }
โ 200 { allowed: true, remaining: 41, reset_ms: 730 }
โ 200 { allowed: false, remaining: 0, retry_after_ms: 730 }
// caller translates false โ HTTP 429 + Retry-After
// batched variant amortizes RPC overhead:
POST /v1/check:batch { checks: [...] }
Return remaining and reset_ms always โ callers surface them as X-RateLimit-* headers, which is what stops clients from hot-looping into you.
Control path (rare)
PUT /v1/rules/{rule_id} { algo, capacity, ... }
GET /v1/rules?tenant=ads
POST /v1/rules/{id}:shadow // evaluate, never deny
DELETE /v1/rules/{rule_id}
GET /v1/stats?rule_id=..&window=5m
โ { allowed, denied, p99_check_ms, top_subjects[] }
Shadow mode is the feature that gets you the offer. Nobody ships a new limit straight to DENY โ you run it in shadow, look at what would have been denied, then flip.
4. Architecture
Two paths: the common one never leaves the sidecar (local tokens pre-fetched in bulk); the miss path does one Redis Lua round-trip. Everything else โ rules, stats โ is off the critical path by construction.
5. Deep dive โ the atomic decision
Token bucket in one Lua script
-- KEYS[1] = bucket key
-- ARGV = now_ms, capacity, refill_per_ms, cost
local b = redis.call('HMGET', KEYS[1], 'tk', 'ts')
local tokens = tonumber(b[1]) or ARGV[2]
local ts = tonumber(b[2]) or ARGV[1]
local delta = math.max(0, ARGV[1] - ts)
tokens = math.min(ARGV[2], tokens + delta * ARGV[3])
local ok = tokens >= tonumber(ARGV[4])
if ok then tokens = tokens - ARGV[4] end
redis.call('HMSET', KEYS[1], 'tk', tokens, 'ts', ARGV[1])
redis.call('PEXPIRE', KEYS[1], 2 * ARGV[5])
return { ok and 1 or 0, tokens }
CAPACITY = burst allowance; refill_per_ms = steady-state rate. Two knobs, and being able to say which one a customer complaint maps to is the whole point of picking this algorithm.
Why it's correct, and where it isn't
Redis executes a Lua script atomically on the shard, so read-refill-decrement-write cannot interleave. That kills the classic race where two concurrent GET/SET pairs both see 1 token left and both allow.
Defend it: the script must be deterministic โ never call TIME inside it, always pass now_ms from the caller, or replicas diverge from the primary. That in turn means clock skew across limiter nodes leaks quota: a node 500 ms fast refills early. Bound it with NTP + reject now_ms more than ~1 s off the shard's own clock.
Why not sliding-window-log (a sorted set of timestamps per key)? Exact, but memory is O(requests in window) โ a 10k-QPS tenant on a 60 s window is 600k members in one key. Token bucket is O(1) memory and the accuracy loss is a bounded burst.
6. Deep dive โ hot keys and hierarchical limits
The hot-key problem
Rate limiter traffic is maximally skewed by design โ the abusive tenant is the one generating the most checks, and all of its checks hash to one Redis shard. That shard saturates and you have taken down limiting for every tenant on it.
Fix 1 โ key sharding:
key = {tenant}:{scope}:{h(subject) % N}
each sub-bucket gets capacity/N
cost: burst accuracy degrades ~Nร
(a client can get N bursts)
Fix 2 โ local pre-allocation (preferred):
sidecar leases 100 tokens at a time
spends locally, refreshes at 20% left
cost: up to (leases ร lease_size) over-admit
on abrupt traffic drop; leases TTL out
Hierarchical limits
Real rules nest: org 100k/s, team 10k/s, key 1k/s. Naรฏvely you evaluate all three, which triples Redis load and creates a partial-decrement bug โ decrement org, then deny at team, and org's tokens are gone.
Fix: evaluate all levels in one Lua script, cheapest/narrowest first, and only commit decrements if every level passes. Same-slot placement via hash tags: {tenant}:org, {tenant}:team:x โ the braces force Redis Cluster to co-locate them so a multi-key script is legal.
Belt-and-suspenders: if levels genuinely cannot be co-located, decrement narrowest-first and refund on a later-level deny. Refunds are idempotent (add back cost, cap at capacity) so a lost refund costs one wasted token, not a stuck bucket.
7. Deep dive โ multi-region
Don't replicate counters synchronously
A globally-accurate limit needs consensus per decision: ~100 ms cross-region RTT against a 5 ms budget. Non-starter. Three options, in order of what you should actually propose:
| Model | Accuracy | Cost |
|---|---|---|
| Per-region quota split (limit/R per region) | Never over-admits globally; under-admits on skew | Free |
| Split + periodic rebalance (every ~10 s, redistribute unused) | Good | One background job |
| Global consensus (Spanner counter) | Exact | Blows the latency SLO |
What to say
"I'd start with static per-region split because it's strictly safe โ the global limit is never exceeded. The failure mode is that a region with 90% of traffic gets throttled while others sit idle, so I'd add an async rebalancer that redistributes unused quota on a 10-second loop. That converts a correctness problem into a utilization problem, which I'd rather have."
Defend it: when a region is partitioned, its quota share is stranded. The rebalancer must treat a silent region as still holding its share (don't redistribute), or you over-admit globally the moment it comes back. Conservative on partition, aggressive on healthy โ that asymmetry is the whole trick.
8. Follow-ups โ answers to have ready
Redis dies. Now what?
Fail-open, loudly. This is an internal limiter protecting against accidental overload, not a security control โ failing closed converts a cache outage into a total outage of every calling service, which is strictly worse than briefly unlimited traffic. Concretely: sidecar serves from its last local lease, then admits everything, emits a limiter_degraded metric, and pages. The exception is any rule tagged security (login attempts, password reset) โ those fail closed, and the rule schema carries that flag so the decision is made at config time, not during the incident.
Two limiter nodes both serve the same tenant. Do they double-admit?
No, because neither node holds state โ the counter lives on one Redis shard and the Lua script is atomic there. The only double-admit is via the local-lease optimization, and it's bounded: at most nodes ร lease_size extra requests, which you tune by shrinking leases for tight limits and growing them for loose ones.
A tenant complains they're getting 429s below their stated limit.
Three usual causes, in order of likelihood. (1) Fixed-window boundary effect โ they sent 2ร limit across a window edge on a previous design; sliding/token-bucket fixes it. (2) Key sharding split their capacity N ways and their traffic isn't uniformly distributed across subjects. (3) Clock skew on one limiter node making refill run slow. The remaining/reset_ms in every response plus the sampled decision log is what lets you answer this in minutes rather than guessing.
How do you roll out a new limit without breaking a team?
Shadow mode: evaluate, log, never deny. Look at the would-be-denied distribution for a week, notify the top offenders, then flip to enforce at a generous capacity and ratchet down. Rules are versioned so rollback is a config push, not a deploy.
Why not just do this at the load balancer?
You should, for the crude global limits โ L7 LB rate limiting is cheaper and closer to the client. But it can't see application-level identity (which tenant, which API key, which cost weight), can't do hierarchical limits, and can't be reconfigured per-team without touching shared infra. Layered: LB for the blunt DDoS ceiling, this service for the semantic limits.
Requests have different costs. Does that break the model?
No โ cost is already a parameter, so an expensive search costs 10 tokens and a health check costs 0. The subtlety is that you don't know the true cost until after execution. Reserve an estimated cost up front, then reconcile with an async adjust call. Over-reserving is safe; under-reserving lets one expensive request slip, which is acceptable.
9. Numbers to drop
Load
- 10โถ checks/sec peak aggregate.
- With sidecar leases at 90% local hit rate โ 10โต Redis ops/sec reaching the cluster.
- Redis shard sustains ~100k ops/s โ ~4 shards + replicas covers it with headroom. That's a startlingly small cluster, and saying so is the point.
- Limiter service: ~5k checks/s/pod โ ~200 pods, stateless, trivially autoscaled.
Storage & latency
- Counter entry ~100 B; 50 M active keys โ ~5 GB. Fits in memory on one shard set; TTL keeps it flat.
- Rules: 10โต rules ร 1 KB = 100 MB. Fits in every node's memory โ that's why you push, not fetch.
- Budget: local hit <0.1 ms; Redis round-trip ~1 ms p99; total p99 under 5 ms with room to spare.
- Decision log at 1% sampling โ 10k events/s โ ~1 Kafka partition-pair. Full logging would be 100ร and isn't worth it.
10. 30-second recap
Callers hit a sidecar that holds a small lease of tokens locally, so ~90% of checks never leave the process. On a miss it calls a stateless limiter service, which runs one atomic Lua script against a sharded Redis Cluster โ read tokens, refill by elapsed time, decrement if sufficient, all in one shot, so there's no read-modify-write race. Rules live in a versioned SQL store and are pushed to nodes, never fetched on the hot path. Multi-region is a static quota split with an async rebalancer, because global consensus per decision blows the 5 ms budget. If Redis dies we fail open and page, since this protects against accidental overload rather than attackers โ except for rules explicitly tagged security, which fail closed. The bounded inaccuracy is nodes ร lease_size extra admits, and I'd shrink leases for tight limits to trade Redis load for precision.
Global real-time notification system full
L6 onsite SD round 2, 2026-01. "Design a global real-time notification system โ push + email + SMS." Stated constraints: 10โธ+ users, multi-channel fallback, latency SLA, dedup and idempotency. The poster's read: this round tested system-owner thinking more than the infra round did.
1. Requirements in one line
Functional
- Producer services submit a notification event; the platform decides channel, renders, delivers, retries.
- Channels: mobile push (APNs/FCM), email, SMS, in-app inbox. Fallback chain per event type.
- User preferences and quiet hours honored; unsubscribe is legally binding for email/SMS.
- Fan-out: one event may target one user, a segment, or all users (announcement).
- Exactly-once user-visible delivery: at-least-once transport + dedup at the edge.
Out of scope: content authoring UI, ML send-time optimization, marketing campaign scheduling.
Non-functional
- 10โธ users; steady ~50k notifications/s, burst 1 M/s on a broadcast.
- Transactional (2FA, security alert): p99 < 5 s end-to-end. Marketing: minutes are fine.
- Durability > latency for transactional โ never silently drop.
- Multi-region, survives loss of one region and of any single third-party vendor.
Core tension: the burst is 20ร steady and the third-party channels (APNs, an SMS vendor) are both rate-limited and unreliable. So the system's real job is absorbing bursts and metabolizing partner failure, not sending messages.
2. Core entities & schema
NotificationEvent -- what the producer asked for (immutable)
event_id (idempotency key, producer-supplied)
type -- ORDER_SHIPPED | SECURITY_ALERT | PROMO_X
audience -- {user_id} | {segment_id} | ALL
payload -- template vars, NOT rendered text
priority -- TRANSACTIONAL | STANDARD | BULK
created_at, ttl_sec
Delivery -- one row per (event, user, channel). The unit of work.
delivery_id = hash(event_id, user_id, channel) <-- dedup key
state -- PENDING|SENT|DELIVERED|FAILED|SUPPRESSED
attempt, next_attempt_at, last_error
vendor_msg_id
UserPrefs -- read-heavy, cached hard
user_id, channel_enabled{}, quiet_hours, locale, tz
device_tokens[] (token, platform, last_seen)
suppression: unsubscribed_types[], hard_bounce, complaint
Why split Event from Delivery โ and why this store
Modelling only "notification" is the mistake that sinks this round. One event fans out to N users ร M channels, each with its own retry clock and terminal state; without a Delivery row you cannot answer "did user X get it?" or retry one channel without re-sending the others. It also makes dedup a primary-key property: delivery_id is a deterministic hash, so a duplicated event replays into the same row instead of a second send.
Delivery is high-volume, append-then-update-once, always accessed by key or by (state, next_attempt_at) โ that's a wide-column store (Bigtable/Cassandra) with a TTL, not a relational DB; there are no joins and no transactions across users. UserPrefs is small (10โธ ร ~1 KB = 100 GB), read on every delivery, and must be correct for unsubscribes, so: replicated SQL as source of truth, aggressively cached, with suppression checked again at send time.
3. API interfaces
Producer (write path)
POST /v1/notify
Idempotency-Key: {event_id}
{ type, audience: {user_id | segment_id},
payload: {...}, priority, ttl_sec }
โ 202 { event_id, status: "ACCEPTED" }
POST /v1/notify:broadcast
{ type, segment_id, payload, rate_limit_per_sec }
โ 202 { event_id, estimated_recipients }
202, never 200. Accepting into a durable log and returning immediately is what decouples the producer from APNs being down. A producer that blocks on delivery has coupled its own availability to a third party.
Consumer & ops (read path)
GET /v1/inbox?user_id=&cursor= // in-app feed
POST /v1/deliveries/{id}:ack // client read receipt
GET /v1/events/{event_id} // fan-out status rollup
โ { total, sent, delivered, failed, suppressed }
POST /v1/webhooks/vendor/{vendor} // bounce/delivery callbacks
PUT /v1/prefs/{user_id}
POST /v1/unsubscribe { token } // one-click, no login
4. Architecture
Green is the durable happy path: accept โ log โ fan-out โ per-channel queue. Amber is the recovery loop that everything eventually falls into. The in-app inbox is written unconditionally, so there is always one channel that cannot fail.
5. Deep dive โ fan-out and the broadcast burst
Two fan-out modes
Targeted (99% of events, 1..1000 users):
fan-out inline in the worker
write Delivery rows in a batch
enqueue per channel
Broadcast (segment or ALL, 10^8 users):
DO NOT expand in one worker
1. materialize segment -> sharded user-id ranges
2. emit N "fan-out chunk" tasks (chunk = 10k users)
3. each chunk expands independently, resumable
via (event_id, chunk_id) checkpoint
4. token-bucket the emit rate per channel
A broadcast at 10โธ users writing 10โธ Delivery rows is ~100 GB and hours of vendor throughput. It must be a paced background job, not a spike.
Why chunking, and where it still hurts
Chunk-level checkpoints make fan-out resumable and idempotent: a worker crash replays one 10k chunk, and the deterministic delivery_id means the replayed rows collide with the originals instead of duplicating. Without this, a crash 80% through a 10โธ fan-out leaves you with no safe action.
Defend it โ the priority inversion. A broadcast dumped into the same queue as 2FA codes will bury them behind hours of backlog, and a 2FA code delivered in 40 minutes is worse than useless. Fix: separate physical queues per priority class, not just a priority field. Transactional gets its own partitions, its own worker pool, and a reserved share of vendor rate budget that bulk traffic can never borrow. This is the single most likely follow-up in this round; volunteer it.
6. Deep dive โ exactly-once as the user perceives it
The chain of idempotency
1. Producer supplies event_id (Idempotency-Key).
Ingest does INSERT ... IF NOT EXISTS.
Retry of the same HTTP call -> same event, no dup.
2. delivery_id = hash(event_id, user_id, channel)
Fan-out replay writes the same row. Dedup by PK.
3. Send is guarded by a conditional state transition:
UPDATE Delivery
SET state='SENDING', attempt=attempt+1
WHERE delivery_id=? AND state IN ('PENDING','FAILED')
Only the winner calls the vendor.
4. Vendor call carries its own idempotency token
where supported (SES, Twilio); where not, accept
at-least-once and dedup on the DEVICE by event_id.
Why true exactly-once is impossible, and the honest answer
Step 3 has an unavoidable gap: you set SENDING, call the vendor, and crash before recording the result. You cannot know whether the vendor sent it. Any system that claims exactly-once across a third-party boundary is lying.
So decide per channel, by consequence. Push and in-app: retry (duplicate push is mildly annoying, and the client dedups on event_id anyway). Email: retry (dup email is tolerable). SMS: do not blind-retry โ it costs real money and duplicate 2FA codes confuse users and can invalidate the first code. For SMS, reconcile against the vendor's delivery-status API before a retry, and cap at one reconciliation attempt.
Belt-and-suspenders: the client-side dedup on event_id with a 24 h local cache is what actually delivers the user-visible guarantee. Server-side you only ever promise at-least-once.
7. Deep dive โ retries, fallback, and vendor failure
Retry policy that doesn't amplify an outage
backoff = min(2^attempt * 1s, 15m) * jitter(0.5..1.5)
max_attempts = 5, then DLQ
hard stop at event.ttl_sec -- a 2FA code has ttl 300s;
-- retrying at t+10m is harmful
Circuit breaker per (channel, vendor):
error_rate > 50% over 30s -> OPEN
OPEN: stop calling, park in queue, alarm
half-open probe every 10s
Jitter is not optional. Without it, a vendor recovering from a 60 s outage gets every parked message at the same instant and immediately falls over again.
Channel fallback โ and why it's usually wrong
The prompt asks for multi-channel fallback, so propose it, then bound it. Fallback chain lives on the event type: SECURITY_ALERT: [push, sms, email], PROMO: [push] (never escalate marketing to SMS โ it costs money and generates complaints).
Defend it: escalating on "not delivered within N seconds" is dangerous, because push delivery receipts are unreliable and slow โ you'll double-notify constantly. Only escalate on a hard signal: no registered device token, token rejected by APNs, or circuit breaker open for that channel. Timeout-based escalation should have a generous floor (minutes, not seconds) and must be off for anything bulk.
Vendor redundancy: two vendors per paid channel, weighted routing, automatic weight shift when a breaker opens. This is the concrete answer to "SMS vendor dies" and it's cheap to say.
8. Follow-ups โ answers to have ready
APNs goes down for an hour. What does the user see?
Nothing missing, because the in-app inbox row is written at fan-out time regardless of push outcome โ open the app and it's there. Push deliveries park behind an open circuit breaker; the retry scheduler holds them until TTL. When APNs recovers, the half-open probe reopens the valve and the backlog drains with jitter. Anything whose TTL expired in the meantime is dropped to DLQ and counted โ and for security alerts specifically, the fallback chain escalates to SMS on breaker-open rather than waiting.
Kafka loses a partition mid-fan-out.
Events are durably committed before the 202, and consumer offsets only advance after Delivery rows are written, so a partition failover replays from the last committed offset. Replay is safe precisely because delivery_id is deterministic โ re-processed chunks collide on the primary key instead of duplicating. The visible cost is latency during failover, not lost notifications. Retention of 7 days also means we can rewind and replay a window if we ship a fan-out bug.
A hot user โ say a celebrity account triggering millions of follower notifications.
That's the broadcast path, not the targeted path, and the trigger should classify it as such above a threshold (say >10k recipients). Practically: materialize the follower list into chunks, pace the emit with a token bucket, and โ importantly โ collapse. If the same producer generates 50 events for one user in a minute, digest them into one notification rather than sending 50. Collapse keys are the standard tool (collapse_key = (user_id, type), keep latest), and APNs/FCM support it natively.
User unsubscribes. How fast does that take effect, and what if a fan-out is in flight?
Suppression is checked twice: once at fan-out (cheap, from cache) and again in the sender immediately before the vendor call (authoritative read). The second check exists exactly for the in-flight case โ a broadcast queued an hour ago must not send to someone who unsubscribed ten minutes ago. For email/SMS this isn't a nicety; it's CAN-SPAM/TCPA exposure, and saying that out loud signals you've shipped this before.
How do you handle quiet hours across time zones?
Store the user's tz and quiet window in prefs; at fan-out, transactional messages ignore quiet hours entirely, and non-transactional ones get next_attempt_at set to the end of the quiet window rather than being dropped. Consequence: a global broadcast doesn't go out at one instant, it rolls around the planet over ~24 h, which you should state explicitly because it changes the "burst" math from 1 M/s to something far gentler.
How do you know delivery is actually working?
Per (type, channel, vendor, region): submitted โ sent โ delivered โ opened funnel, with alerting on ratio shifts, not absolute counts. The signal that catches real incidents is delivered/sent dropping while sent stays flat โ that means the vendor is accepting and silently dropping, which no error rate would show. Plus a synthetic canary: send to a handful of owned test accounts on every channel every minute and alert on the round-trip.
Region loss.
Producers are geo-routed and events are replicated cross-region in the log. Fan-out and senders are active-active with the Delivery store as the arbiter โ the conditional PENDING โ SENDING transition prevents two regions from sending the same delivery, assuming the store does per-key linearizable writes (Spanner or a quorum config). If it can't, accept the duplicate and rely on client-side dedup; say which one you're assuming.
9. Numbers to drop
Throughput
- Steady 50k notifications/s โ ~4 B/day. At ~1 KB of Delivery row each, ~4 TB/day before TTL.
- 30-day TTL on Delivery โ ~120 TB. Wide-column with compression, ~20 nodes. Cheap.
- Kafka: 50k msg/s ร 1 KB = 50 MB/s โ ~10 partitions at 10 MB/s, call it 32 for headroom and priority isolation.
- Broadcast to 10โธ: at a paced 100k sends/s that's ~17 minutes per channel โ quote this number, it reframes "real-time" honestly.
Latency budget (transactional, 5 s p99)
- Ingest + commit to log: ~20 ms
- Log โ fan-out worker pickup: ~100 ms
- Prefs lookup (cached): ~2 ms
- Delivery write + queue: ~30 ms
- Vendor call (APNs): ~200 ms p99
- Device receipt: ~1 s, outside our control
- Slack: ~3.6 s โ which is what pays for exactly one retry inside the SLA.
10. 30-second recap
Producers POST an event with an idempotency key and get a 202 as soon as it's durable in a partitioned log โ nothing blocks on a third party. Fan-out workers expand the audience, apply prefs and suppression, and write one Delivery row per (event, user, channel) with a deterministic ID, which is what makes the whole pipeline replay-safe. Deliveries go into per-channel, per-priority physical queues so a 10โธ-user broadcast can't bury a 2FA code. Senders take a conditional PENDINGโSENDING transition before calling the vendor, so only one worker sends; retries use exponential backoff with jitter behind a per-vendor circuit breaker, and hard-stop at the event TTL. Fallback escalates on hard signals only โ no device token, or breaker open โ never on a soft delivery timeout. The honest limit is that exactly-once across a vendor boundary is impossible, so we promise at-least-once server-side and dedup on the device by event ID, and we don't blind-retry SMS because it costs money.
Anomaly / metrics monitoring system full
L6 GCP phone-screen SD, 2026-05 (reported as "SD: ๅผๅธธ็ๆต็ณป็ป"). The question bank frames it as deliberately under-specified โ clarify the domain first (metric stream? user behavior? infra alerts?). The expected reading: a metrics platform that ingests, stores, detects, alerts, and visualizes time series.
1. Requirements in one line
Functional
- Ingest multi-dimensional time series:
metric{service,region,host,env,segment} โ value@ts. - Query: aggregate over label selectors and time ranges, for dashboards and detectors.
- Detection, at least two mechanisms: static threshold and dynamic baseline (seasonality / change-point).
- Alerting: routing by team, dedup, grouping, silencing, escalation, multiple channels.
- Visualization: dashboards, plus the ability to slice by dimension after an alert fires.
Out of scope: log search, distributed tracing, incident management workflow (page a PagerDuty-equivalent, don't build it).
Non-functional
- 10โท active series, 10โท datapoints/s ingest.
- Detection latency < 60 s from event to page for critical metrics.
- Query p99 < 1 s for a 24 h dashboard panel.
- The monitoring system must not depend on the systems it monitors, and must degrade before it lies.
Core tension: alert quality is a precision/recall trade-off, and it is a product decision, not a modelling one. A detector that never misses will page so often that humans stop reading, at which point recall is effectively zero. Everything below is about being deliberately less sensitive in exchange for being believed.
2. Core entities & schema
Series -- identity is the full label set
series_id = hash(metric_name, sorted(labels))
metric_name, labels{}, type (gauge|counter|histogram)
first_seen, last_seen
Sample -- the hot, huge one
series_id, ts, value
stored COLUMNAR + delta-of-delta on ts,
XOR/Gorilla on value -> ~1.4 bytes/sample
Detector
detector_id, selector (PromQL-ish), algo,
params{threshold | sensitivity | season_len},
for_duration, -- must hold N min before firing
severity, route_to (team), runbook_url
AlertInstance -- the dedup unit
fingerprint = hash(detector_id, grouping_labels)
state -- PENDING|FIRING|RESOLVED|SILENCED
started_at, last_eval, value, notification_count
Why a purpose-built TSDB and not "just Bigtable"
Three properties make time series a special case, and naming them is the point of this section. (1) Writes are append-only and time-ordered, so you never update a sample โ that kills the need for compaction-heavy general stores. (2) Adjacent values are highly correlated, so Gorilla-style XOR + delta-of-delta compression gets you ~1.4 bytes per sample versus ~16 raw; at 10โท points/s that's the difference between 500 TB/year and 4 PB/year. (3) Queries are range scans over one series, so the physical layout must be series-major, time-minor โ the exact opposite of what you'd get by keying on timestamp first (which also creates a write hot spot on "now").
Concretely: shard by hash(series_id) so writes spread evenly and one series' history is contiguous. An inverted index maps label=value โ [series_id] so a selector resolves to a series set before touching sample data.
3. API interfaces
Write path
POST /v1/ingest (protobuf, batched, gzip)
{ samples: [{ labels{}, ts, value }, ... ] }
โ 204, or 429 with Retry-After when over quota
// agents push every 10s; server does NOT pull.
// per-tenant series-cardinality quota enforced HERE.
Push, not pull, at this scale โ pull requires the monitoring system to hold a service-discovery view of 10โท targets and creates a fan-in storm. Pull's advantage (you know when a target is missing) is recovered by a staleness detector.
Read / control path
GET /v1/query_range
?selector=http_errors{service="ads",env="prod"}
&start=&end=&step=60s&agg=sum by (region)
โ { series: [{labels, points:[[ts,v]...]}] }
PUT /v1/detectors/{id}
POST /v1/silences { matcher, until, reason }
GET /v1/alerts?state=FIRING
POST /v1/detectors/{id}:backtest { window: "30d" }
Backtest is the feature that makes detectors shippable โ replay a candidate detector over 30 days of history and show how many times it would have paged. Nobody should tune a threshold in production.
4. Architecture
Two detection paths deliberately: a streaming one on the raw feed for the <60 s threshold alerts, and a batch one reading the TSDB for seasonal baselines that need hours of context. Both converge on one alert manager, which is the only thing allowed to page.
5. Deep dive โ the two detectors
Static threshold (streaming)
for each eval tick (15s):
v = agg(selector, last 5m)
if compare(v, threshold):
pending[fp] = pending[fp] or now
if now - pending[fp] >= for_duration: # e.g. 5m
emit FIRING
else:
clear pending[fp]; emit RESOLVED if firing
for_duration is the anti-flap knob, and it's the whole reason threshold alerting is usable. A metric crossing a line for 15 s is noise; crossing it for 5 minutes is an incident. Cost: you've added for_duration to your detection latency, so critical detectors run a short duration on a tight threshold and a long duration on a loose one.
Dynamic baseline (batch)
Seasonal decomposition per series:
baseline(t) = median over last K weeks
at same (weekday, time-of-day)
band(t) = baseline(t) ยฑ k * MAD(t)
fire when actual outside band for for_duration
Change point (fast, cheap):
compare rolling mean/var of last 5m vs prior 60m
robust z = (m1 - m0) / (1.4826 * MAD_0)
|z| > 4 -> candidate
Median + MAD, not mean + stddev. Mean and standard deviation are wrecked by the very outliers you're trying to detect โ one past incident inflates the band and hides the next one. MAD is robust to ~50% contamination, which is the difference between a detector that degrades gracefully and one that silently stops working after its first real incident.
Defend it โ why not "just use ML"
Reach for a learned model only when you can state what it buys. Here it buys multivariate correlation (error rate up and latency up and only in one region) at the cost of: training pipelines, a cold-start period per new service, opacity at 3 a.m. when the on-call asks why it fired, and a whole new failure mode where the model drifts and nobody notices. Seasonal median + MAD covers most of the value at a fraction of the operational cost. Propose ML as a ranking layer over already-fired alerts (which of these 40 alerts is the root cause), not as the trigger โ that framing is far more defensible than "add an LSTM."
6. Deep dive โ cardinality, the actual failure mode
How these systems really die
Not from datapoint volume โ from series count. Every distinct label combination is a new series with its own index entry and memory-resident head block. One engineer adds user_id or request_id or a raw URL as a label and cardinality goes from 10โท to 10โน overnight, the index blows out of RAM, and the whole cluster falls over. This takes down monitoring for everyone, during whatever incident prompted them to add the label.
Defenses, in order:
1. Per-tenant series quota, enforced at INGEST.
Over quota -> reject new series, keep existing.
(Never reject existing: partial data is worse
than no new dimensions.)
2. Label-value cardinality cap per label name.
>1000 distinct values -> auto-drop label, warn.
3. Cardinality explosion detector on the metadata
itself: d(series_count)/dt alert.
4. Chargeback: teams see their series cost.
The other structural risks
Hot shard: hashing by series_id spreads writes evenly by construction, but a query for sum by (region) over a million series fans out to every shard. Mitigate with pre-computed recording rules โ materialize the common aggregations at write time so dashboards read one cheap series instead of a million expensive ones.
Alert storm: one bad deploy fires 400 detectors at once and pages a human 400 times. The alert manager must group by a shared label set (e.g. all alerts for service=ads,region=us-east become one notification with 400 lines) and apply an inhibition rule โ if ServiceDown is firing, suppress every dependent HighLatency alert for that service. Without inhibition rules the system is technically correct and operationally useless.
Silent failure: the worst outcome isn't a false page, it's no page because ingest stopped. Hence the dead-man's switch below.
7. Deep dive โ not lying, and not depending on yourself
Staleness > absence
An absent series and a healthy series look identical to a threshold detector: neither crosses the line. So every detector needs an explicit staleness clause โ absent_over_time(selector[5m]) fires its own alert. Otherwise a crashed agent reads as "all clear," which is the single most dangerous bug a monitoring system can have.
Dead-man's switch: a synthetic detector that fires continuously and routes to an external service which pages you when the heartbeat stops. It's the only construct that catches "the monitoring system itself is down," and it costs about ten lines.
Circular dependency
If the monitoring stack runs on the same Kubernetes control plane, the same service mesh, and the same object store as production, then a production outage blinds you exactly when you need sight. State this explicitly: the meta-monitoring stack runs in a separate failure domain โ different region, different account, minimal dependencies, and its notification path does not traverse the primary network. It monitors ~50 series about the main system rather than 10โท, so it can be small and boring.
Degrade before lying: under ingest overload, shed low-priority tenants and keep critical ones at full fidelity, rather than uniformly downsampling everyone. Uniform degradation quietly reduces detection sensitivity across the board with no signal that it happened.
8. Follow-ups โ answers to have ready
Kafka/WAL backs up. What gives?
Ingest applies backpressure with 429s and agents buffer locally (bounded, ~15 min, then drop oldest). Critically, the streaming detector reads from the same log, so a backup delays detection โ that's the real cost, not lost data. So: prioritized partitions, with detector-relevant series on a reserved partition set that bulk dashboard-only metrics can't saturate. And the lag itself is a first-class alert on the meta stack.
Someone asks for a 2-year retention dashboard.
Not at raw resolution โ 10โท points/s for 2 years is ~900 TB compressed and no one needs 10-second granularity from 18 months ago. Downsample on compaction: raw for 15 days, 5-minute rollups for 90 days, 1-hour rollups for 2 years, each with min/max/sum/count so you can still reconstruct averages and see spikes. The query engine picks resolution from the requested range automatically. Say the trade-off out loud: you permanently lose the ability to see a 30-second spike from last year.
How do you stop alert fatigue?
Treat it as the primary metric of the system, not a side effect. Track per-detector precision (fired โ was there an actual incident?) via a required post-alert disposition, and auto-quarantine detectors below ~30% precision into a non-paging channel until the owner fixes them. Combine with grouping, inhibition, and for_duration. The organizational move matters as much as the technical one: pages should have a runbook link, and a page with no runbook shouldn't be allowed to page.
Two regions both evaluate the same detector. Double page?
Detector evaluation is sharded by fingerprint with a lease, so exactly one evaluator owns a detector at a time; on lease expiry another takes over. Even if two do evaluate during a partition, the alert manager dedups on fingerprint and the notification has its own dedup window, so the human sees one page. Duplicate evaluation is cheap and safe; duplicate notification is what you must prevent, and you prevent it at the last hop.
A metric is anomalous but it's Black Friday.
This is why the baseline is seasonal by (weekday, time-of-day) rather than a flat threshold โ but a once-a-year event isn't in the seasonal profile. Two mechanisms: scheduled silences tied to known events, and a relative comparison against a control group (this region vs all regions, this version vs previous). Relative detectors survive traffic-shape changes that absolute ones can't, which is a strong point to volunteer unprompted.
How would this evolve at 10ร?
The ingest and TSDB tiers scale horizontally by series hash, so 10ร is more shards โ boring, which is the point. The parts that don't scale linearly are the inverted index (a selector matching millions of series gets slow, so push more work into recording rules) and the human alert budget, which doesn't scale at all. At 10ร the re-architecture is on the alerting side: correlate and rank rather than route more pages.
9. Numbers to drop
Ingest & storage
- 10โท series ร 1 sample / 10 s = 10โถ datapoints/s steady; size for 10โท peak.
- At ~1.4 bytes/sample compressed: 10โถ/s โ ~120 GB/day raw-resolution.
- 15-day raw retention โ ~1.8 TB. Genuinely small โ the compression is the headline number.
- Index: ~10โท series ร ~200 B = 2 GB of index in RAM, which is why cardinality (not volume) is the binding constraint.
- Ingest node handles ~200k samples/s โ ~50 nodes at peak.
Detection latency budget (60 s)
- Agent scrape interval: 10 s (worst-case age at emit)
- Agent โ ingest โ log commit: ~2 s
- Streaming detector eval tick: 15 s
- Alert manager group_wait: 10 s (batches the storm)
- Notifier โ page: ~3 s
- ~40 s, leaving ~20 s of slack โ which is exactly the budget that
for_durationeats, so critical detectors runfor_duration=0on high thresholds.
10. 30-second recap
Agents push batched samples to a stateless ingest tier that enforces per-tenant series-cardinality quotas โ cardinality, not volume, is how these systems actually die. Samples land in a partitioned log, which feeds two consumers: a streaming detector for sub-60-second threshold alerts, and a TSDB sharded by series ID storing Gorilla-compressed columnar blocks at about 1.4 bytes a sample. Batch detectors read the TSDB to compute seasonal baselines using median and MAD rather than mean and stddev, because the outliers you're hunting would otherwise inflate your own band. Everything funnels into one alert manager that dedups on fingerprint, groups by shared labels, applies inhibition rules so a service-down alert suppresses its dependents, and honors silences. Two things I'd insist on: an explicit staleness detector, because absent data and healthy data look identical to a threshold, and a dead-man's switch on a separate stack in a separate failure domain โ otherwise the monitoring system goes blind precisely when production breaks.
Find-My-Device full
L6 GCP onsite SD, 2026-05. "่ฎพ่ฎก find my iphone๏ผ้่ฆๅ็่่้็งๆง๏ผๅฎๅ จๆงๅๆ็." The prompt names privacy, security, and efficiency as the graded axes โ that's unusual and it means a correct-but-generic location-tracking design fails this round.
1. Requirements in one line
Functional
- Owner sees near-real-time location of their devices from web or another device.
- Lost mode: play sound, mark as lost with a message, remote lock, remote wipe.
- Offline devices: report position when they come back online; optionally crowd-sourced reporting via nearby devices.
- Family sharing: explicitly authorized members can view, with revocation.
- Audit log the owner can read: who queried my location, when.
Out of scope: the map rendering stack, device activation lock enforcement, carrier integration.
Non-functional
- 10โธ devices; peak 10โถ location updates/sec.
- Query QPS is tiny by comparison โ people check rarely. Write-dominated by ~1000:1.
- Battery: location reporting must cost single-digit mW average.
- The server must not be able to read device locations. Treat that as a hard requirement, not a nice-to-have.
Core tension: a normal design puts plaintext locations in a database so the service can index and query them efficiently. But a database of where 10โธ people are is the highest-value breach target you could build, and an insider-abuse vector. The entire design is about giving that up โ and then paying for it in query complexity and lost server-side features.
2. Core entities & schema
Device device_id, owner_id, model, enrolled_at public_key (device keypair, private key never leaves device) status -- NORMAL | LOST | WIPED LocationReport -- what the server actually stores report_key = H(rotating_public_key) <-- opaque, no device_id! ciphertext -- E(location, timestamp) under device key received_at -- server clock, for TTL only reporter_hint -- coarse, for abuse throttling only TTL = 7 days // server CANNOT link report_key -> device_id -> user Command -- the write-back channel command_id, device_id, type (SOUND|LOCK|WIPE|MARK_LOST) issued_by, issued_at, nonce signature (signed by owner's key) state -- PENDING | DELIVERED | ACKED | EXPIRED AccessGrant grantor_id, grantee_id, device_id, expires_at, revoked_at AuditEntry device_id, actor_id, action, ts, ip_coarse
Why the location store is a dumb blob store, and why that's the answer
The instinct is a geospatial index โ S2 cells, geohash, "so we can query by area." Resist it. Nobody queries "which devices are near this point"; the only query is "give me the reports for my device." That's a point lookup on a key the owner can compute. Since the sole access pattern is keyโblob, you can afford the strongest possible privacy posture: the server holds ciphertext it cannot decrypt, keyed by a rotating identifier it cannot link to a user.
Store shape: a wide-column / KV store sharded by report_key, 7-day TTL, no secondary indexes at all. Writes are blind appends; reads are a multi-get of the ~2016 candidate keys the owner derives locally. The device registry and grants live in a normal replicated SQL store, because those do need joins, revocation semantics, and audit โ but they contain no location.
3. API interfaces
Report path (huge volume, unauthenticated by design)
POST /v1/reports // from ANY device, incl. finders
{ report_key, ciphertext, ts }
โ 204
// No auth header tied to identity โ attaching one
// would let the server correlate reporter to subject.
// Abuse control is anonymous rate-limiting:
// device attestation token + per-token quota.
POST /v1/commands/{device_id}
Signed-By: owner_key
{ type: LOCK, nonce, expires_at, signature }
โ 202 { command_id }
Query path (tiny volume)
POST /v1/reports:fetch
{ report_keys: [k1..kN] } // client-derived
โ { found: [{key, ciphertext, received_at}] }
// client decrypts locally; server sees nothing
GET /v1/devices // my devices + status
POST /v1/grants { grantee, device_id, expires_at }
DELETE /v1/grants/{id}
GET /v1/audit?device_id= // who looked, when
Fetch is a multi-get of opaque keys. The server can't tell whose device it is, or even that the N keys belong to one device. That property is the whole design.
4. Architecture
Green is the crowd-sourced report path โ a stranger's phone relays an encrypted blob it cannot read, to a server that also cannot read it. Red is the command path, signed by the owner's key and verified on the device, so a compromised server still can't wipe your phone.
5. Deep dive โ the privacy construction
Rotating keys, derived by both sides
At pairing (owner device + lost device share a secret):
master_secret ms (never leaves either)
Every 15 minutes, epoch i:
sk_i = KDF(ms, i) # both sides derive
pk_i = curve_pub(sk_i) # advertised over BLE
key_i = SHA256(pk_i) # the report_key
Finder:
sees pk_i over BLE
loc_ct = ECIES_encrypt(pk_i, {lat,lng,acc,ts})
POST { report_key: SHA256(pk_i), ciphertext: loc_ct }
Owner looking for the device:
for i in last 7 days (672 epochs @15m):
keys.append(SHA256(pk_i))
POST /reports:fetch { keys } # ~672 keys
decrypt each hit with sk_i # locally
What each property buys, and what it costs
- Server can't decrypt โ it never has any private key. A full database dump leaks nothing but "some blobs existed."
- Server can't link reports to a device โ
pk_irotates every 15 min, so consecutive reports look unrelated. This is what stops the server (or an insider) reconstructing anyone's movement history. - Finder can't track the lost device โ it sees only a rotating pseudonym and can't decrypt what it just uploaded.
- Stalking defense โ receiver devices detect an unknown beacon that persists across epochs while moving with you and alert the user. This is a real, shipped requirement and volunteering it is a strong signal.
Cost, stated honestly: ~672 key lookups per query instead of one; no server-side features at all (no "notify me when it moves", no geofencing, no server-side history search); and if the owner loses every device holding the master secret, past reports are unrecoverable. That last one is a genuine product trade and worth naming.
6. Deep dive โ command authenticity (the wipe problem)
Sequence
1. Owner authenticates (with 2FA) to the account service.
2. Client builds:
cmd = {device_id, WIPE, nonce, expires_at}
sig = sign(owner_private_key, cmd)
3. Server stores cmd, does NOT need to trust itself:
it also checks the owner's session, but that is
defense in depth, not the control.
4. Push wakes the device.
5. DEVICE verifies sig against the owner pubkey it
pinned at enrollment, checks nonce not seen and
expires_at in the future, then executes.
6. Device ACKs with a signed receipt -> audit log.
Why verification must happen on the device
If the server decides who may wipe, then a server compromise, a malicious insider, or a support-tool bug can wipe 10โธ devices. Verifying the owner's signature on the device reduces the server to an untrusted transport: the worst it can do is refuse to deliver commands (a denial of service, recoverable) rather than forge them (catastrophic, irreversible).
Defend it โ replay. Without nonce + expires_at, anyone who captures a valid LOCK command can replay it forever. The device keeps a small seen-nonce window and rejects anything expired, which bounds that memory.
Defend it โ the offline wipe race. A stolen device that's kept offline never gets the command. WIPE is queued with a long TTL and fires the moment it gets network; meanwhile the real protection is local full-disk encryption, which the remote wipe merely accelerates by destroying the key. Say this โ candidates who present remote wipe as the primary protection are wrong, and the interviewer knows it.
7. Deep dive โ efficiency (the third named axis)
Battery on the device
| Choice | Why |
|---|---|
| BLE advertise, don't GPS-fix | The lost device only broadcasts a rotating pubkey โ no GPS, no radio uplink. Micro-amp scale; a dead-battery phone can beacon for hours on reserve power. |
| Finder supplies the location | The finder was already computing its own position. Marginal cost โ one small HTTPS POST, batched. |
| Batch + opportunistic upload | Finders buffer reports and flush on Wi-Fi / when the radio is already up. Never wake the modem just for this. |
| Adaptive rate on the owner's own devices | Report position on significant-change, not on a timer. Stationary device โ near-zero traffic. |
Efficiency on the server
Write path is 10โถ/s of ~200-byte blind appends with no index maintenance and no read-modify-write โ the cheapest possible write. Sharded by report_key, which is a hash, so distribution is uniform by construction and there is no hot key possible: a popular location doesn't concentrate, because keys are per-device-epoch, not per-place.
Dedup: twenty phones in a train carriage all report the same beacon in the same epoch. Twenty near-identical ciphertexts under one key. Cap at ~K reports per (key, epoch) โ first-K-wins โ and drop the rest at ingest. That's a ~10ร write reduction in dense areas and it costs nothing, since the owner only needs one good fix per epoch.
7-day TTL is doing heavy lifting: it bounds storage, and it bounds the blast radius of any future cryptographic weakness. Longer retention has no product value โ nobody finds a phone with a 3-week-old fix.
8. Follow-ups โ answers to have ready
A malicious actor floods the report endpoint with garbage under someone's key.
They'd have to know H(pk_i), which requires BLE proximity during that 15-minute epoch โ so the attack is inherently local and short-lived. Beyond that: the endpoint requires a device-attestation token (Play Integrity / DeviceCheck equivalent) with a per-token quota, so anonymous doesn't mean unlimited. And garbage ciphertext simply fails to decrypt on the owner's client and is discarded โ the poisoning doesn't corrupt anything, it just adds noise, which the K-per-epoch cap bounds.
The report store dies. What breaks?
Fail-open on the write path: ingest keeps accepting and buffers to a durable log, because a dropped report is permanently lost โ there's no retry from a stranger's phone that has already walked away. Reads simply return nothing, and the client shows "no recent location" rather than a stale one, which matters: showing an hour-old position as current sends someone to the wrong address. Registry and command paths are on a separate store, so lock and wipe still work during a location outage โ that separation is deliberate and worth calling out.
How does family sharing work if the server can't decrypt?
Key sharing, not server-side authorization. The owner wraps the master secret (or a derived, time-boxed viewing key) to the grantee's public key and hands it over through the registry. The server stores an opaque wrapped blob and enforces the grant record for UX and audit, but the actual cryptographic capability lives with the grantee. Revocation is therefore not instant for already-fetched data โ you must rotate the master secret on revoke, which invalidates future reports but not the ones the grantee already downloaded. State that limitation plainly; pretending revocation is instant is the wrong answer.
Law enforcement requests a user's location history.
The design's answer is that we can't produce it โ we hold ciphertext keyed by unlinkable rotating identifiers, with a 7-day TTL. That's not evasion, it's an architectural property chosen up front, and it's the same posture Apple ships. What we can produce is the registry: account, device enrollment, and the audit log of who queried. Being explicit about where the boundary sits is the mature answer here.
672 keys per fetch seems wasteful. Optimize it?
In practice the client fetches newest-first and stops early โ the common case is "where is it right now," which hits within the first few keys, so the tail only gets scanned when a device has been offline for days. You can also widen the epoch (15 min โ 1 h) to cut keys 4ร, at the cost of coarser unlinkability, since a longer-lived pseudonym is easier to correlate. That's the exact trade to name: epoch length is the privacy/efficiency dial. And the fetch is a batched multi-get on a hash-sharded KV store, so 672 keys is a handful of parallel shard reads, not 672 round trips.
Scale to 10ร devices.
The write path scales linearly โ hash-sharded blind appends with no cross-shard coordination, so 10โท/s is more shards and nothing else. The parts that don't scale for free are the anonymous abuse quota (attestation token issuance becomes the bottleneck) and BLE spectrum congestion in dense areas, which is a client-side problem solved by advertise-interval backoff. The server design genuinely doesn't need re-architecting, and being able to say why โ no indexes, no joins, no hot keys, uniform hashing โ is the point.
9. Numbers to drop
Write volume
- 10โถ reports/s ร ~250 B = 250 MB/s ingest, ~21 TB/day.
- With per-epoch dedup at K=5 in dense areas: ~10ร reduction โ ~2 TB/day effective.
- 7-day TTL โ ~15 TB steady state. Small enough to keep entirely on SSD.
- Shard at ~50k writes/s/node โ ~20 ingest shards at peak. Uniform by construction.
Read & crypto
- Queries: say 10โธ users ร 1 check/week โ 165 QPS. Three orders of magnitude below writes โ which justifies optimizing everything for write.
- 672 keys/fetch ร 165 QPS = ~110k key lookups/s. Trivial for a hash-sharded KV.
- ECIES encrypt on a finder: ~1 ms, negligible against a BLE scan already running.
- Key derivation on the owner client: 672 ร KDF โ tens of ms, done once per open.
10. 30-second recap
The lost device does nothing but broadcast a Bluetooth pseudonym derived from a shared master secret, rotating every 15 minutes. Any nearby stranger's phone picks it up, encrypts its own GPS position to that rotating public key, and posts the blob under a key that is just the hash of the pseudonym. The server stores ciphertext it cannot decrypt, under identifiers it cannot link to a user or to each other, with a 7-day TTL โ so a full breach leaks nothing and there's no location history to subpoena or for an insider to abuse. The owner derives the same key sequence locally, multi-gets the last week of candidate keys, and decrypts on-device. Commands like lock and wipe are signed with the owner's key and verified on the device, so a compromised server can delay a wipe but never forge one. Efficiency falls out of the same design: blind appends with no index, hash-uniform sharding so no hot keys, and the beacon costs microamps because the finder supplies the location. The price I'm paying is real โ 672 key lookups instead of one, no server-side geofencing or history, and revocation that requires rotating the master secret rather than flipping a flag.
Street View image ingest full
L6/L7, Google Cloud Storage org, phone screen, 2026-01. Verbatim: "Design a system that supports Google Map Street View storage. The images will be uploaded from a taxi. Each taxi has a camera; images are processed by downstream โ image understanding, display to users, map generation." Reported follow-ups: auth design and what to do when a token is compromised; upload response design; behavior on poor network. Interviewer was near-silent and let the candidate drive for the full hour.
1. Requirements in one line
Functional
- Fleet vehicles capture panoramic frames + GPS/IMU metadata and upload them.
- Uploads survive intermittent connectivity: resumable, deduplicated, eventually complete.
- Ingested images are durably stored and published to downstream consumers (blur/PII, image understanding, tiling, map generation).
- Operators can see per-vehicle ingest status and re-drive coverage gaps.
Out of scope: the ML models themselves, the map-tile serving CDN, camera firmware.
Non-functional
- 10k vehicles; each ~1 frame/sec while driving, ~8 h/day. ~10 MB/panorama.
- Ingest is throughput-bound, not latency-bound โ nobody needs the photo in a second.
- Durability is the hard requirement: a lost frame means re-driving a street. Target 11 nines.
- Cost matters at petabyte scale; storage tiering is part of the design, not an afterthought.
Core tension: the client is a moving vehicle on flaky cellular, generating far more data than its link can carry, and it cannot be trusted (a stolen credential lets someone poison the map). So the design is really about a well-behaved offline-first client plus an untrusted-upload security model, with the storage tier being the easy part.
2. Core entities & schema
CaptureSession
session_id, vehicle_id, driver_id, started_at, route_id
status -- ACTIVE | UPLOADING | COMPLETE | PARTIAL
Frame -- metadata, in a wide-column store
frame_id = uuidv7(capture_ts) -- time-sortable
session_id, seq
content_hash (sha256 of the image bytes) <-- dedup + integrity
captured_at, lat, lng, heading, accuracy
s2_cell_l16 -- geo index key
blob_uri, bytes, camera_id
state -- PENDING | UPLOADED | VERIFIED | PUBLISHED | QUARANTINED
Upload -- the resumable session
upload_id, frame_id, total_bytes
received_ranges [] -- byte ranges committed
expires_at
Blob -- object storage, immutable
path = /raw/{s2_cell}/{yyyymmdd}/{content_hash}
// content-addressed: same bytes uploaded twice = one object
Why object storage + a metadata store, and why content-addressing
Images never belong in a database โ 10 MB blobs are what object storage exists for: cheap, replicated, tiered, with no query needs. Metadata is small, queried by geography and by session, and needs a secondary index, so it lives separately in a wide-column store keyed by frame_id with an s2_cell index for "what have we covered here?"
Content-addressing the blob path is the load-bearing choice. The client hashes bytes before upload; the path is the hash. That gives you three things for free: (1) a retry that re-uploads the same frame overwrites itself instead of duplicating โ idempotency without any server-side dedup logic; (2) end-to-end integrity, since the server recomputes the hash and rejects a mismatch, catching both corruption and tampering; (3) natural dedup when a vehicle stops at a light and the camera captures near-identical frames โ though see the follow-ups, byte-identical is rarer than you'd hope.
uuidv7 for frame_id rather than random UUIDv4 because it's time-prefixed and therefore time-sortable, which makes session-ordered scans a contiguous range read rather than a scatter.
3. API interfaces
Upload path (resumable, the actual question)
POST /v1/uploads:init
Authorization: Bearer {short-lived vehicle token}
{ session_id, seq, content_hash, bytes, captured_at,
lat, lng, heading }
โ 200 { upload_id,
upload_url, // signed, direct-to-storage
expires_in: 3600,
already_have: false } // <-- dedup short-circuit
PUT {upload_url}
Content-Range: bytes 5242880-10485759/10485760
โ 308 { received: "0-10485759" } // resume point
โ 201 when complete
GET /v1/uploads/{upload_id} // "where was I?"
โ { received_ranges: [[0, 5242879]] }
Downstream & ops
// consumers subscribe, they don't poll
topic: frames.uploaded { frame_id, blob_uri, geo }
topic: frames.published { frame_id, after blur+verify }
GET /v1/coverage?s2_cell=&since=
โ { frames, gaps: [{lat,lng,reason}] }
GET /v1/sessions/{id}/status
โ { captured: 28800, uploaded: 28112, pending: 688 }
POST /v1/frames/{id}:quarantine { reason }
Answering the "upload response design" follow-up: return 308 with the exact committed byte range, plus already_have on init so a client that crashed after uploading never re-sends 10 MB. The response's job is to make the client's retry decision unambiguous.
4. Architecture
Bytes never traverse the ingest service โ it only issues signed URLs, so the control plane stays small and cheap while 10 MB objects go straight to storage. Downstream consumers subscribe to an event log rather than being called synchronously, so a slow ML pipeline can never back-pressure a vehicle on a cellular link. Note the blur/PII stage is a gate, not a parallel consumer.
5. Deep dive โ the offline-first client (the "poor network" follow-up)
The client is the interesting system
Capture loop (never blocks on network):
frame -> local SSD spool + WAL entry
compute sha256 while writing
spool is a bounded ring: 500 GB โ 14 h of capture
Upload loop (independent, opportunistic):
while spool not empty:
pick oldest UNSENT frame
if link_quality == CELLULAR and not urgent:
upload at throttled rate (leave headroom)
if link == WIFI (depot / known AP):
upload at full rate, drain aggressively
on 5xx / timeout:
exponential backoff + jitter, keep frame
on 201:
mark SENT, delete local copy lazily
(keep until VERIFIED event or 24h)
Chunk size adapts: start 8 MB, halve on each
timeout down to 256 KB, grow back on success.
Why this shape
Capture and upload must be decoupled queues. If capture blocks on upload, a tunnel or a dead cell means lost coverage and a re-drive โ the single most expensive failure in this system. Spooling to local disk turns a network problem into a storage problem, which is far cheaper.
Adaptive chunk size is the concrete answer to "network is bad." On a marginal link, a 10 MB PUT that dies at 90% wastes 9 MB of the vehicle's uplink; 256 KB chunks lose almost nothing per failure. The trade is more round-trips and more per-request overhead, so you grow the chunk back when the link proves itself.
Deletion is delayed until VERIFIED, not until the 201. A 201 means bytes landed; it doesn't mean the server rehashed them successfully. Deleting the only copy on a 201 is how you lose data to a silent corruption. Cost: local disk holds a day of already-uploaded frames.
Belt-and-suspenders โ the bounded spool. 14 hours of buffer exceeds an 8-hour shift, so a vehicle that never sees network all day still loses nothing; it drains at the depot on Wi-Fi. If the ring does wrap, drop by lowest marginal coverage value (a frame from an already-well-covered cell) rather than oldest-first. Volunteering a smarter eviction policy than FIFO reads well.
6. Deep dive โ auth, and the compromised-token follow-up
Layered credentials
Layer 1 โ hardware identity (long-lived, never sent)
Per-vehicle key in a TPM/secure element at fleet
provisioning. Private key cannot be extracted.
Layer 2 โ short-lived access token (the thing on the wire)
Vehicle signs a challenge with the TPM key
-> token service returns JWT, TTL 1 hour,
scoped: {vehicle_id, session_id, WRITE only,
geo_bbox of today's route}
Never a long-lived API key on the device.
Layer 3 โ per-object signed URL (TTL 1 hour, single path)
Scoped to exactly /raw/{cell}/{date}/{content_hash}
So even a leaked URL can write ONE object,
whose path is the hash of its own contents โ
so it cannot write anything but that exact image.
"The token is compromised. What do you do?"
Answer in the order the interviewer is grading: contain, revoke, detect, prevent.
- Blast radius is already small by design โ the token is write-only, expires in an hour, is scoped to one vehicle and one geographic box, and the signed URLs it can obtain are content-addressed. An attacker cannot read anyone's imagery, cannot overwrite existing frames with different bytes, and cannot write outside today's route area.
- Revoke: add the token's
jtito a deny list checked at the ingest control plane (small, because TTL is short โ the list self-cleans in an hour). Revoke the vehicle's provisioning cert to stop it minting new tokens. - Detect: the signals that catch this are behavioral โ uploads from an IP/ASN that doesn't match the vehicle's cellular carrier, frames whose GPS is inconsistent with the session's route continuity or physically impossible speed, upload rate exceeding the camera's capture rate, or two concurrent sessions for one
vehicle_id. Any of these quarantines the frames rather than dropping them, so a false positive is recoverable. - Prevent: mutual TLS with the device cert on the control plane, so the bearer token alone is insufficient; anything anomalous goes to
QUARANTINEDand is human-reviewed before publish.
The point to land: the right answer to a compromised credential isn't a faster revocation loop, it's designing so that a compromised credential can't do much. Short TTL, narrow scope, content-addressing, and write-only are all doing that work.
7. Deep dive โ downstream fan-out and the PII gate
Event-driven, not synchronous
Ingest publishes frames.uploaded and stops caring. Consumers โ blur, image understanding, tiling, map generation โ subscribe independently with their own offsets, their own scaling, and their own failure domains. A stalled ML pipeline creates consumer lag, not upload failures.
Reprocessing is the reason raw is immutable and kept forever. Models improve; when a better sign-detection model ships you replay the last N years of frames.uploaded from the log (or scan the metadata store by S2 cell) rather than re-driving the planet. Derived artifacts โ tiles, extracted geometry โ are therefore disposable, and treating them that way is what makes the storage bill defensible.
The blur stage must be a gate, not a peer
Faces and license plates must be blurred before anything user-facing exists. If blur is just another parallel subscriber, there's a window where the tiling pipeline has already published an unblurred panorama โ a privacy incident with legal consequences in the EU and elsewhere.
So: two topics. frames.uploaded is consumed only by blur/verify. Blur writes a redacted derivative and emits frames.published; everything user-facing subscribes to that one. Internal-only consumers may read raw under separate authorization and audit.
Defend it: blur models miss things, so you need a takedown path โ a user reports their face, you re-blur that frame and purge derived tiles, which is only tractable because frame_id โ tile lineage is recorded. Without lineage tracking, a single takedown request means rebuilding a region. Say this; it's the kind of operational detail that separates L5 from L4 here.
8. Follow-ups โ answers to have ready
Two frames have the same content hash. Is that dedup or a bug?
Usually a bug worth investigating. Real-world sensor noise means two genuinely separate captures essentially never produce identical bytes โ so a collision means either the camera emitted a duplicate frame (stuck buffer), or the client retried and re-uploaded, or someone is replaying captured data. The write is idempotent either way, so nothing breaks; but I'd count it as a metric and alert on the rate. For perceptual dedup of a vehicle stopped at a red light, content hashing is the wrong tool โ use a perceptual hash or simply drop frames when GPS shows the vehicle hasn't moved, which is cheaper and done on the client.
Object storage in one region goes down mid-shift.
Vehicles keep capturing โ that's the whole point of the local spool โ and uploads fail into backoff. The ingest control plane should detect the regional failure and start issuing signed URLs for a secondary region, since blobs are content-addressed and the region is just a prefix; the metadata store records where each blob actually landed. Recovery cost is an asynchronous cross-region copy job later. What you must not do is let the vehicle agent block or drop frames while the server sorts itself out.
A vehicle uploads frames with plausible but wrong GPS.
This is the poisoning case and it's the reason for the verifier stage. Checks: route continuity (does this position follow the previous one at a physically possible speed?), IMU cross-check, and consistency with existing imagery for that S2 cell. Suspicious frames go to QUARANTINED โ never auto-deleted, because a false positive on a legitimately unusual route (a ferry, a new road) would silently lose coverage. Quarantine plus human review is the right severity.
How much does this cost, and how do you cut it?
Raw is the dominant line item and it grows forever. Lifecycle tiering does the heavy lifting: hot for 30 days while the pipelines chew on it, nearline for a year, coldline archive after. Raw is never deleted because re-driving costs vastly more than storing. Derived artifacts go the other way โ regenerate tiles on demand from raw rather than storing every historical version. And compression: raw sensor frames stored as lossless for reprocessing, but the panorama derivatives as aggressive lossy, since human viewers can't tell.
Why not have the vehicle upload straight to the bucket with no control plane at all?
Because you'd need a long-lived broad credential on a physically accessible device โ exactly the thing the layered auth avoids. The control plane exists to mint narrow, short-lived, per-object grants and to record metadata transactionally with the intent to upload, so you can detect a frame that was initiated but never completed. It stays cheap because it never touches bytes: 10k vehicles at ~1 init/sec is ~10k QPS of tiny JSON, which is a handful of pods.
Scale from 10k vehicles to 1 M (consumer dashcams).
The storage and event tiers scale horizontally without change. Three things break. First, trust: consumer devices have no TPM and no fleet provisioning, so poisoning becomes the dominant risk and you need reputation scoring and heavy cross-validation between contributors. Second, redundancy: a million dashcams will cover the same popular streets thousands of times, so you need server-side coverage-aware admission โ reject the upload at init time when that cell is already well covered, which also saves the contributor's bandwidth. Third, cost per useful frame collapses unless you do that admission control. That's the re-architecture point, and it's on the ingest control plane, not on storage.
9. Numbers to drop
Ingest volume
- 10k vehicles ร 1 frame/s ร 8 h = 288 M frames/day.
- At 10 MB each: ~2.9 PB/day raw. This number is the design constraint โ say it early.
- Sustained ingest if spread over 24 h: ~33 GB/s. Per vehicle that's only ~10 Mbps, which is exactly why depot Wi-Fi drain matters โ cellular can't carry it all.
- Control plane: ~3.3k init QPS average, ~10k peak. Trivially small next to the byte volume.
Storage & metadata
- ~1 EB/year raw. With lifecycle tiering (30d hot / 1y nearline / then coldline) the blended rate is roughly 4โ5ร cheaper than all-hot.
- Metadata: 288 M rows/day ร ~500 B = ~145 GB/day โ 0.005% of the blob volume, which is why it can afford indexes.
- Vehicle spool: 500 GB SSD รท 10 MB = 50k frames โ 14 h of capture, comfortably more than one shift.
- Depot drain: 14 h of backlog at 1 Gbps Wi-Fi โ 70 min. Sizes the depot link.
10. 30-second recap
The vehicle agent treats capture and upload as two independent queues joined by a local disk spool sized for a full shift, so a tunnel or a dead cell never costs coverage โ that's the expensive failure here, since a lost frame means re-driving a street. The client hashes each frame, calls a small control plane to initialize a resumable upload, and gets back a signed URL scoped to exactly one content-addressed path; bytes then go straight to object storage, so the control plane never touches a 10 MB payload. Content-addressing gives idempotent retries, end-to-end integrity, and free dedup. On a bad link the client halves its chunk size down to 256 KB so a failure wastes almost nothing, and it doesn't delete its local copy until the server confirms it rehashed the bytes. Auth is layered โ a TPM key that never leaves the vehicle mints one-hour, write-only, geo-scoped tokens โ so a stolen token can't read anyone's imagery, can't overwrite existing frames, and expires before you've finished revoking it. Ingest publishes an event; blur and PII redaction is a gate rather than a parallel consumer, so nothing user-facing can ever see an unblurred face, and everything downstream subscribes to the post-blur topic. Raw is immutable and kept forever because models improve and reprocessing is far cheaper than re-driving.
Global chain restaurant menu sync full
็ๅฎถ SD, 2026-02. Verbatim: "็ปไธไธชๅ จ็่ฟ้้คๅ ่ฎพ่ฎก่ๅๆดๆฐ็ณป็ปใๅ่ฎพๆฏๅฎถ้คๅ ๅฏ่ฝๆๅ็งๅฑ็คบ่ๅ็่ฎพๅคใ ่ฟ้้คๅ ๅจๅๅฝๅฎถ็่ๅๆ็ธๅ็้จๅ๏ผไนๆๆฌๅฐๅ็้จๅใ่ๅ็ๆดๆฐ็ฑๆป้จๆงๅถ๏ผๆฏๅคฉๆฉไธญๆ้คๅฏ่ฝ่ทๅไธๅคฉไธไธๆ ทใ" Read the prompt carefully โ it hands you three hard requirements: heterogeneous display devices, inheritance with localization, and time-scheduled variants.
1. Requirements in one line
Functional
- HQ authors a global menu; countries/regions/stores override parts of it (price, availability, language, local items).
- Menus vary by daypart: breakfast / lunch / dinner, and can differ day to day.
- Push updates to heterogeneous devices: digital menu boards, kiosks, tablets, drive-thru displays, the mobile app, and third-party delivery aggregators.
- Scheduled activation ("this menu goes live Monday 6am local") and emergency rollback.
- A store must show a correct menu even with no network.
Out of scope: order placement, inventory/POS integration beyond an availability flag, the authoring UI itself.
Non-functional
- 40k stores ร ~5 devices = 200k endpoints, ~100 countries.
- Update propagation: minutes is fine for planned changes; seconds matter for a takedown (allergen recall, wrong price).
- Availability >> consistency at the edge: a device must never show a blank screen.
- Correctness of price is the one thing that's legally and financially sensitive.
Core tension: the update volume is laughably small โ this is not a scale problem. It's a configuration-correctness and safe-rollout problem across 200k unreliable, heterogeneous, intermittently-connected endpoints. Candidates who reach for Kafka and sharding here have misread the question.
2. Core entities & schema
MenuItem -- the catalog (global)
item_id, canonical_name, category, image_set,
allergens[], nutrition{}, default_price_cents
MenuLayer -- the inheritance chain. THE key entity.
layer_id, scope_type -- GLOBAL | COUNTRY | REGION | STORE
scope_id -- null | "JP" | "kanto" | "store_8842"
parent_layer_id
patch { -- sparse: only what this layer CHANGES
add: [item_id...],
remove: [item_id...],
override: { item_id: {price, name_i18n, image} }
}
MenuVersion -- immutable, content-addressed
version_id = hash(resolved menu bytes)
layer_set[], daypart, effective_from, effective_to
state -- DRAFT | APPROVED | STAGED | ACTIVE | ROLLED_BACK
approved_by, approved_at
DeviceRegistration
device_id, store_id, type (BOARD|KIOSK|DRIVE_THRU|APP),
capabilities {screen, aspect, supports_video, locales},
current_version_id, last_heartbeat_at
Why layered patches instead of a menu per store
The naรฏve model โ one full menu document per store โ means 40k documents, and a global price change requires rewriting all of them with no atomicity and no way to tell an intentional local difference from a stale copy. The layered model makes inheritance explicit and diffable: HQ edits the GLOBAL layer, and every store that hasn't overridden that item inherits the change automatically.
Resolution is a deterministic fold โ GLOBAL โ COUNTRY โ REGION โ STORE, later layers win. Determinism is what lets you content-address the resolved output: same layers in, same version_id out. That in turn gives you free change detection (a device compares hashes), free dedup across identical stores, and an unambiguous rollback target.
Storage is boring on purpose: a replicated SQL store for layers and versions (you need transactions, approval workflow, and audit on price), plus object storage + CDN for the resolved bundles that devices actually download. Resolution happens once at publish time, not per-device โ 40k stores ร 3 dayparts is ~120k resolutions per publish, which is a few seconds of compute, and it means a device does zero logic.
3. API interfaces
Device path (pull, not push)
GET /v1/devices/{device_id}/manifest
If-None-Match: {current_version_id}
โ 304 Not Modified // the 99.9% case
โ 200 {
schedule: [
{ daypart: "breakfast", from: "06:00",
to: "10:30", version_id: "a1b2",
bundle_url: "https://cdn/.../a1b2.tar",
sha256: "...", size: 4_200_000 },
{ daypart: "lunch", ... }
],
poll_after_sec: 300,
emergency_channel: "wss://..."
}
POST /v1/devices/{device_id}/heartbeat
{ active_version_id, staged_version_ids[],
render_errors[], clock_skew_ms }
Devices pull on a jittered interval and hold a long-lived websocket only for emergencies. Pull scales trivially through a CDN and survives a device being offline for a week; push-only would require the server to track 200k connections just to deliver a change that isn't urgent.
Authoring & ops path
PUT /v1/layers/{layer_id} { patch }
POST /v1/versions:preview
{ layer_set, daypart, store_id }
โ fully resolved menu + diff vs current
POST /v1/versions/{id}:approve // 2-person for price
POST /v1/rollouts
{ version_id, cohort: "canary_50_stores",
effective_from: "2026-03-01T06:00 LOCAL" }
POST /v1/rollouts/{id}:halt
POST /v1/emergency/takedown
{ item_id, scope, reason } // seconds, not minutes
GET /v1/fleet/status?version_id=
โ { on_target: 39_812, stale: 188, unreachable: 40 }
4. Architecture
Publishing is a fold + validate + freeze into an immutable, content-addressed bundle. Distribution is boring CDN pull with a per-store gateway that pre-stages every upcoming version, so activation is a local event that needs no network at all. The red path is the only push channel, reserved for takedowns.
5. Deep dive โ staging and local activation (the daypart problem)
Never activate over the network
Publish (T-24h or more): resolver produces bundles for every (store, daypart) combination for the next 48h โ manifest lists them ALL with local times Store gateway (continuously): downloads every staged bundle verifies sha256 keeps last-known-good + all staged Activation (exactly at 06:00 store-local): gateway compares local clock to schedule swaps active pointer atomically devices re-render from local files Network down at 06:00? Still switches. Network down for 3 days? Runs out of staged versions โ holds last-known-good, alarms.
Why this is the crux of the question
The prompt's "ๆฉไธญๆ้คๅฏ่ฝ่ทๅไธๅคฉไธไธๆ ท" is a trap for anyone who designs a push-at-activation-time system. If the breakfast menu is pushed at 06:00 local, then every store in a timezone hits you simultaneously โ a thundering herd at every hour boundary around the globe โ and any store whose link is down at that moment shows the wrong menu during its busiest period. Pre-staging converts a synchronized network event into an unsynchronized download plus a local clock comparison.
Defend it โ clock skew. The design now depends on device clocks. A gateway with a drifting clock switches to dinner at lunchtime. Mitigations: NTP on the gateway, clock skew reported in every heartbeat with an alert above ~60 s, and a sanity rule that the gateway won't jump more than one daypart forward without a server confirmation. This is the failure mode the interviewer will hunt for once you propose local activation, so raise it yourself.
Timezones and DST: schedules are stored as local wall-clock times with an IANA timezone ID per store, never as UTC offsets โ offsets break twice a year, and a store in a DST-shifting region would show breakfast an hour late every spring.
6. Deep dive โ heterogeneous devices
Ship data, plus per-capability renditions
Bundle contents:
menu.json -- semantic content, device-agnostic
items, prices, i18n strings, order
layouts/
board_16x9.json kiosk_portrait.json
drive_thru.json app_compact.json
assets/
{hash}.webp @1x/2x {hash}.mp4 (boards only)
render_hints: { max_items_per_panel, font_scale }
Device fetches ONLY the layout + assets matching
its registered capabilities โ a drive-thru display
never downloads 4K video it can't show.
The trade
The alternative โ pre-rendering images per device type server-side โ makes devices trivially dumb but multiplies bundle count by device type, and any layout tweak means republishing everything. Shipping semantic JSON plus layout descriptors means one content pipeline and per-device presentation, at the cost of some rendering logic on the endpoint.
Defend it โ the capability lie. Devices misreport capabilities, and old firmware won't understand a new layout field. So: version the bundle schema, require devices to declare a supported schema range in the manifest request, and have the server serve the highest bundle version that device can parse. Unknown fields must be ignored, not fatal. Without this, one new field bricks every unpatched screen in the fleet โ and menu boards get firmware updates roughly never.
Third-party aggregators (delivery apps) are just another consumer type, but they pull via API rather than bundle, and they're the one consumer where a stale price becomes a customer-facing refund. Give them a webhook on version change plus a short cache TTL.
7. Deep dive โ safe rollout and the takedown path
Two speeds, deliberately
| Planned change | Emergency takedown | |
|---|---|---|
| Trigger | Approved version, scheduled | Allergen error, wrong price, recall |
| Path | CDN pull, jittered, staged | Websocket push + poll interval collapses to 10 s |
| Granularity | Whole version | Single item_id: hide it |
| Target | Minutesโhours | < 30 s to 99% of fleet |
| Fallback | โ | If push fails, next poll catches it |
A takedown is a subtractive overlay, not a new version โ "hide item X everywhere" is a tiny message that any device can apply to whatever bundle it's currently running, including a stale one. That's why it's fast and why it works on devices that haven't synced.
Canary the way you'd canary code
A bad menu version is a production incident with revenue impact. Roll out by cohort: 50 stores โ one region โ country โ global, with automatic halt on error signals (render failures in heartbeats, device version-adoption stalling, or a spike in POS price mismatches). POST /rollouts/{id}:halt freezes advancement; rollback re-points to the previous version_id, which is instant because bundles are immutable and still cached at the edge.
Validation gates before any of that: no item priced at 0 or negative; no price change greater than X% without a second approver; every item has allergen data in every locale it's published to; every referenced asset exists; total items fit the smallest target screen. These catch the realistic disaster โ a fat-fingered price rolled to 40k stores โ far more reliably than a canary does, because a canary only catches what the canary stores happen to sell.
Defend it โ partial fleet state. During any rollout, different stores are legitimately on different versions. That's fine for menus but not for a nationally advertised promotion, so the schedule's effective_from is what gates display, while the rollout gates distribution. Separating those two is what lets you pre-stage for days and still have every store flip together.
8. Follow-ups โ answers to have ready
A store's internet has been down for two days. What's on the screen?
The correct menu, because the gateway pre-staged 48 hours of bundles and activation is a local clock comparison. Past that horizon it holds last-known-good and keeps switching dayparts using the most recent schedule it has โ a slightly stale menu beats a blank board. The heartbeat gap raises the store in the fleet dashboard as unreachable, and there's a manual USB/local-admin path for a store manager to sideload a bundle in a true emergency. The design decision to state plainly: we choose availability over freshness at the edge, because a dark menu board stops sales entirely.
HQ changes a global price but Japan has overridden that item. What happens?
Japan keeps its override โ that's the whole point of the layer model, and it's the correct default. The risk is silent divergence: HQ thinks it changed a price globally and doesn't realize 12 countries have pinned it. So the authoring UI must show, at edit time, how many downstream layers override this item, and the publish diff must list them explicitly. Optionally offer a "force, clearing overrides" action, which requires a higher approval level. The technical model is easy here; the product affordance around it is what makes it correct.
The resolver has a bug and produces a garbage menu.
Validation gates run after resolution and before the bundle is frozen, so structurally-invalid output never gets a version_id. For semantically-wrong-but-valid output, the canary cohort plus adoption monitoring catches it, and rollback is instant because the previous immutable bundle is still in the CDN and in every gateway's local cache. Worth noting the resolver is a pure function of the layer set, so it's the easiest component in the system to test exhaustively โ golden-file tests over real store configurations.
Why not just push everything over websockets? 200k connections isn't that many.
It isn't, and you do need the channel for takedowns. But making it the primary path means the server owns delivery state for 200k endpoints, reconnect storms after any deploy, and no natural way to serve a device that was offline for a week. Pull with a CDN inverts that: the server publishes an immutable artifact and forgets, and the edge absorbs the load. Push is reserved for the case where seconds matter, which is rare, small, and subtractive.
How do you know the fleet is actually correct right now?
Heartbeats carry active_version_id, so the fleet endpoint answers "how many stores are on the version I intended, right now" โ and that ratio, not the publish succeeding, is the definition of a successful rollout. Alert on stale-device count and on adoption curves that plateau. Add a synthetic: a handful of reference devices that screenshot and diff against the expected render, which catches the class of bug where the data is right and the rendering is wrong.
Scale it to a franchise model where each owner can edit their own menu.
The layer model already supports it โ franchisees author at the STORE layer. What changes is governance: you now need per-layer permissions, a policy engine restricting which fields a franchisee may override (price yes, allergen data absolutely not, brand assets no), and much stronger validation because the authors are no longer a small trusted team. Authoring volume goes from a handful of HQ edits a day to tens of thousands, which finally makes the write path non-trivial โ but distribution is unchanged.
9. Numbers to drop
Fleet & distribution
- 40k stores ร ~5 devices = 200k endpoints; ~5 dayparts ร 2 days staged = ~10 bundles live per store.
- Manifest poll every 5 min, jittered โ 200k / 300 s โ 670 QPS, and ~99.9% are
304s. This is a small service. - Bundle ~5 MB (mostly images, incremental after the first). Full fleet refresh = 200k ร 5 MB = 1 TB, spread over hours via CDN. Trivial.
- Resolution cost: 40k stores ร 3 dayparts = 120k folds per publish, ~ms each โ a few seconds on one machine.
Latency targets
- Planned change โ 99% of fleet: ~15 min (3 poll intervals).
- Emergency takedown โ 99%: < 30 s via websocket, with the 10 s collapsed poll as backstop.
- Daypart switch: 0 ms of network, purely local.
- Rollback: instant (re-point to a cached immutable bundle).
- Data volume overall is so small that the entire design is driven by correctness and offline behavior โ say this explicitly, it shows you sized the problem.
10. 30-second recap
Menus are modelled as a chain of sparse patch layers โ global, country, region, store โ that resolve by a deterministic fold, so HQ edits propagate automatically to everyone who hasn't overridden that item, and local customization is a first-class concept rather than a copy. Resolution happens once at publish time and produces an immutable, content-addressed bundle per store and daypart, which passes validation gates for price sanity and allergen completeness before it's allowed a version ID. Distribution is CDN pull with jittered polling โ mostly 304s โ and each store gateway pre-stages 48 hours of upcoming bundles, so the breakfast-to-lunch switch is a local clock comparison that needs no network at all. That's the key move: it avoids a global thundering herd at every hour boundary and means a store with a dead link still shows the right menu. Devices get semantic JSON plus a layout matching their declared capabilities, with schema versioning so a new field doesn't brick unpatched menu boards. Rollout is canaried by cohort with automatic halt, and there's a separate fast path for emergency takedowns โ a subtractive "hide this item" overlay that applies to whatever bundle a device is already running, so it works in under 30 seconds even on a stale device. The thing I'd watch hardest is clock skew on the gateways, since local activation makes correctness depend on their clocks.
Dictionary range query store condensed
็ๅฎถ L6, 2026-02. Verbatim: "่ฎพ่ฎกไธ็งๅญๅจ็ณป็ป๏ผ่ฝๅคๆๅญๅ ธๅบๆฅ่ฏขไธไธชๅบ้ดๅ ็ๆๆ่ฏใ ๆฏๅฆๆฐๆฎๆ {aa, aaa, ac, f, z}: ่พๅ ฅ [a, b] โ {aa, aaa, ac}; ่พๅ ฅ [a, aa] โ {aa}; ่พๅ ฅ [] โ {f, z}; ่พๅ ฅ [ab, b] โ {ac}." This is the rare SD prompt that's really a data-structure choice with a distributed-systems tail. Nail the local structure first, then scale it.
1. Requirements & the semantics trap
Clarify before designing
- Boundary semantics: the examples imply
[a, b)half-open โ[a, aa]returns{aa}but not{aaa, ac}, and[ab, b]returns{ac}. Confirm this out loud; getting it wrong invalidates everything after. - Result size: could a range match 10โธ words? If yes, the API must paginate with a cursor, not return a list.
- Mutability: read-only corpus (build once, serve forever) or live insert/delete? This is the biggest fork in the design.
- Scale: 10โถ words fits in memory on one box; 10ยนยน words does not. Ask.
- Latency SLO, and whether prefix/autocomplete queries are also needed (they change the structure choice).
Assume, and say so
- 10ยนโฐ words, ~20 bytes each โ ~200 GB. Doesn't fit one machine.
- Mostly reads; writes are bulk-loaded plus a live trickle.
- p99 < 50 ms for a range returning โค 1000 results.
- Results must be returned in sorted order and paginated.
Core tension: range queries want data sorted and contiguous; distribution wants data hashed and uniform. Those are directly opposed, and choosing sorted (range-partitioned) means accepting hot shards. Everything below follows from that.
2. Structure choice โ say why, not just what
| Option | Range query | Cost | Verdict |
|---|---|---|---|
| Sorted array + binary search | O(log n) to find start, then scan | Immutable; insert is O(n) | Right for a static corpus. Say this first โ it's the simplest thing that works. |
| B+ tree | O(log n) descend, then walk the leaf linked list | Node splits on write; ~1.3ร space | The answer for a mutable corpus. Leaves are sorted and linked, which is exactly a range scan. |
| LSM tree (SSTables) | Merge-iterate across sorted runs | Read amplification; compaction | Right when writes dominate. Range scan touches every level. |
| Trie / radix tree | Prefix-natural; range needs bounded DFS | High pointer overhead, poor cache locality | Only if prefix queries are the real requirement. For arbitrary [a, b) it's worse than a B+ tree. |
| Hash index | Impossible โ hashing destroys order | โ | Naming why this fails is a cheap point. Say it. |
Pick B+ tree for the mutable case and be explicit about why the trie loses: a trie shines when you're matching a shared prefix, but [ab, b) has no single prefix โ you'd descend to ab, DFS its subtree, then walk siblings up to b, which is a lot of pointer chasing for what a B+ tree does as one sequential leaf walk.
3. Architecture
Range-partition, don't hash-partition. The routing table maps lexicographic boundaries to shards, so range(a, b) touches only the shards that overlap [a, b) โ often one. Hashing would fan every query to every shard.
4. Trade-offs worth stating
The hot-shard problem you created
Real word distributions are wildly skewed โ a shard covering [s, t) holds far more English words than one covering [x, z), and query traffic is skewed too. Range partitioning guarantees this; it's the price of ordered scans.
- Split by load, not by key space. Choose boundaries so each shard holds roughly equal bytes and QPS, not equal alphabet width. Boundaries come from sampling the actual corpus.
- Dynamic split/merge like Bigtable tablets: a shard exceeding a size or QPS threshold splits at its median key and the registry updates.
- Read replicas absorb read skew without resharding โ cheap and usually sufficient, since this workload is read-dominated.
Streaming, not materializing
A range can match arbitrarily many words, so the coordinator must stream: open an iterator per overlapping shard, k-way merge them through a min-heap, emit in sorted order, stop at the page limit. Never collect the full result set in the coordinator โ one range("", "") would OOM it.
cursor = base64(last_word_returned) next page: range(cursor, b) exclusive of cursor // stateless, resumable, survives coordinator restart
Defend it: a key-based cursor is stable under concurrent inserts, whereas an offset-based one silently skips or repeats words when the corpus changes mid-pagination.
Space: prefix compression is the one optimization worth naming
Sorted leaves mean adjacent words share long prefixes (aa, aaa, aabโฆ). Store each as (shared_prefix_len, suffix) โ front-coding โ with a full key every 16 entries so you can still binary-search within a block. Typical dictionary corpora compress 3โ5ร this way, which is the difference between leaves fitting in page cache and not. Cost: a decode step per entry, negligible against the SSD read it saves.
5. Follow-ups โ answers to have ready
How do you handle deletes during a scan?
MVCC: reads take a snapshot version at query start and see a consistent view for the whole scan, including across pages if you pin the snapshot in the cursor. Deletes write a tombstone that compaction later reclaims. Without snapshots, a long paginated scan can return a word that was deleted an hour earlier and miss one inserted before it started, which is confusing rather than merely stale.
A shard dies mid-query.
Each shard is a replica group (3 replicas, Raft or primary/backup); the coordinator retries the iterator against another replica from the last returned key, which is safe because the cursor is key-based and the scan is idempotent. The partial results already streamed to the client are still valid and in order โ you just resume. If a whole group is unavailable, return an explicit partial result with the missing range flagged rather than silently returning an incomplete set; silently-incomplete is the worse failure for a search-like API.
Unicode? Locale-specific collation?
"Lexicographic" is underspecified the moment you leave ASCII. Byte-order on UTF-8 is not the same as human alphabetical order in most languages, and locales genuinely disagree (in Swedish, รค sorts after z). Practical answer: store an ICU collation key alongside each word and range-partition on that, with the collation locale as part of the index identity โ one index per locale you support. It's a real cost and worth surfacing as a clarifying question early rather than a surprise later.
Would you use an existing system?
Yes, and saying so is the senior answer. Bigtable, HBase, and CockroachDB are all range-partitioned, sorted-key stores that give you exactly this: ordered scans, dynamic tablet splitting, replication. I'd build on one of those and spend my effort on the collation and pagination semantics. Building a distributed B+ tree from scratch is only justified if the workload has a property those don't serve โ and I'd want to name that property before signing up for it.
6. Numbers & recap
Sizing
- 10ยนโฐ words ร 20 B = 200 GB raw; ~50 GB after front-coding.
- ~10 shards at 5 GB each, ร3 replicas = 30 nodes. Small.
- B+ tree fanout ~200 โ depth 5 for 10ยนโฐ keys; internal nodes ~1 GB, cached in RAM, so a lookup is one SSD read.
- Scan of 1000 results โ a few contiguous leaf pages โ < 5 ms.
30-second recap
Range queries need sorted, contiguous data, so I range-partition rather than hash-partition โ accepting hot shards as the explicit cost. Each shard is a replicated B+ tree with internal nodes in RAM and prefix-compressed, linked leaves on SSD, so a range is one descent plus a sequential leaf walk. A coordinator consults a boundary registry, opens iterators only on overlapping shards, and k-way merges them into a streamed, key-cursor-paginated response โ never materializing the full result. Skew is handled by choosing boundaries from a sample of the real corpus and splitting tablets dynamically on size or QPS, with read replicas absorbing query skew. The things I'd pin down first are the interval semantics โ the examples imply half-open โ and the collation, because lexicographic order is locale-dependent the moment the corpus isn't ASCII.
Local business search service condensed
็ๅฎถ senior, 2024-09. "่ฎพ่ฎกไธไธช local business search service." Given API params: geolocation (lat, lng), radius.
Explicitly a senior-level prompt where the details are yours to elicit โ the interviewer gave the minimum and expected the candidate to drive.
1. Requirements โ what to elicit
Functional
search(lat, lng, radius, query?, filters?)โ ranked businesses.- Filters: category, open-now, rating, price band.
- Ranking blends distance, relevance, and quality โ ask which dominates; it's a product decision.
- Business data ingest: owner edits, third-party feeds, user-suggested corrections.
Ask early: is this "find me sushi nearby" (text + geo) or "what's within 2 km" (pure geo)? The first needs an inverted index; the second only needs a spatial index. Most candidates assume one and design the wrong system.
Non-functional
- ~10โธ businesses globally; ~10โต QPS peak, heavily skewed to dense metros.
- p99 < 200 ms โ it's an interactive search box.
- Read-dominated ~10 000:1. Business data changes slowly; open-now and busy-ness change constantly.
- Stale results acceptable for hours on descriptions, not on permanent closure.
Core tension: geo-filtering and text-relevance want different index structures, and joining them at query time is expensive. The design is about doing the cheap filter first and the expensive scoring on a small candidate set.
2. Geo indexing โ the part they're actually testing
| Approach | How radius search works | Trade-off |
|---|---|---|
| Geohash | Prefix = bounding box. Query the covering cell + 8 neighbors, then filter by exact distance. | Simple, string-prefix friendly. Boundary artifacts: two nearby points can differ at the first character, hence the neighbor scan. |
| S2 cells (pick this) | Hilbert curve โ 64-bit cell IDs. A radius becomes a small set of variable-level cell ranges (S2RegionCoverer), each a contiguous ID range. | Better locality, no pole/meridian distortion, and ranges map directly onto a sorted-key store. It's also what Google actually uses, which lands well here. |
| Quadtree / R-tree | Tree descent to the query rectangle. | Adapts to density, but tree rebalancing under writes and awkward to shard across machines. |
| PostGIS / naive lat-lng box | WHERE lat BETWEEN โฆ AND lng BETWEEN โฆ | Fine to ~10โถ rows on one box. At 10โธ with 10โต QPS it falls over โ say why, then discard it. |
The move to state: S2 turns a 2-D proximity problem into 1-D range scans over sorted integers, which is something a distributed store can shard and serve. That sentence is most of the value of this section.
3. Architecture
Cheap filters first: S2 cell ranges cut 10โธ businesses to a few thousand candidates before anything expensive runs. Intersect with the text posting list, apply attribute filters, then two-stage rank. Hydration of full business documents happens last, on ~20 IDs.
4. Trade-offs worth stating
Density skew is the defining problem
Manhattan has orders of magnitude more businesses per kmยฒ than rural Montana, and far more queries. A fixed S2 level is therefore wrong everywhere: at level 13 a Manhattan cell holds thousands of businesses, and a Montana cell holds none.
- Variable-level cells: subdivide until a cell holds โค K businesses. Dense areas get deep cells, sparse areas shallow ones.
S2RegionCovererhandles mixed levels natively โ that's the reason to prefer S2 over flat geohash. - Adaptive radius: if a small radius in a dense area already returns hundreds of results, don't expand. If a rural query returns two, expand the radius progressively โ but tell the user you did, or "nearest" results 40 km away look like a bug.
- Shard by cell prefix but balance by load, with hot metro cells split across more replicas. Geographic sharding without load balancing puts all of Tokyo on one machine.
Two-stage ranking, and why
L1 (cheap, on thousands of candidates):
score = w1*exp(-dist/d0)
+ w2*log(1+review_count)*rating
+ w3*text_match_bm25
โ keep top 200
L2 (expensive, on 200):
GBDT / neural with personalization,
popularity-at-this-hour, click history
โ top 20
Running the expensive model on every candidate blows the 200 ms budget; running only the cheap one loses quality. The staged funnel is the standard answer and stating the candidate-set sizes is what makes it credible.
Defend it โ freshness split. Static attributes (name, category, location) go in the nightly-rebuilt index. Volatile ones (open now, temporary closure, live busy-ness) must not be baked into the index โ compute open-now at query time from cached hours plus the local timezone, and keep a small, fast-updating overlay for closures. Baking hours into the index means a business shows as open at 3 a.m. until the next rebuild.
5. Follow-ups โ answers to have ready
Why not just PostGIS with a GiST index?
For 10โถ businesses and modest QPS, that is the right answer and I'd say so โ don't build a distributed index you don't need. It stops working at 10โธ rows and 10โต QPS with a text-relevance join, because you can't shard a single PostGIS instance by geography without building the coordinator layer anyway, and the ranking model doesn't belong in the database. The migration path is: start with PostGIS, move the geo index to S2-over-a-sorted-store when the metro shards get hot.
A business permanently closes. How fast does it disappear?
Fast, via the delta path, not the nightly rebuild. Closure writes a tombstone into a small, always-consulted overlay that the query service intersects out of results โ seconds, not hours. This is the one data change where staleness is genuinely user-harmful (someone drives to a closed restaurant), so it gets its own fast path while descriptions and photos ride the batch rebuild.
Query at a cell boundary โ do you miss a business 10 m away in the next cell?
No, because the coverer generates the cells covering the whole circle, not the cell containing the center โ the covering inherently spans boundaries. Then every candidate gets an exact haversine distance check, so cell coverage only ever over-fetches, never under-fetches. That "cover generously, filter exactly" pattern is the correctness argument, and it's worth stating explicitly since boundary handling is the classic geohash bug.
Hot metro traffic melts a shard.
Read replicas first, since this is 10 000:1 read-dominated and replicas are cheap. Then a query-result cache keyed on quantized inputs โ snap lat/lng to a ~100 m grid and radius to buckets, so nearby users share a cache entry; that alone gives a very high hit rate in dense areas where queries cluster. Finally split the hot cells to more shards. Notably the index is nearly static, so replicating it aggressively costs storage but no consistency headache.
6. Numbers & recap
Sizing
- 10โธ businesses ร ~2 KB doc = 200 GB. Index postings far smaller.
- Geo index: 10โธ (cell_id, biz_id) pairs โ 1.6 GB โ small enough to replicate widely.
- 10โต QPS ร ~3 KB response = 300 MB/s out. Cache-friendly.
- Budget: coverer ~1 ms, index scan ~20 ms, L1 ~10 ms, L2 on 200 docs ~30 ms, hydrate ~10 ms โ ~70 ms, well inside 200 ms.
30-second recap
The core move is turning 2-D proximity into 1-D range scans: an S2 region coverer converts (lat, lng, radius) into a handful of contiguous cell-ID ranges, which a sorted, range-partitioned index can serve directly. Cells are variable-level so a dense metro subdivides deeper than farmland, which is why S2 beats a flat geohash here. Covering is generous and then every candidate gets an exact haversine check, so boundaries can over-fetch but never miss. Text queries intersect the cell candidates with an inverted-index posting list, filters apply, then a two-stage ranker โ a cheap linear pass down to 200, an expensive model down to 20 โ keeps us inside 200 ms. Volatile attributes like open-now and permanent closures ride a fast overlay rather than the nightly index rebuild, because a user driving to a closed restaurant is the failure that actually matters. And I'd start on PostGIS if the scale were 10โถ โ this design only earns its complexity at 10โธ.
Quota limiter (Drive-style usage quota) condensed
From the 1p3a Google question bank, L6 System Design cluster. The bank's own framing is the point: clarify whether the requirement is request throttling, per-user/per-tenant resource quota, or storage-usage enforcement before choosing counters, windows, and reconciliation semantics. Guessing wrong here means you design a rate limiter when they wanted an accounting system.
1. Requirements โ the disambiguation is the first grade
| Rate limiting | Resource quota | Storage quota (assume this) | |
|---|---|---|---|
| Unit | Requests per window | Concurrent slots (VMs, connections) | Cumulative bytes stored |
| Resets? | Yes, every window | On release | Never โ it's a running balance |
| Error of over-count | Self-heals next window | Self-heals on release | Permanent drift โ needs reconciliation |
| Right primitive | Token bucket in Redis | Semaphore / lease | Durable counter + async audit |
| Failure stance | Fail open | Fail closed | Fail open on read, closed on write |
Say all three out loud, pick one, and justify: "Drive-style implies cumulative storage, so I'll design that โ it's the hardest of the three because errors accumulate forever rather than washing out."
Functional (storage quota)
reserve(user, bytes)before an upload;commitorreleaseafter.- Quota is hierarchical: org โ team โ user, and the tightest binding limit applies.
- Deletes and trash restore adjust usage; trash counts against quota until purged (a real product decision โ state it).
- Users see accurate usage, broken down by category.
Non-functional
- 10โน users; ~10โต quota ops/sec; p99 < 50 ms on the upload path.
- Accuracy over availability on the write path โ the inverse of a rate limiter. Letting a user exceed quota costs real money and is hard to claw back.
- Displayed usage may lag by seconds; enforced usage may not drift permanently.
Core tension: you need a durable, monotonic, per-user counter that survives crashes and never double-counts, on a path that also needs to be fast. That's a transaction, not a cache.
2. Design โ reserve / commit with reconciliation
The two-phase sequence
reserve(user, bytes, op_id):
BEGIN
row = SELECT used, reserved, limit
FROM quota WHERE user=? FOR UPDATE
if row.used + row.reserved + bytes > row.limit:
ROLLBACK; return 507 Insufficient Storage
INSERT INTO reservation(op_id, user, bytes,
expires_at=now+1h)
ON CONFLICT (op_id) DO NOTHING -- idempotent
UPDATE quota SET reserved = reserved + bytes
COMMIT
commit(op_id): -- after bytes are durable in blob store
BEGIN
r = DELETE FROM reservation WHERE op_id=? RETURNING *
if r is null: return OK -- already committed
UPDATE quota SET used = used + r.bytes,
reserved = reserved - r.bytes
COMMIT
sweeper (every 5 min):
expire reservations past expires_at -> release
Why reserve-then-commit, not just increment
A single increment at upload start over-counts every abandoned upload; a single increment at upload end lets a user start a thousand concurrent uploads that each individually fit under the limit and blow past it together. The reservation makes the check and the claim atomic, which is the actual race in this problem.
op_id makes both phases idempotent, so a client retry after a timeout can't double-charge. The expiry sweeper is what stops a crashed uploader from permanently stranding quota โ without it, reserved-but-never-committed bytes leak until the user is locked out of their own account with no way to recover.
Why a transactional store (Spanner/SQL) rather than Redis: this counter is the source of truth for something users pay for. It must survive a cache flush, support a real transaction across the reservation and the balance, and be auditable. This is the exact inverse of the rate-limiter argument in panel 1 โ and being able to explain why the same-looking problem gets the opposite answer is the whole point of having both.
3. Trade-offs worth stating
Hot row on shared quotas
An org-level quota is one row that every member's upload contends on with FOR UPDATE. A 10 000-person org serializes on it.
- Sharded counters: split the org balance into N sub-rows; a reserve picks one at random, and only when a sub-row is exhausted does it consult siblings. Reads sum all N. Cost: the "am I over?" check is approximate near the boundary, so keep a small reserve buffer.
- Lease blocks: a service instance leases 1 GB of org quota and hands out sub-allocations locally. Same pattern as the rate limiter's token lease, but here the lease must be durable and reclaimable, not best-effort.
Drift, and the audit job
Because errors are permanent, you need a periodic reconciler that recomputes true usage from the authoritative object metadata (sum of object sizes per owner) and compares it to the counter. Discrepancies get logged, and small ones auto-corrected; large ones page, because a systematic drift means a bug that's silently mischarging users.
Defend it: the reconciler is scanning billions of objects, so run it incrementally โ per user, on a rolling schedule, prioritized by users near their limit and by accounts with recent anomalies. A full-fleet nightly recompute doesn't scale and isn't necessary; what matters is that every account gets reconciled on some bounded horizon.
4. Follow-ups โ answers to have ready
Quota service is down. Can users upload?
Split by direction. Reads and display of usage fail open with a cached value โ showing a slightly stale number is harmless. Writes fail closed for users near their limit and open for users far below it: if the last known usage is under, say, 80%, admit the upload and reconcile later, because the worst case is a small, correctable overage. If they're near the limit, reject with a retryable error. That graded stance is much better than a blanket answer, and it's exactly the kind of "which is right for this system" reasoning the L6 rubric asks for.
A user deletes 10 GB. When does their quota free up?
Depends on the product decision about trash, which you should surface rather than assume: if deleted files sit in trash for 30 days and remain restorable, they must still count, otherwise you've promised storage you can't guarantee on restore. So the counter decrements at purge, not at delete, and the UI must say so โ this is the single most common user complaint about storage quotas and naming it shows product sense.
Deduplication โ two users store the same file. Who pays?
Both, for their logical usage, even though you store one copy. Charging by physical bytes would leak information (usage dropping tells you someone else has the same file) and produce a quota that changes without the user doing anything. So logical accounting for quota, physical accounting for cost โ and keeping those two ledgers separate is the right architecture.
Why is this different from the rate limiter you designed earlier?
Because the error term behaves differently. A rate limiter's over-count vanishes at the next window, so you can trade accuracy for latency with local leases and Redis, and fail open. A storage quota's over-count is permanent and monetary, so it needs a durable transactional counter, idempotent two-phase updates, an expiry sweeper, and a reconciliation job โ and it fails closed at the boundary. Same shape, opposite answers, and the reason is entirely about whether the error self-heals.
5. Numbers & recap
Sizing
- 10โน users ร ~200 B quota row = 200 GB. Sharded SQL/Spanner, easily.
- 10โต quota ops/s, each a short transaction on one row โ ~2 000 rows/s/shard across 50 shards. Comfortable.
- Reservations in flight: 10โต ops/s ร ~60 s mean upload = 6 M open rows. Small, and the sweeper keeps it bounded.
- Reconciler: 10โน users on a 7-day rolling horizon โ 1 650 accounts/s of background scan.
30-second recap
First I'd pin down which quota this is, because request throttling, concurrency slots, and cumulative storage need genuinely different machinery. Assuming Drive-style storage: the counter is a durable, transactional per-user balance, not a cache, because an over-count here is permanent and monetary rather than washing out at the next window. Uploads take a reserve-then-commit path โ one transaction atomically checks the limit and claims the bytes, keyed by an operation ID so retries are idempotent, and an expiry sweeper releases reservations from crashed uploads so quota can't leak. Shared org quotas would serialize on one hot row, so I'd shard the balance into sub-counters or hand out durable lease blocks, accepting approximate enforcement near the boundary in exchange for concurrency. Because errors accumulate, there's a rolling reconciler that recomputes true usage from object metadata and pages on systematic drift. And the failure stance is graded rather than binary: fail open for users well below their limit, closed for users near it.
Distributed key-value store + throughput estimation condensed
Reported in an earlier Google SD round (the "System Design ๆ็ป" thread): design a distributed KV system, with the interviewer pushing on throughput estimation. That second half is the tell โ this round is at least as much about capacity math out loud as about architecture.
1. Requirements
Functional
get(key),put(key, value),delete(key); optionallycas(key, value, expected_version).- Keys โค 256 B, values โค 1 MB.
- Optional TTL per key.
- Ask: do we need range scans? If yes the partitioning answer flips from hash to range โ see panel 7.
Non-functional โ pin these numerically
- 10โถ QPS, 90/10 read/write.
- p99 < 10 ms read, < 20 ms write.
- 10 TB of data, growing 2ร/year.
- Durability: no acknowledged write may be lost. Availability target 99.99%.
- Consistency: ask, don't assume. Linearizable, read-your-writes, or eventual? This single answer determines the entire replication design.
Core tension: the CAP choice isn't philosophical here, it's a product question โ a session store can be eventually consistent and stay up through a partition; a counter or a lock service cannot. Make the interviewer pick, then design decisively.
2. The throughput math โ do this out loud, early
From QPS to node count
Reads: 9e5 QPS. Writes: 1e5 QPS.
Per-node capability (commodity, NVMe):
memory hit ~200k ops/s
NVMe random read ~100k IOPS, ~100 us
network 10 Gbps = 1.25 GB/s
Data: 10 TB. Cache budget: 10% hot = 1 TB RAM.
at 128 GB/node -> 8 nodes just to hold cache
say 20 nodes -> 500 GB data + 50 GB RAM each
Read path:
90% cache hit -> 8.1e5 ops/s from RAM
10% disk -> 9e4 IOPS spread over 20 nodes
= 4.5k IOPS/node (of ~100k) OK
Write path with RF=3, quorum W=2:
1e5 client writes -> 3e5 physical writes
= 15k writes/s/node, ~1 KB each = 15 MB/s
plus WAL fsync: batch-commit groups of ~100
-> 150 fsync/s/node, trivially fine
Network per node:
(8.1e5/20)*1KB read + replication traffic
~ 40 MB/s + 45 MB/s = 85 MB/s of 1250 MB/s
-> network is NOT the bottleneck; RAM is.
What the math is for
The conclusion โ RAM for the working set, not IOPS or network, sets the node count โ is the thing to say. Any candidate can propose consistent hashing; few can tell you which resource actually binds, and that's what "throughput estimation" is probing.
State the assumptions as you go and keep them round: 100k IOPS, 200k memory ops/s, 10 Gbps, 1 KB average value. Nobody is checking your arithmetic to two decimals; they're checking that you know which numbers matter and roughly how big they are.
Then use it. The math should change a decision: here it says 20 nodes is right, that we're memory-bound so a bigger cache tier is the cheapest lever, and that RF=3 quorum writes cost us nothing we can't afford. Math that doesn't change a decision is decoration.
3. Architecture
Consistent hashing with virtual nodes (~256 per physical node) so adding a machine moves ~1/N of the data instead of half the ring, and so a heterogeneous fleet can be weighted by capacity. Replication factor 3 placed across racks and zones โ replica placement is where availability actually comes from, not the replication factor itself.
4. Trade-offs worth stating
Quorum, and the honest caveat
W + R > N guarantees the read set intersects the write set, so a read sees the latest acknowledged write. With N=3: W=2/R=2 is the balanced default; W=3/R=1 favors read-heavy workloads at the cost of write availability; W=1/R=1 is fast and eventually consistent.
Defend it: quorum is not linearizability. Concurrent writers can still produce conflicting versions, a failed write that reached one replica may later surface, and read-repair timing is not deterministic. If the interviewer wants true linearizability, say so plainly and switch to Raft/Paxos per partition with reads served by the leader (or via lease reads) โ that costs you a leader election window on failover, which quorum-only designs don't have. Naming this distinction is a strong senior signal; asserting "quorum gives strong consistency" is a common and visible error.
Conflicts, hot keys, LSM cost
- Conflict resolution: last-write-wins on a timestamp is simple and silently loses data under clock skew. Vector clocks / version vectors preserve causality but push merge logic to the client. State the choice โ LWW is defensible for a cache-like store, unacceptable for a shopping cart.
- Hot key: hashing spreads keys, not traffic to one key. One viral key still lands on 3 replicas. Fixes: client-side caching with short TTL, or replicating that key to extra nodes on demand. Say that consistent hashing doesn't solve this โ many candidates think it does.
- LSM trade: writes are sequential and fast; reads may touch multiple levels (mitigated by bloom filters), and compaction consumes background IO and causes p99 spikes. If the workload were read-dominated with in-place updates, a B-tree would be the better pick โ and knowing when not to use an LSM is the point.
5. Follow-ups โ answers to have ready
A node dies. What happens to writes targeting it?
With W=2 of N=3, writes still succeed on the surviving two โ availability is preserved by construction. The coordinator stores a hinted handoff for the dead replica and replays it when the node returns. If the node is gone for longer than the hint window, an anti-entropy repair using Merkle trees reconciles the ranges by exchanging hashes rather than data, so only the differing subranges are shipped. The cost of getting this wrong is silent permanent divergence, which is why the repair job is not optional.
How do you add a node without a latency spike?
Virtual nodes mean the new machine claims ~1/N of the ranges from many existing nodes rather than a contiguous half from one neighbor, so the migration is spread and no single donor saturates. Stream ranges in the background at a throttled rate, serve reads from the old owner until a range is fully transferred, then flip ownership atomically in the membership version. The throttle is the important knob โ an unthrottled rebalance is itself an outage, and I'd bound it to a fraction of each node's disk and network budget.
Read-your-writes for a specific user?
Three options, cheapest first: route that user's requests to the same coordinator and replica set via sticky routing; or have the client carry the version it last wrote and require the read to see at least that version, retrying elsewhere if not; or use W=3. The middle option โ a monotonic version token in the client โ is the general and correct one, and it's the same mechanism as a session consistency token in Spanner or DynamoDB.
Where does this design stop working?
At 2ร growth per year, the memory-bound conclusion holds until the hot working set stops fitting economically in RAM โ around 10ร current size. At that point the lever isn't more nodes, it's a tiered cache with a dedicated hot tier, or accepting a lower cache hit rate and re-doing the IOPS math. The other break point is multi-region: everything above assumes one region, and crossing regions forces an explicit choice between async replication with conflict handling and synchronous quorums at ~100 ms. I'd make that a separate design conversation rather than hand-waving it.
6. Numbers & recap
The headline figures
- 10 TB ร RF 3 = 30 TB physical; 20 nodes ร ~2 TB NVMe. Comfortable.
- Hot set 10% = 1 TB RAM across the fleet โ 50 GB/node โ this is the binding constraint.
- Reads: 9e5 QPS, ~90% from memory; disk load only ~4.5k IOPS/node of ~100k available.
- Writes: 1e5 ร RF 3 = 3e5/s = 15k/s/node โ 15 MB/s, with batched fsync at ~150/s.
- Network ~85 MB/s/node against 1.25 GB/s โ an order of magnitude of headroom.
30-second recap
Consistent hashing with a few hundred virtual nodes per machine spreads both data and rebalancing cost, and replication factor 3 placed across racks and zones is where availability actually comes from. Each node runs an LSM engine โ write-ahead log, memtable, SSTables with bloom filters โ which suits a write-heavy path at the cost of read amplification and compaction-driven p99 spikes. Quorum with W=2, R=2 of N=3 means a read intersects the latest acknowledged write, but I'd be explicit that this is not linearizability: for that you'd want Raft per partition with leader reads, trading a failover window for stronger guarantees. Failed nodes are covered by hinted handoff and Merkle-tree anti-entropy repair. On the sizing: at a million QPS with 10 TB and a 10% hot set, it's memory that sets the node count at around twenty โ disk IOPS and network both have an order of magnitude of headroom โ so the cheapest lever if we need more throughput is a bigger cache tier, not more machines.