Skip to content

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:

  1. Like / Unlike Operations

  2. Users can like and unlike:

    • Posts
    • Reels
    • Comments
    • Each user can like an item at most once
    • Repeated likes/unlikes should be idempotent
  3. Counter Retrieval

  4. Fetch like counts for:

    • Single item
    • Bulk items (e.g., feed with 50 posts)
    • Reads must be low-latency
  5. Real-Time Feedback

  6. UI should update immediately after a like/unlike

  7. Backend counts may be eventually consistent

  8. APIs

  9. Like API

  10. Unlike API
  11. GetCount API (single & batch)
  12. Optional: HasUserLiked API

  13. Scope & Extensibility

  14. System should support:

    • Other counters (views, shares, saves)
    • Future aggregation (daily likes, per-region likes)

Non-Functional Requirements

Your design must satisfy:

  1. Scale

  2. Tens of billions of items

  3. Hundreds of millions of active users
  4. Peak write throughput: millions of likes/unlikes per second
  5. Peak read throughput: tens of millions of reads per second

  6. Latency

  7. Like/Unlike P99 ≤ 50 ms

  8. Read P99 ≤ 20 ms

  9. Availability

  10. ≥ 99.99% uptime

  11. Must survive:

    • Data center failures
    • Network partitions
  12. Consistency

  13. Eventual consistency is acceptable

  14. Temporary over/under-counting is acceptable
  15. No permanent count corruption

  16. Durability

  17. Likes must not be permanently lost

  18. Counters must be recoverable after crashes

What You Should Deliver

Provide a production-grade system design that includes:

  1. Requirement clarification & assumptions

  2. Accuracy guarantees

  3. Acceptable error bounds
  4. Read-after-write expectations

  5. High-Level Architecture

  6. Core services

  7. Data flow:

    • Like → write path
    • Read → aggregation path
  8. Data Modeling

  9. How likes are stored

  10. User-like relationships
  11. Counter representation

  12. Counter Update Strategy

  13. Synchronous vs asynchronous updates

  14. Handling extreme write contention
  15. Idempotency and deduplication

  16. Storage Choices

  17. Hot counters vs cold storage

  18. In-memory vs persistent stores
  19. Sharding and partitioning strategy

  20. Caching Strategy

  21. What is cached

  22. Cache invalidation/update model
  23. Handling hot keys (viral posts)

  24. Scalability & Performance

  25. Horizontal scaling approach

  26. Write amplification control
  27. Hot partition mitigation techniques

  28. Failure Handling

  29. Retry strategy

  30. Data recovery
  31. Handling partial writes and duplicates

  32. Approximation & Optimization (If Any)

  33. Probabilistic counters (if used)

  34. Batching and aggregation windows
  35. Trade-offs between accuracy and performance

  36. Capacity Estimates

    • Likes per second (average & peak)
    • Storage growth
    • Memory usage for hot counters
  37. 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:

  1. 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.)
  2. 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.)
  3. 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.)
  4. 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.