Scaling the System for 2x Concurrent Users
Problem The platform's concurrent user load doubles. Describe how you would scale the system to absorb it without degrading latency or availability.
Functional requirements
- Sustain all existing functionality (browse, search, order, track, pay) at 2x load.
- Identify and remove the binding bottleneck rather than uniformly over-provisioning.
- Degrade gracefully rather than failing hard if 2x is exceeded.
- Validate the design under load before rollout.
- Keep the change incrementally deployable and reversible.
Non-functional requirements
- Current: ~20M DAU, ~1M concurrent at peak, ~40k QPS read / ~1.5k orders/sec. Target: ~2M concurrent, ~80k QPS read / ~3k orders/sec.
- Latency budget unchanged: browse p99 < 300 ms, order placement p99 < 500 ms.
- Availability unchanged at 99.95%.
- Primary DB currently ~12k writes/sec against a ~15k/sec ceiling — this is the binding constraint, and it will break first at 2x.
- Cache hit rate currently ~85%; a drop to 70% would triple origin DB load, so cache sizing must scale with the working set.
Key components
- Bottleneck analysis first: app server CPU and connection pools, DB read/write throughput, cache hit rate, and downstream third parties (payments, delivery-partner services) that will NOT scale just because you asked.
- Stateless app/API tier scaled horizontally behind the load balancer, with autoscaling keyed on request latency and CPU.
- Read scaling: additional read replicas for browse/search/menu traffic; expand Redis for hot restaurant listings and menus to keep origin load flat.
- Write scaling: shard the write-heavy order tables by region or restaurant_id once a single primary tops out.
- Async offload: move notifications, analytics, and non-critical order-status fan-out onto Kafka so the synchronous request path shortens.
- CDN/edge caching for static assets and cacheable listing responses.
- Load testing at 2x plus capacity headroom, and graceful degradation paths.
Deep dives / trade-offs
- Where 2x actually breaks: doubling users does not double every subsystem uniformly. Reads scale with replicas and cache almost trivially; the write primary does not. Show the arithmetic that identifies the order-write path as the failure point, and resist the instinct to just add app servers — that only moves the queue.
- Read replicas introduce replication lag, so a user who just placed an order may not see it (read-your-writes violation). Route post-write reads to the primary or pin a session to the primary for a few seconds.
- Sharding orders is a one-way door: it breaks cross-shard reporting queries and multi-entity transactions. Compare it against first exhausting vertical scaling, connection pooling (PgBouncer), and moving cold orders to an archive — often 2x is reachable without sharding at all, and that is the better answer.
- Cache is the highest-leverage lever and the most dangerous: an 85%→95% hit rate cuts origin load by 3x for the price of memory, but a cache-tier failure or a cold start after deploy now delivers 100% of 2x traffic to a DB sized for 15% of 1x. Discuss request coalescing, staggered TTLs with jitter to prevent synchronized expiry, and whether the DB can survive a cold cache at all.
- Downstream limits: the payment provider's rate limit doesn't double on request. Queue, backpressure, and negotiate — or the checkout path fails at exactly peak.
- Graceful degradation: serving slightly stale restaurant listings under extreme load is nearly free and invisible; failing checkout is not. Rank what you shed.
- Validate with load tests against production-shaped data before rollout — a design that holds on a whiteboard and dies on a hot partition is common.
asked …