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)
| Question | Why 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.
| Metric | Calculation | Value |
|---|---|---|
| Total likes / day | Given | 5B |
| Likes / sec (avg) | 5B ÷ 86400 | ~58K |
| Peak likes / sec (viral post) | Given viral post peak design target | 1M+ |
| Avg post likes | Given (typical workload assumption) | 50 |
| Celebrity post likes | Given (assumption documented in value) | 10M - 100M |
| Like record size | Given (assumption documented in value) | 32 bytes (user_id + post_id + timestamp) |
| Storage / day | 5B x 32B logical event payload | 160 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.
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.-- 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.
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.
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.
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.
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.
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.
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.
# 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 -> normalCassandra: 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.
-- 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
| Concern | Solution |
|---|---|
| Redis node crash | Use 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 submissions | Use 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 layers | Run 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 node | Promote viral posts to 8 counter shards, use deterministic user based shard placement, and optionally batch effective deltas at the service layer. |
| Kafka consumer lag | Serve 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 + 1pattern 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.
| Approach | Memory assumption | Accuracy | Latency assumption | Best For |
|---|---|---|---|---|
| Redis SET reference | 16 MB | 100% exact | < 1 ms | Legacy or small-post membership sets |
| Bloom Filter prefilter | 125 MB per 100M entries | ~99% membership estimate and not authoritative | < 1 ms | Mega-posts as a negative-membership prefilter |
| Cassandra exact state | Disk | 100% exact | ~5 ms (illustrative) | Durable fallback and reconciliation |
| Hybrid ⭐ | Redis state + BF prefilter for large posts | Exact state with approximate prefilter | < 1 ms cache path | Production 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
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.