System Design Problem

Design Like Count for High Profile Users

Commonly Asked By:TwitterMetaByteDanceLinkedIn

Interview Setup

Interview Prompt

Design a like system that handles both normal posts and celebrity posts with 100M+ likes. Users can like/unlike, see counts, check if they liked something, and view recent likers.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Do displayed counts need to be exact, or is approximate OK?Allows the system to use eventual consistency and formatted numbers like 1.5M instead of locking for exact counts.
What's the peak like rate on a viral celebrity post?A post hitting 1M likes/sec requires a tiered architecture, whereas normal posts can use a simpler synchronous path.
Does the liker need to see their own like immediately on refresh?Read-your-writes for deduplication state requires a dedicated routing path separate from the aggregate counter.
Do we need to notify the post owner of every like?Celebrity posts require batched notifications handled by an asynchronous Kafka consumer rather than the tap path.

Scope

In scope

  • Like/unlike with idempotent dedup
  • Display counts at celebrity scale (100M+)
  • Tiered architecture for normal vs hot posts
  • Sharded Redis counters
  • Read-your-writes for liker state
  • Cache warming for anticipated viral posts

Out of scope (state explicitly)

  • Live-stream reaction broadcast (covered in Live Likes & Reactions)
  • Comment system or share counts
  • ML-based like prediction or ranking
  • Full notification delivery infrastructure (batching strategy only)

Functional Requirements

Scope static post like counts at celebrity scale rather than live stream reactions, which are covered in the Live Likes & Reactions problem. Clarify whether exact counts or ±1% approximations are acceptable for display.

In an interview, Justin Bieber posting is the canonical hot key example, so bring it up if the interviewer remains vague.

This problem represents an architectural variant of the Live Likes & Reactions system. While live events require transient reaction streaming, static post counters require durable persistence, tiered ingestion for high profile accounts, and dedicated read your writes guarantees.

  • Like and Unlike: Users can like or unlike posts, photos, and videos.
  • Display Count: Show total like counts on each piece of content.
  • Check Liked Status: Indicate whether the current user has already liked a given post.
  • Celebrity Scale: Handle viral posts with 100M+ total likes from high profile accounts.
  • Recent Likers: Display a sample of recent likers along with aggregate counts, such as "Liked by Alice, Bob, and 2.3M others".
  • Like Notifications: Notify content owners of new likes, batching and grouping updates for high profile creators.

Non-Functional Requirements

The system is read heavy with massive write spikes whenever a celebrity posts. Focus on exact per-user state transitions, hot key mitigation for aggregate counters, cache coherency, and durable event processing for billions of daily interactions.

  • Ultra-High Throughput: Sustain 1M+ likes per second on a viral celebrity post.
  • Low Latency: Like actions complete in under 100 ms at p99.
  • Eventual Consistency: Displayed counts can lag by several seconds because approximate numbers suffice for feed rendering.
  • High Availability: Maintain 99.99% availability because likes represent a critical engagement driver.
  • Idempotent Processing: Rapid double-taps or client retries must not generate duplicate likes.
  • Abuse and Retry Protection: Rate limit by user, device, and IP so automated storms are rejected before expensive state transitions and durable event enqueue.
  • Massive Scalability: Support 1B+ active users and posts accumulating over 100M likes.

Capacity Estimations

A single celebrity post can trigger millions of concurrent writes, which requires sizing sharded counters and asynchronous aggregation queues before designing the system topology. The 160 GB/day storage figure is logical payload only and excludes Cassandra replication, indexes, metadata, and compaction overhead.

MetricCalculationValue
Total likes / dayGiven5B
Likes / sec (avg)5B ÷ 86400~58K
Peak likes / sec (viral post)Given viral post peak design target1M+
Avg post likesGiven (typical workload assumption)50
Celebrity post likesGiven (assumption documented in value)10M - 100M
Like record sizeGiven (assumption documented in value)32 bytes (user_id + post_id + timestamp)
Storage / day5B x 32B logical event payload160 GB raw logical payload/day before Cassandra overhead, indexes, compaction, and replication

Architecture Diagram

The architecture separates exact user like state from derived aggregate counts. User scoped Redis state provides immediate read your writes, Kafka provides the durable event path, Cassandra stores durable authoritative state and recent event history, and 8 way counter sharding prevents a single viral post from becoming one Redis hot key. Redis aggregate counts remain repairable projections rather than the source of truth.

This system design focuses on static post likes, featuring tiered write paths for normal and viral posts, sharded counter keys, and dedicated read your writes consistency for the acting user. For live-stream reaction fan out with 500ms WebSocket deltas and animation sampling, refer to the Live Likes & Reactions design.

In an interview setting, mention that sharding hot post counters into N sub-keys prevents single Redis nodes from saturating when high profile creators post.

Loading...

Component Deep Dives

The Hot Counter Problem

Tracking likes requires tracing the complete path from API ingress to displayed client counts. Hot key counter contention and write behind aggregation represent the core architectural pivots. When a high profile user publishes a viral update, 1M likes can arrive within 60 seconds. If every tap executes UPDATE posts SET like_count = like_count + 1, all write transactions target a single row. Row-level locks serialize these updates, exhausting the database connection and worker thread pools while causing cascading timeouts across unrelated services.

Celebrity publishes post: "Ronaldo scores goal" generates 1M likes/sec for 60 seconds

Naive relational approach:
  UPDATE posts SET like_count = like_count + 1 WHERE post_id = ?;

Cascading failure sequence under viral write load:
  1. 1M writes/sec target a single database row.
  2. Row-level lock contention forces database threads to serialize updates.
  3. Database thread and connection pools become completely exhausted.
  4. All other queries on the database slow down, triggering cascading timeouts.

Tiered Architecture: Normal vs Hot Posts

The system addresses hot key contention by segregating posts into distinct ingestion tiers based on write velocity. Normal posts use an exact user scoped Redis state transition and a low latency count path, while the durable event is persisted asynchronously through Kafka and Cassandra. Hot posts switch to sharded Redis count projections, with Kafka consumers batching Cassandra persistence every 5 seconds. Real-time velocity tracking promotes active posts to the hot tier, and a cooldown window prevents rapid tier oscillation.

Tier 1: Normal Posts (< 10K likes)
  Like request atomically updates user scoped Redis like state, durably enqueues the event,
  and serves the count from the normal post counter projection. Cassandra persistence is asynchronous.

Tier 2: Hot Posts (10K+ likes, detected via velocity)
  Like request atomically updates user scoped Redis state, durably enqueues the same event,
  and serves the aggregate count from 8 Redis counter shards. Cassandra consumers batch persistence every 5 seconds.

Hot post velocity detection:
  Track like velocity per post in Redis:
    INCR like_velocity:{post_id}:{minute}
  Promotion rule:
    If velocity exceeds 1,000 likes per minute, mark post_id as warming.
    Capture the current aggregate and Kafka offsets, seed the 8 counter shards, replay events after those offsets,
    then switch reads and new projections to hot mode at the catch-up barrier.
  Demotion rule:
    When velocity stays below 100 likes per minute for a cooldown window, keep sharded writes active,
    materialize one counter from the shard sum, and switch back only after an offset barrier.
  Hysteresis prevents rapid promotion and demotion when traffic fluctuates.

Liked-By Deduplication

Rapid client double-taps and network retry attempts must not inflate aggregate counts. The authoritative decision is an atomic user scoped state transition keyed by user_id, which supports exact like and unlike semantics. For mega-posts, a Bloom Filter can reduce negative membership work, but positive results still require an exact state check because Bloom Filters have false positives. Cassandra stores durable current state in bounded post and user buckets so exact recovery and reconciliation remain possible at celebrity scale.

Exact user-state path:
  Redis stores the acting user's current like state on a user scoped shard.
  Atomic state transition returns:
    +1  -> first like changes false -> true
     0  -> requested state already matches current state
    -1  -> unlike changes true -> false

Bloom Filter prefilter for mega-posts:
  Memory footprint: 100M entries x 10 bits = 125 MB
  Configured false positive rate: ~1%
  Bloom Filter is a prefilter, not the authority:
    EXISTS = false -> definitely not present in the filter, so the exact state check can often be skipped only after backfill is complete and every new like updates the filter before this fast path is used
    EXISTS = true  -> possibly present, so consult exact user state before treating the request as a duplicate
  False positives increase exact-check traffic but cannot silently drop a legitimate first like.

Tiered policy:
  Smaller workloads can use exact user state without a Bloom Filter.
  Mega-posts can add a sharded Bloom Filter for cold membership paths while Cassandra remains the durable source of exact state.
  Because exact user state is user scoped, Bloom Filter membership is never the correctness mechanism.

Sharded Counters

Even with tiered ingestion, a single Redis counter key such as like_count:{post_id} creates a hot key bottleneck at 1M writes per second. Counter sharding partitions the counter across 8 sub-keys. Each state-changing event deterministically maps to a shard using the user ID, which keeps retries and repeated actions for the same user and post consistent. Reads aggregate all 8 sub-keys in parallel and cache the sum briefly to trade a small amount of freshness for much higher write parallelism.

Counter sharding layout across 8 sub-keys:
  like_count:<post_id>:0 = 15234
  like_count:<post_id>:1 = 15189
  ...
  like_count:<post_id>:7 = 15122
  Redis Cluster note: do not use a hash tag such as {post_id} on these counter keys, or all 8 keys would map to the same cluster slot.

Write operation:
  shard = hash(user_id) % 8
  INCRBY like_count:<post_id>:<shard> effective_delta
  A deterministic shard keeps repeated operations for the same user/post on one counter shard.
  With balanced traffic and suitable Redis Cluster placement, sharding can provide roughly 8x aggregate write capacity at the counter layer, though actual gains depend on node placement and workload.

Read operation:
  SUM(like_count:<post_id>:0 through :7) = 121,234
  The application aggregates all 8 sub-keys in parallel and caches the sum locally for 1 to 2 seconds.

Read Your Writes for the Liker

While global aggregate counts can safely lag by several seconds, the user who tapped the like button must immediately see their filled heart icon on page reload. The system provides read your writes consistency by routing both the exact user-state transition and the subsequent read to the same Redis shard determined by a user ID hash tag. The acting user's state is therefore visible immediately after a successful request even while the aggregate count remains eventually consistent.

Write path:
  1. Route the request to the user's home Redis shard using a user_id hash tag.
  2. Compute a request_fingerprint from post_id + action and atomically validate operation_id reuse.
  3. If a prior operation for the same user/post is still pending durable enqueue, return blocked_pending. Do not assign a later sequence while an earlier sequence is pending.
  4. If the requested state already matches current state, record the operation as a committed no-op. Do not emit a count-changing event.
  5. For a state change, assign operation_seq and record:
       like_state:{user_id}:<post_id> -> { liked, operation_seq, operation_id }
       like_op:{user_id}:<operation_id> -> pending|operation_seq|desired_state|effective_delta|request_fingerprint
       pending_like:{user_id}:<post_id> -> operation_id|operation_seq|request_fingerprint
  6. Publish the event with the same operation_id and operation_seq to Kafka or a durable outbox. The producer uses idempotence and an explicit acknowledgment policy before the API reports success.
  7. After durable enqueue succeeds, mark the operation committed and clear pending_like. If the client retries while it is pending, reuse the same operation_id and sequence and re-enqueue the same event. A different request with the same operation_id is an idempotency conflict. A recovery worker must never silently let a reserved operation_seq expire: it either re-enqueues the original event or emits a durable no-op tombstone carrying that sequence before allowing later sequences to advance. The 24-hour operation record is the guaranteed retry horizon.

Read path (check if liked):
  1. Route by the same user_id hash tag to the user's home Redis shard.
  2. Execute the equivalent GET/HMGET for like_state:{user_id}:<post_id>.

Consistency rationale:
  Because the write and immediate read use the same user scoped Redis shard, the acting user's state is visible immediately after the successful request.
  The aggregate count is a separate derived value and may lag.

Fallback path:
  If the user-state cache misses or the shard is unavailable, first check the durable like_operation_by_day ledger for a retried operation_id in the current and previous UTC day buckets. If present and unexpired, reuse its exact sequence and payload. Otherwise route a point lookup to Cassandra:
    like_state_by_post_bucket WHERE post_id = ? AND user_bucket = hash(user_id) % 1024 AND user_id = ?
  If Redis failed after a successful durable enqueue, replay retained Kafka events for the affected user and post before clearing recovered pending state.
  If a pending operation cannot be recovered from its durable ledger or Kafka record, resolve it explicitly with a durable no-op tombstone or compensating rollback before clearing the pending gate. Never abandon a sequence gap silently.
  After Redis shard recovery, initialize the local sequence counter from the recovered highest operation_seq before accepting new state changes.
  The ~5ms value is an illustrative latency assumption, not a guarantee.
LUA
-- Illustrative atomic state transition on the user's Redis shard
-- KEYS[1] = like_state:{user_id}:<post_id>
-- KEYS[2] = like_seq:{user_id}:<post_id>
-- KEYS[3] = like_op:{user_id}:<operation_id>
-- KEYS[4] = pending_like:{user_id}:<post_id>, all sharing the user hash tag
-- ARGV[1] = action (like|unlike)
-- ARGV[2] = operation_id
-- ARGV[3] = request_fingerprint
-- Operation records define a 24-hour retry horizon. A recovery worker must resolve pending operations before expiry so their reserved sequence is not lost.
-- Values are encoded as: status|operation_seq|desired_state|effective_delta|request_fingerprint
local prior = redis.call("GET", KEYS[3])
if prior then
  local parts = {}
  for value in string.gmatch(prior, "[^|]+") do
    parts[#parts + 1] = value
  end
  if parts[5] ~= ARGV[3] then
    return {"idempotency_conflict"}
  end
  return {parts[1], tonumber(parts[2]), tonumber(parts[3]), tonumber(parts[4])}
end

local pending = redis.call("GET", KEYS[4])
if pending then
  local parts = {}
  for value in string.gmatch(pending, "[^|]+") do
    parts[#parts + 1] = value
  end
  return {"blocked_pending", tonumber(parts[2]), parts[1]}
end

local current = redis.call("HGET", KEYS[1], "liked") or "false"
local desired = ARGV[1] == "like"
local operationId = ARGV[2]
local fingerprint = ARGV[3]
local currentBool = current == "true"
local TTL = 86400

if currentBool == desired then
  local existingSeq = tonumber(redis.call("HGET", KEYS[1], "operation_seq") or "0")
  redis.call("SETEX", KEYS[3], TTL,
    "committed|" .. existingSeq .. "|" .. (currentBool and "1" or "0") .. "|0|" .. fingerprint)
  return {"committed", existingSeq, currentBool and 1 or 0, 0}
end

local seq = redis.call("INCR", KEYS[2])
local delta = desired and 1 or -1
redis.call("HSET", KEYS[1],
  "liked", desired and "true" or "false",
  "operation_id", operationId,
  "operation_seq", seq)
redis.call("SETEX", KEYS[3], TTL,
  "pending|" .. seq .. "|" .. (desired and "1" or "0") .. "|" .. delta .. "|" .. fingerprint)
redis.call("SETEX", KEYS[4], TTL, operationId .. "|" .. seq .. "|" .. fingerprint)
return {"pending", seq, desired and 1 or 0, delta}

Cache Warming

Celebrity posts appear in millions of follower feeds before the like spike arrives. Without proactive warming, initial read requests create a cache stampede against persistent storage. The system pre-populates Redis counter keys when a post enters high-velocity tracking or fans out across follower feeds.

Cache Warming Strategy for Viral Celebrity Content:

1. Proactive warming on feed fan out:
   When a celebrity publishes a post and it enters follower feeds, publish post_id to a warm-cache Kafka topic.
   Workers populate the Redis aggregate count before the first read spike.

2. Reactive warming on hot-post promotion:
   When like_velocity:{post_id} crosses 1,000 likes per minute, move post_id to a warming state.
   Capture a verified count snapshot at a projector checkpoint boundary together with the Kafka partition offset map through which that count is complete. Treat the snapshot and offsets as one logical replay barrier.
   Seed the 8 count shards using base = floor(count / 8) and remainder = count % 8, then add one to the first remainder shards.
   Replay events strictly after the captured offsets. Mark the post hot only after the replay reaches the promotion barrier.

3. Demotion after a cooldown window:
   Keep the post sharded until velocity stays below 100 likes per minute for the configured cooldown.
   Drain the projector through a known offset barrier, materialize the single counter from the 8 shard sum, record the matching offset map, and switch the read and write mode only after the new projection has caught up through that same barrier.
   This prevents lost or double-applied deltas during migration.

4. Eviction and TTL policy:
   Retain hot post keys without TTL while active.
   Evict normal posts after 24 hours of inactivity to reclaim memory.

Event Bus (Kafka)

The user facing API can respond after the user's state transition and durable event enqueue succeed. Secondary operations including Cassandra persistence, aggregate count projection, recent liker indexing, creator notification delivery, and analytics ingestion execute asynchronously through independent Kafka consumer groups. This keeps downstream processing off the request path while preserving a durable replay source.

YAML
topic: like-events
partitions: 128
partition_key: post_id + user_bucket # user_bucket = hash(user_id) % 1024. Preserves per-user/post ordering while distributing one viral post across many partition keys
retention_days: 7
replication_factor: 3
min_insync_replicas: 2

producer:
  idempotence: true # enable.idempotence=true prevents duplicate publications from producer retries
  acks: all # successful API response requires the event to be acknowledged by the required in-sync replicas
  event_schema:
    event_id: "uuid"
    operation_id: "string"
    user_id: "uuid"
    post_id: "uuid"
    action: "like | unlike"
    operation_seq: "monotonic per user/post, contiguous for state-changing events"
    effective_delta: "+1 | -1"
    timestamp: "server_assigned_epoch_ms"

consumer_groups:
  db_writer: "Persist authoritative current like state, time-bucketed operation ledger, and recent-like events to Cassandra in batches"
  count_projector: "Apply effective deltas in operation_seq order, hold sequence gaps, and deduplicate retries by operation_id. Redis is a repairable derived projection"
  count_reconciler: "Periodically compare derived counts against authoritative like state and repair drift"
  notification: "Batched push notifications dispatched to post creators"
  analytics: "Sink event stream to ClickHouse for analytics and reporting"

execution_paths:
  sync_path: "Validate request, atomically transition user like state on its home Redis shard, validate operation_id reuse, assign operation_seq for state changes, durably enqueue the event, and return the latest available cached count"
  async_path: "Persist state and event history to Cassandra, update aggregate counter shards, generate notifications, and reconcile drift without blocking the user response"
  fault_handling: "Retry failed consumers with idempotent operation_id handling, route poison events to like-events-dlq after 3 retries, hold sequence gaps for ordered replay, and alert when consumer lag exceeds 30 seconds"

durability_note:
  "The API must not acknowledge a state-changing operation merely because Redis changed. Kafka durability or an equivalent durable outbox must be confirmed before a successful response. Redis provides the low latency read model."

API Design

Service Interfaces and Domain Types

The application layer exposes typed RPC interfaces for setting like state, querying user like status, fetching batch count totals, and retrieving recent liker profiles. Each state changing request carries an operation ID, and the operation ID in the request body must agree with the Idempotency-Key header. Retries can then return the original result instead of creating a second state transition. The 24-hour ledger TTL defines the guaranteed idempotency horizon.

TYPESCRIPT
type PostId = string;
type UserId = string;
type OperationId = string;
type OperationSequence = number;
type LikeCount = number;
type EpochMs = number;
type CursorToken = string;
type PageLimit = number;

export interface LikeRequest {
  postId: PostId;
  action: "like" | "unlike";
  operationId: OperationId;
}

export interface LikeResponse {
  liked: boolean;
  likeCount: LikeCount;
  operationId: OperationId;
  operationSequence: OperationSequence;
  countIsEventuallyConsistent: boolean;
  countReadTimestampMs: EpochMs;
}

export interface CheckLikedResponse {
  liked: boolean;
}

export interface BatchCountRequest {
  postIds: PostId[];
}

export interface BatchCountResponse {
  counts: Record<PostId, LikeCount>;
  countIsEventuallyConsistent: boolean;
  countReadTimestampMs: EpochMs;
}

export interface LikedUserSummary {
  id: UserId;
  name: string;
}

export interface RecentLikersResponse {
  users: LikedUserSummary[];
  totalCount: LikeCount;
  nextCursor?: CursorToken;
  countIsEventuallyConsistent: boolean;
  countReadTimestampMs: EpochMs;
}

export interface LikeCountService {
  toggleLike(userId: UserId, request: LikeRequest): Promise<LikeResponse>;
  checkLiked(userId: UserId, postId: PostId): Promise<CheckLikedResponse>;
  getBatchCounts(request: BatchCountRequest): Promise<BatchCountResponse>;
  getRecentLikers(postId: PostId, limit?: PageLimit, cursor?: CursorToken): Promise<RecentLikersResponse>;
}

Like and Unlike Endpoints

The like endpoint mutates user engagement state idempotently and returns the latest available aggregate like count for client side rendering. The response explicitly marks the count as eventually consistent.

HTTP
POST /api/v1/likes HTTP/1.1
Host: api.example.com
Content-Type: application/json
Idempotency-Key: op-9d0b6a7c

{
  "post_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
  "action": "like",
  "operation_id": "op-9d0b6a7c"
}

HTTP/1.1 200 OK
Content-Type: application/json

{
  "liked": true,
  "like_count": 1523401,
  "operation_id": "op-9d0b6a7c",
  "count_is_eventually_consistent": true,
  "count_read_timestamp_ms": 1778179200000
}

POST /api/v1/likes HTTP/1.1
Host: api.example.com
Content-Type: application/json
Idempotency-Key: op-9d0b6a8e

{
  "post_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
  "action": "unlike",
  "operation_id": "op-9d0b6a8e"
}

HTTP/1.1 200 OK
Content-Type: application/json

{
  "liked": false,
  "like_count": 1523400,
  "operation_id": "op-9d0b6a8e",
  "count_is_eventually_consistent": true,
  "count_read_timestamp_ms": 1778179200000
}

Check If Liked

Verifies whether the authenticated user has previously liked a given post, routing by user ID to guarantee immediate read your writes accuracy.

HTTP
GET /api/v1/likes/check?post_id=9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d HTTP/1.1
Host: api.example.com

HTTP/1.1 200 OK
Content-Type: application/json

{
  "liked": true
}

Get Like Count (Batch)

Enables feed generation services to fetch the latest available like totals for up to several hundred posts in a single round trip. These aggregate values are eventually consistent.

HTTP
POST /api/v1/likes/counts HTTP/1.1
Host: api.example.com
Content-Type: application/json

{
  "post_ids": [
    "post-101",
    "post-102",
    "post-103"
  ]
}

HTTP/1.1 200 OK
Content-Type: application/json

{
  "counts": {
    "post-101": 1523401,
    "post-102": 42,
    "post-103": 89234
  },
  "count_is_eventually_consistent": true,
  "count_read_timestamp_ms": 1778179200000
}

Get Recent Likers

Retrieves the most recent users who currently like a post using a paginated recent event index and exact state checks, along with the latest available aggregate count to support social proof displays.

Recent liker query for a hot post:
  1. Read the newest entries from recent_likers:<post_id>:<event_shard> across all 8 Redis candidate shards in parallel.
  2. Merge candidates by liked_at and keep an opaque cursor containing the last timestamp plus per-shard positions.
  3. Deduplicate candidate user_id values in memory.
  4. Point-check authoritative current state using like_state_by_post_bucket and user_bucket = hash(user_id) % 1024.
  5. Keep only users whose current state is liked. Unlike events remain in history but never appear as current recent likers.
  6. If fewer than the requested page size remain, continue into older hour buckets across the 64 Cassandra event shards until enough users are collected or the seven-day retention horizon is reached.
  7. Return the latest available aggregate count plus the next cursor.

The recent event index is a read optimization. Exact current user state remains authoritative, so stale or duplicate history cannot make an unliked user appear as currently liking the post.
HTTP
GET /api/v1/likes/recent?post_id=9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d&limit=5 HTTP/1.1
Host: api.example.com

HTTP/1.1 200 OK
Content-Type: application/json

{
  "users": [
    { "id": "u-101", "name": "Alice" },
    { "id": "u-102", "name": "Bob" }
  ],
  "total_count": 1523401,
  "next_cursor": "eyJzY2FuX3NldCI6MX0=",
  "count_is_eventually_consistent": true,
  "count_read_timestamp_ms": 1778179200000
}

Common Error Responses

Standardized HTTP error responses returned across like endpoints during validation, rate limiting, and dependency degradation scenarios.

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

Redis: Real Time In Memory Storage

Normal posts use a single integer counter key, whereas viral posts partition counts across 8 sub-keys (:0 through :7) once write velocity exceeds 1,000 likes per minute. Exact user state and retry records are kept on user scoped Redis shards, while a Bloom Filter can act only as a memory efficient prefilter for mega-post membership checks after backfill is complete.

REDIS
# Acting user's exact like state, routed by user_id
like_state:{user_id}:<post_id>:
  type: hash
  fields: [liked, operation_seq, operation_id]
  routing: user_id hash tag keeps read your writes on one Redis shard

# Per-user/post sequence state. Retain it for at least the Kafka replay and recovery horizon.
like_seq:{user_id}:<post_id>:
  type: string (monotonic integer)
  retention: "At least the Kafka retention and recovery horizon so replayed events cannot reuse an old sequence"

# Durable retry gate on the same user hash slot
like_op:{user_id}:<operation_id>:
  type: string
  ttl: 86400s
  value: status|operation_seq|desired_state|effective_delta|request_fingerprint

pending_like:{user_id}:<post_id>:
  type: string
  ttl: 86400s
  value: operation_id|operation_seq|request_fingerprint
  purpose: "Blocks a later state change until the earlier event is durably enqueued"

# Like counter for normal posts (< 10K likes)
like_count:<post_id>:
  type: string (integer counter)
  operations: [INCRBY, GET]

# Like counter for hot posts (8 sub-keys)
like_count:<post_id>:<0..7>:
  type: string (integer counter)
  operations: [INCRBY, GET x 8 in parallel]
  shard_selector: "hash(user_id) % 8"
  aggregation: "Summed across all 8 sub-keys on read"
  cluster_note: "Do not use a shared {post_id} hash tag or all 8 counters map to one Redis Cluster slot"

# Optional negative-membership prefilter for mega-posts, sharded by user
bf:liked:<post_id>:<0..7>:
  type: bloom_filter (RedisBloom)
  operations: [BF.ADD, BF.EXISTS]
  shard_selector: "hash(user_id) % 8"
  error_rate: 0.01
  readiness: "Backfill must complete before a negative result can safely skip the exact state check"
  rule: "A positive result never proves a duplicate, so exact state must confirm"

# Recent liker candidate index for hot posts
recent_likers:<post_id>:<0..7>:
  type: sorted_set
  score: liked_at_epoch_ms
  member: user_id
  shard_selector: "hash(user_id) % 8"
  operations: [ZADD, ZREM, ZREVRANGE]
  retention: "Keep a bounded recent window and expire or trim inactive keys"
  purpose: "Candidate index only, where exact current state confirms the liker"

# Hot post velocity tracking
like_velocity:{post_id}:{minute}:
  type: string (integer counter)
  ttl: 120s
  operations: [INCR]

# Hot post state machine
hot_posts:
  type: set
  values: post_id currently in warming, hot, or cooldown state

hot_state:<post_id>:
  type: string
  values: normal | warming | hot | cooldown
  transitions: normal -> warming -> hot -> cooldown -> normal

Cassandra: Durable Source of Truth

Cassandra stores durable current like state, the durable operation ledger, and append only like event history using bounded partitions. Current state is bucketed by post_id and a deterministic user bucket for point lookups. Recent events are bucketed by hour and event shard for ordered recent reads, and recent liker pagination scans the newest required buckets with exact current state checks. Aggregate counts are stored as reconciled snapshots rather than Cassandra counters so they can be rebuilt from durable state and a matching Kafka partition offset map.

SQL
-- Authoritative current like state, bucketed to bound celebrity partitions
-- user_bucket = hash(user_id) % 1024
CREATE TABLE like_state_by_post_bucket (
    post_id        UUID,
    user_bucket    SMALLINT,
    user_id        UUID,
    liked          BOOLEAN,
    operation_seq  BIGINT,
    operation_id   TEXT,
    updated_at     TIMESTAMP,
    PRIMARY KEY ((post_id, user_bucket), user_id)
);

-- Durable idempotency ledger for retry-safe operation replay.
-- operation_day = UTC day of creation, and operation_bucket = hash(operation_id) % 4096.
-- The bucket count is an illustrative bound for the stated 5B/day workload and should be tuned to the target partition size.
-- request_fingerprint prevents reuse of one operation_id for a different request.
-- status records pending vs committed. expires_at defines the 24-hour guaranteed idempotency horizon.
-- The application checks the current and previous operation-day buckets on retry, then purges expired buckets asynchronously.
CREATE TABLE like_operation_by_day (
    operation_day        DATE,
    operation_bucket     SMALLINT,
    operation_id         TEXT,
    user_id              UUID,
    post_id              UUID,
    action               TEXT,
    request_fingerprint  TEXT,
    operation_seq        BIGINT,
    desired_state        BOOLEAN,
    effective_delta      SMALLINT,
    status               TEXT,
    created_at           TIMESTAMP,
    updated_at           TIMESTAMP,
    expires_at           TIMESTAMP,
    PRIMARY KEY ((operation_day, operation_bucket), operation_id)
);

-- Recent like events used to serve recent liker lists.
-- event_shard = hash(user_id) % 64 keeps hourly partitions bounded for the stated viral burst.
-- Tune the shard count from peak event rate, burst duration, and the target partition size.
CREATE TABLE recent_like_events_by_post (
    post_id        UUID,
    hour_bucket    TIMESTAMP,
    event_shard    SMALLINT,
    liked_at       TIMESTAMP,
    user_id        UUID,
    operation_id   TEXT,
    action         TEXT,
    operation_seq  BIGINT,
    PRIMARY KEY ((post_id, hour_bucket, event_shard), liked_at, user_id, operation_id)
) WITH CLUSTERING ORDER BY (liked_at DESC)
  AND default_time_to_live = 604800;

-- Per-user like event history for account-level queries.
CREATE TABLE like_events_by_user_bucket (
    user_id        UUID,
    month_bucket   DATE,
    event_at       TIMESTAMP,
    post_id        UUID,
    operation_id   TEXT,
    action         TEXT,
    operation_seq  BIGINT,
    PRIMARY KEY ((user_id, month_bucket), event_at, post_id, operation_id)
) WITH CLUSTERING ORDER BY (event_at DESC);

-- Verified aggregate snapshots. These are derived projections and can be rebuilt.
-- A snapshot is valid only with the matching Kafka partition offsets through which its count is complete.
-- The snapshot and offset map form one logical replay barrier for promotion, demotion, and recovery.
CREATE TABLE like_count_snapshots (
    post_id                         UUID PRIMARY KEY,
    like_count                      BIGINT,
    applied_partition_offsets_json  TEXT,
    updated_at                      TIMESTAMP
);

Fault Tolerance

ConcernSolution
Redis node crashUse Redis Cluster across AZs. Rebuild user like-state cache and hot count projections from Cassandra state plus replay from the durable Kafka offset map after failover. Do not treat Redis replicas as the only durable source.
Duplicate like submissionsUse operation_id plus an atomic user scoped state transition in Redis, reject idempotency key reuse with a different request fingerprint, and deduplicate durable events by operation_id and ordered operation_seq. Bloom Filter results are never authoritative.
Count drift between layersRun periodic reconciliation from authoritative Cassandra like state or a verified snapshot plus its Kafka offset map, compare with Redis shard SUMs, and repair drift above 0.1%.
Hot post overloads single nodePromote viral posts to 8 counter shards, use deterministic user based shard placement, and optionally batch effective deltas at the service layer.
Kafka consumer lagServe low latency state and count reads from Redis while consumers catch up. Alert at >30s lag and retain Kafka long enough to replay safely from durable offsets.

Additional Considerations

Count Display Formatting

Client applications format large numbers into abbreviated strings, which naturally masks minor eventual consistency drift between Redis cache shards and persistent storage.

Threshold Range      Display Format    Example Display
------------------------------------------------------
< 1,000              Exact integer     "842"
1,000 - 999,999      Thousands (K)     "15.2K"
1,000,000 - 999M     Millions (M)      "1.5M"
1,000,000,000+       Billions (B)      "1.2B"

Note: Because UI formatting rounds counts to three significant digits,
a drift of ±100 likes is imperceptible to end users.

Celebrity Notification Batching

When a celebrity post attracts 1M likes per minute, individual push alerts would overwhelm client devices and mobile push gateways. The system groups notifications into dynamic time window digests based on engagement velocity.

Celebrity Notification Batching Strategy:

Cumulative Likes     Notification Delivery Policy
--------------------------------------------------------------------------------
1 - 10 likes         Send real time individual push alerts ("Alice liked your post")
11 - 100 likes       Batch every 5 minutes ("42 people liked your post")
101 - 1,000 likes    Batch every 15 minutes ("850 people liked your post")
1,000+ likes         Hourly digest notification ("1.2M people liked your post")

Related Problems

The Live Likes & Reactions system focuses on live-streaming broadcasts that require 500ms batched deltas over WebSockets, client animation sampling, and high throughput reaction ingestion. In contrast, this design targets static post feeds through tiered ingestion paths, distributed counter sharding, and durable Cassandra persistence. While both architectures rely on Redis atomic increments, they diverge significantly in real time broadcast fan out and memory trade-offs for user deduplication.

Interview Walkthrough

Structure the discussion by starting with single row database bottlenecks before introducing tiered write paths and counter sharding.

  • 25-Minute Interview Strategy

    Prioritize tiered write paths and counter sharding before exploring staff-level multi region or Bloom Filter trade-offs.

    • Core Requirements and Bottlenecks of Relational Counters (5 min)
    • Tiered Ingestion Architecture: Normal vs Hot Posts (6 min)
    • Counter Sharding Across Redis Keys (5 min)
    • Deduplication: Exact Redis state vs Bloom Filter vs Cassandra (5 min)
    • Read Your Writes Guarantees and Cache Warming (4 min)
  • Start with the naive UPDATE like_count = like_count + 1 pattern and explain why single row lock contention exhausts the database connection pool at 1M writes per second.
  • Introduce tiered ingestion paths that use the same durable Kafka event path for normal and hot posts, with one counter projection for normal posts and 8 sharded counter projections for hot posts detected by like velocity.
  • Shard hot counters into 8 Redis keys when single key write throughput saturates an individual Redis node.
  • Reference Redis Patterns for Interview Systems to explain atomic increments and distributed caching mechanics.
  • Explain deduplication trade-offs using exact user scoped Redis state for smaller workloads, Bloom Filter prefiltering for mega-posts with 100M+ likes, and Cassandra as the durable source of truth. A Bloom Filter positive is never sufficient to reject a like.
  • Separate read your writes guarantees, which ensure the acting user immediately sees their filled heart icon, from eventual aggregate count display.
  • Describe proactive cache warming when celebrity posts enter follower feeds to prevent thundering herd cache misses.
  • Batch celebrity push notifications into aggregated digests once cumulative likes exceed 1,000 to avoid overwhelming creator devices.
  • Emphasize the common pitfall of relying on transactional relational counters or synchronous fan out for viral social media content.

Engineering Trade-offs

Redis SET vs Bloom Filter vs Cassandra-Only Dedup

Selecting the appropriate deduplication strategy requires balancing memory consumption against accuracy and latency requirements.

ApproachMemory assumptionAccuracyLatency assumptionBest For
Redis SET reference16 MB100% exact< 1 msLegacy or small-post membership sets
Bloom Filter prefilter125 MB per 100M entries~99% membership estimate and not authoritative< 1 msMega-posts as a negative-membership prefilter
Cassandra exact stateDisk100% exact~5 ms (illustrative)Durable fallback and reconciliation
Hybrid ⭐Redis state + BF prefilter for large postsExact state with approximate prefilter< 1 ms cache pathProduction systems at scale

Latency and memory values are illustrative planning assumptions. Actual Redis and Cassandra performance depends on hardware, topology, request shape, replication, and workload.

Why Not PostgreSQL for Likes?

Relational databases fail under viral engagement workloads because row locking on shared counter records creates severe transaction bottlenecks.

Relational Database (PostgreSQL):
  Bottleneck at 1M likes/sec:
    Row-level locking on one shared counter row creates severe contention. The exact sustainable limit is workload-dependent, so ~10K writes/sec is an illustrative capacity point rather than a universal PostgreSQL limit.
  Bottleneck on count aggregations:
    Running COUNT(*) across 100M like records requires scanning or maintaining expensive indexes unless a materialized counter is maintained separately.

Distributed NoSQL (Cassandra):
  Write efficiency:
    Partitioned writes avoid a single shared relational row lock, but the workload must still be modeled around bounded partitions and predictable access paths.
  Partitioning strategy:
    Current like state is bucketed by post_id and a user bucket, so a celebrity post does not create one unbounded Cassandra partition.
  Read efficiency:
    Point lookups use post_id + user_bucket + user_id. Aggregate counts come from Redis projections and durable count snapshots rather than a COUNT(*) scan.

Recommended Hybrid Architecture:
  User-scoped Redis state provides immediate read your writes. Kafka provides the durable event path. Cassandra stores authoritative current state and recent event history, while Redis stores low latency aggregate count projections.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...