Distributed counter
System Design Task: Instagram-Scale Distributed Like Counter¶
Problem Statement¶
Design a highly scalable, fault-tolerant distributed counter system to support Instagram-scale “Like” counts on posts, reels, comments, and stories.
The system must handle massive write throughput, provide fast reads, and maintain acceptable accuracy under extreme concurrency, network partitions, and regional failures.
This counter system will be used across multiple products and surfaces, so it must be generic, reusable, and production-ready.
Functional Requirements¶
Your system must support:
-
Like / Unlike Operations
-
Users can like and unlike:
- Posts
- Reels
- Comments
- Each user can like an item at most once
- Repeated likes/unlikes should be idempotent
-
Counter Retrieval
-
Fetch like counts for:
- Single item
- Bulk items (e.g., feed with 50 posts)
- Reads must be low-latency
-
Real-Time Feedback
-
UI should update immediately after a like/unlike
-
Backend counts may be eventually consistent
-
APIs
-
Like API
- Unlike API
- GetCount API (single & batch)
-
Optional: HasUserLiked API
-
Scope & Extensibility
-
System should support:
- Other counters (views, shares, saves)
- Future aggregation (daily likes, per-region likes)
Non-Functional Requirements¶
Your design must satisfy:
-
Scale
-
Tens of billions of items
- Hundreds of millions of active users
- Peak write throughput: millions of likes/unlikes per second
-
Peak read throughput: tens of millions of reads per second
-
Latency
-
Like/Unlike P99 ≤ 50 ms
-
Read P99 ≤ 20 ms
-
Availability
-
≥ 99.99% uptime
-
Must survive:
- Data center failures
- Network partitions
-
Consistency
-
Eventual consistency is acceptable
- Temporary over/under-counting is acceptable
-
No permanent count corruption
-
Durability
-
Likes must not be permanently lost
- Counters must be recoverable after crashes
What You Should Deliver¶
Provide a production-grade system design that includes:
-
Requirement clarification & assumptions
-
Accuracy guarantees
- Acceptable error bounds
-
Read-after-write expectations
-
High-Level Architecture
-
Core services
-
Data flow:
- Like → write path
- Read → aggregation path
-
Data Modeling
-
How likes are stored
- User-like relationships
-
Counter representation
-
Counter Update Strategy
-
Synchronous vs asynchronous updates
- Handling extreme write contention
-
Idempotency and deduplication
-
Storage Choices
-
Hot counters vs cold storage
- In-memory vs persistent stores
-
Sharding and partitioning strategy
-
Caching Strategy
-
What is cached
- Cache invalidation/update model
-
Handling hot keys (viral posts)
-
Scalability & Performance
-
Horizontal scaling approach
- Write amplification control
-
Hot partition mitigation techniques
-
Failure Handling
-
Retry strategy
- Data recovery
-
Handling partial writes and duplicates
-
Approximation & Optimization (If Any)
-
Probabilistic counters (if used)
- Batching and aggregation windows
-
Trade-offs between accuracy and performance
-
Capacity Estimates
- Likes per second (average & peak)
- Storage growth
- Memory usage for hot counters
-
Trade-offs & Design Decisions
- What accuracy guarantees are relaxed
- Why certain technologies are chosen
- What is explicitly not optimized
Expectations¶
- Be concrete and realistic
-
Name specific data structures and system components
-
(e.g., Redis, log-based ingestion, sharded counters, write-behind cache)
- Avoid academic-only solutions
- Focus on operability, debuggability, and long-term maintenance
- Assume this system will be used by multiple teams and products
Interview Kit¶
Read first: Solution · Caching §4 stampedes, Case 6 hot key · Sharding §6.4 hot partitions · Load control
Curveballs. The interviewer changes one thing mid-design. The hint in italics is what a strong answer reaches for:
- A celebrity post gets 2M likes in 60 s. Which single key melts, and what do you change? (Counter sharding for that item only, local write aggregation, L1 cache for reads.)
- A bot farm adds 5M fake likes over a week, then gets banned. How do counts heal without manual edits? (Replay unlikes through the normal idempotent path.)
- Redis loses its primary, and the replica missed the last 2 s of increments. What is the source of truth, and how do you reconcile? (Durable like edges plus periodic recount, and the counter is a cache of them.)
- Product wants exact counts for posts under 1,000 likes and "1.2M" style above. Does that simplify anything?
Must answer (security, privacy, operations):
- Who can read who liked an item, and how blocks and private accounts apply to batch
GetCount - What happens to a deleted user's likes (GDPR erasure, counters decremented idempotently)
Phase it (MVP → Growth → Scale): MVP: one Postgres table of like edges, unique on (user, item), plus count(*) cached in Redis. Growth: async counter updates through Kafka. Scale: sharded counters for hot items and a multi-region merge.
Score yourself with the rubric.