Interview Setup
Interview Prompt
Design a live reaction system for streaming events. Viewers tap emoji reactions during a live stream, see counts update in real time, and watch floating reaction animations on screen.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Is this for live streams only, or also static posts? | Live streams have a WebSocket broadcast path and much higher peak throughput than static posts. |
| Do we broadcast every individual reaction, or just count updates? | 1M individual emoji events/sec will melt the WebSocket layer. Batched deltas every 500ms is the realistic path. |
| Can a user change their reaction type (heart to laugh)? | Reaction changes need an atomic state transition that can decrement the previous type, increment the new type, reject stale action sequences, and remain safe across retries. |
| How many concurrent viewers on a peak stream? | 10M viewers x 1 reaction each in 10 seconds = 1M reactions/sec. That drives the entire aggregation design. |
Scope
In scope
- High throughput reaction ingestion (1M+/sec)
- Buffered counter aggregation (500ms windows)
- WebSocket broadcast of batched deltas
- Floating emoji animation sync
- Per user dedup (one reaction per content item)
Out of scope (state explicitly)
- Static post likes at celebrity scale (covered in Like Count for High-Profile Posts)
- Full ML ranking model training pipeline
- Direct messaging / chat
- Ad insertion and monetization
Functional Requirements
Scope live reactions with your interviewer by discussing real time broadcasting, per user deduplication, and whether floating emoji animations are required for this round.
In an interview setting, distinguish this real time streaming problem from Like Count for High-Profile Posts, because live viewer broadcast latency takes precedence over audit grade transactional accuracy.
For live streaming environments, the primary architectural focus is real time reaction broadcast with under one second count delivery across millions of concurrent viewers. For static posts and celebrity-scale like counters where audit durability and write behind reconciliation dominate, refer to Like Count for High-Profile Posts.
- React to content: Allow users to emit live reactions (❤️ 😂 😮 😢 😡) to posts, videos, or active live streams.
- Real time count: Display reaction counts that update through 500ms aggregated broadcast windows for all active viewers.
- Reaction animations: Render floating reaction emojis on live streaming screens in real time.
- Toggle reactions: Allow users to change or withdraw their reaction immediately.
- Aggregate totals: Show the cumulative total counts grouped per reaction type (e.g., 1.2K ❤️, 300 😂).
- Deduplication: Restrict every user to exactly one reaction per content item.
Non-Functional Requirements
Under one second broadcast latency and over one million reactions per second at peak represent the primary non functional requirements. The Redis serving count updates on the 500ms aggregation cadence, while Cassandra durable state can lag by a few seconds because persistence is asynchronous.
- Ultra low latency: Broadcast reaction updates to all viewers within 1 second.
- High throughput: Support over 1M reactions per second during peak viral streams.
- Eventual consistency: Cassandra persistence can lag a few seconds, while the live Redis serving count remains responsive during the event.
- High availability: 99.99% availability because reactions drive live engagement during events.
- Idempotency: Repeated taps and retries must not duplicate counts.
Capacity Estimations
Viral live streams cause massive spikes in reaction volume, making it essential to size WebSocket fan out and counter aggregation before constructing the high level architecture diagram.
| Metric | Calculation | Value |
|---|---|---|
| Peak Reactions / sec (viral live event) | Given (viral live event peak assumption) | 1,000,000 |
| Avg Reactions / sec (active traffic) | Given average active traffic rate during live streams. Daily volume is a separate planning assumption | 50,000 |
| Base Reaction Payload Size | Derived calculation basis | 50 bytes (user_id + content_id + type) |
| Write Throughput (Peak) | Given (peak load assumption) | 50 MB/s |
| Active Live Events | Derived | 10,000 |
| Reactions per Event (Avg) | Given (typical workload assumption) | 100,000 |
I/O and Bandwidth Calculations:
IngestionBandwidthPeak: "1,000,000 reactions/sec x 50 bytes = 50 MB/s incoming bandwidth"
DailyAggregateCounts:
ActiveEvents: 10,000
ReactionsPerEvent: 100,000
TotalReactionsPerDay: "1 Billion reactions per day"
StoragePerDay: "1 Billion x 50 bytes = 50 GB storage space needed per day (raw persistence)"Architecture Diagram
Clients submit reaction mutations over HTTP POST, while WebSocket carries the live count deltas back to viewers. The Reaction Service validates each mutation, publishes the accepted event to Kafka, and returns the acknowledgement without waiting for Cassandra or WebSocket broadcast. The hot path therefore avoids relational database latency and avoids a Redis plus Kafka dual write. Shard aggregators consume Kafka, apply atomic Redis state transitions, and form partial deltas in 500ms windows. A merge stage combines those partials into one logical content window and pushes frames to WebSocket gateways through Redis Pub/Sub. Cassandra persists the accepted current reaction state and rebuildable aggregate counters asynchronously every 5 seconds.
In an interview setting, emphasize that viewers need 500ms batched deltas rather than individual reaction pushes, stating this throughput mismatch before drawing component boxes.
Component Deep Dives
Hot Counter Aggregation
Walk the reaction path from client tap to broadcast. Dedup and idempotency belong on the write path before counters increment. A viral stream can generate 1M reactions per second. Writing each one to the database with UPDATE counters SET count = count + 1 locks rows and crashes the database within seconds. Redis HINCRBY or its sharded equivalent handles the serving count in memory, while Cassandra gets batched durable writes every 5 seconds.
# Redis serving counters materialized by the Reaction Aggregator
Write on each Redis shard:
HINCRBY reaction_count:{content_id:<shard>} heart 1
HINCRBY reaction_count:{content_id:<shard>} laugh 1
Read:
HGETALL reaction_count:{content_id:<shard>}
# Returns the shard local counts for each reaction type
Important:
- The aggregator applies add, remove, and type change mutations with one Lua transition.
Reaction state and counter deltas change atomically on the same Redis shard.
- For sharded counters, read and sum all N content shards for the displayed total.
- Redis is the low latency serving layer. Kafka remains the durable accepted event stream.
- A promoted stream keeps its shard count stable for the event lifetime, or uses a versioned
mapping during migration so a user mutation cannot move to a different state shard.-- Layer 2: Cassandra persistent state
Write:
-- Kafka db-writer flushes accepted state every 5 seconds using bounded micro-batches.
-- shard_id = hash(user_id) % N for a hot stream. Use N=1 for a normal stream and N=128
-- for a mega-stream with tens of millions of reactors to keep durable partitions bounded.
INSERT INTO user_reactions (content_id, shard_id, user_id, type, action_id, action_seq, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?);
-- The Kafka key preserves ordering for one user within a content shard. The stream processor
-- tracks the latest action_seq for that user and discards stale or duplicate actions before
-- issuing the Cassandra upsert. This avoids an LWT on every hot-stream write.
UPDATE user_reactions
SET type = ?, action_id = ?, action_seq = ?, updated_at = ?
WHERE content_id = ? AND shard_id = ? AND user_id = ?;
-- Derive aggregate deltas only from accepted current-state transitions rather than trusting
-- an event's previousType field. reaction_counts is a rebuildable materialized snapshot.
-- Store absolute snapshot values rather than non-idempotent counter increments. Guard each
-- snapshot write with a monotonically increasing snapshot_seq so an older batch cannot overwrite
-- a newer snapshot when multiple persistence workers handle the same content_id.
-- If a worker crashes during a batch, replay retained Kafka events or rebuild counts from
-- user_reactions and resume from the last committed consumer offset.
-- user_reactions remains the durable current-state source of truth.Deduplication
Each user has one current reaction per content item. The Redis state transition stores the current type and action_seq, so a type change decrements the old type and increments the new type atomically. Retries and stale actions become no ops. Cassandra stores the same current state under PRIMARY KEY ((content_id, shard_id), user_id), and a stale action_seq cannot overwrite a newer state. The aggregate snapshot is updated with absolute values and snapshot_seq so parallel persistence workers cannot publish an older snapshot over a newer one.
# Atomic reaction state transition in Redis
State key:
reaction_state:{content_id:<shard>}
Field:
user_id
Value:
{ type, action_seq, action_id }
For each accepted Kafka event, one Lua script on that shard:
1. Reads the current reaction type, action_seq, and action_id.
2. If action_id was already applied, returns the prior result without changing counters.
3. Ignores the mutation when action_seq is not newer than the stored value.
4. For a first reaction, increments the new type and total.
5. For a type change, decrements the old type and increments the new type.
6. For removal, decrements the current type and total, then clears the reaction type.
7. Stores action_id, type, and action_seq atomically with the counter delta.
# Durable current state in Cassandra
Schema:
PRIMARY KEY: ((content_id, shard_id), user_id)
Columns: type, action_id, action_seq, updated_at
Behavior: Persist the latest accepted reaction state. The stream processor rejects stale
action_seq values before the Cassandra upsert.Real Time Broadcast
Pushing 1M individual reaction events to 10M WebSocket clients is not feasible. The aggregator collects 500ms of events, sums them into one delta frame per content_id, and publishes it through Redis Pub/Sub. Clients use the frame total for the displayed count and the delta for animation. Pub/Sub is a best effort transport, so a missed frame is recovered by fetching a fresh count snapshot.
# Aggregated Delta Broadcasting for Live Streams
WindowDuration: 500ms
AggregatedPayload:
content_id: "live-456"
window_seq: 1842
deltas:
heart: 342
laugh: 89
wow: 23
total: 1523400
BroadcastMechanism:
- Aggregator consumes accepted Kafka events, applies Redis state, and emits one frame per content_id per window.
- WebSocket gateways push the frame to their local viewers for that stream.
- Clients use total for the displayed count and delta values for animation.
- Clients discard duplicate window_seq values, detect gaps, and refresh a full count snapshot when needed.
- Client renders a random sample of 5 floating emojis rather than all 342 items.Animation Sync
Floating emoji animations are a visual effect, not a data contract. The client applies the full total for the displayed count but renders only 5 random emojis from the delta pool. This keeps the browser responsive even when reaction rates spike, and window_seq prevents stale frames from being applied twice.
# Animation synchronization and client rendering load
Challenge:
- One million concurrent viewers reacting simultaneously creates 1M potential emoji animations/sec
Client Side Sampling:
- Server sends a batched delta: { heart: +342, laugh: +89 } together with window_seq and total.
- Client uses total as the displayed count and delta only for animation.
- Client randomly samples 5 floating emojis from the delta pool to animate smoothly.
Server Side Sampling:
- Aggregator includes a sample array of 10 user avatars for bubble bursts when required by the product.
- Remaining count is rendered as "+332 others".
Timing:
- Aggregator dispatches batched deltas every 500ms.
- Client animation loop interpolates visual effects over the 500ms window.
- Messages carry server timestamps and window_seq values so clients can detect stale or missing frames.Event Bus (Kafka)
Kafka decouples the fast acknowledgement path from persistence and broadcast. Normal streams can partition by content_id, while hot streams shard a content item across partitions using a user based shard key. The aggregator groups those events back by content_id and window, so a viral stream can use multiple Kafka partitions without requiring a total order across different users.
Topic: reactions
Partitions: 128
PartitionKey: normal streams use content_id, while hot streams use content_id + hash(user_id) % N
Normal streams keep all events for a content item ordered on one partition.
Hot streams spread one content item across N partitions while preserving per user ordering.
Hot stream shard count N is fixed for the stream generation. If N changes, use a versioned routing
mapping and controlled migration so a user's mutations do not move between ordering domains.
Retention: 24 hours (high volume streaming telemetry)
ReplicationFactor: 3
MinInSyncReplicas: 2
Producer:
Config: enable.idempotence=true, acks=all
Payload:
eventId: string
actionId: string
actionSeq: int64
userId: string
contentId: string
previousType: "heart" | "laugh" | "wow" | "sad" | "angry" | null
type: "heart" | "laugh" | "wow" | "sad" | "angry" | null
action: "add" | "remove" | "change"
timestamp: int64
ConsumerGroups:
1. shard-aggregator:
Window: 500ms tumbling window
Action: Apply atomic Redis state transitions and form partial deltas per content_id and Kafka shard
Note: For hot streams, user hash sharding keeps one user's mutations ordered while spreading event load.
2. content-window-merger:
Window: 500ms tumbling window
Action: Merge shard partials by content_id and publish one windowed frame to Redis Pub/Sub
Note: The merge stage handles a small number of partials rather than the full 1M events/sec input
Ordering: Emits monotonically increasing window_seq values per content_id so clients can detect gaps.
3. db-writer:
Window: 5 second tumbling flush using bounded micro-batches
Action: Apply only newer action_seq state transitions to Cassandra and refresh the rebuildable aggregate snapshot
Ordering: Serialize snapshot publication per content_id or guard low-frequency snapshot writes with snapshot_seq
4. analytics:
Action: Sink the event stream to ClickHouse for historical analytics
Paths:
Sync: Validate request, publish the accepted event to Kafka, return ACK after Kafka acceptance (< 50ms target)
Async: Shard aggregators apply Redis state and form 500ms partial deltas. A merge stage combines partials into one content_id window before WebSocket push, while db-writer persists current state to Cassandra.
Reliability: Replayed events are ignored by action_id or stale action_seq checks. Kafka producer idempotence prevents duplicate producer records, while consumer state checks handle at least once delivery.
DLQ: Route to reactions-dlq after 3 failures, alert operators, and replay after the underlying issue is fixedAPI Design
Service Interfaces and Domain Types
The API handles reaction submission, withdrawal, aggregate count lookups, and WebSocket delta distribution. Mutation endpoints acknowledge accepted Kafka events. Any returned counts are the latest materialized serving snapshot and may not yet include the accepted mutation.
export type UUID = string;
export type ContentId = UUID;
export type UserId = UUID;
export type ActionId = UUID;
export type ActionSequence = number;
// Monotonic per user and content item.
export type ReactionType = "heart" | "laugh" | "wow" | "sad" | "angry";
export interface ReactionMutationRequest {
contentId: ContentId;
type: ReactionType | null;
actionId: ActionId;
actionSeq: ActionSequence;
}
export interface ReactionRequest {
contentId: ContentId;
type: ReactionType;
actionId: ActionId;
actionSeq: ActionSequence;
}
export interface ReactionCounts {
heart: number;
laugh: number;
wow: number;
sad: number;
angry: number;
total: number;
}
export interface ReactionResponse {
accepted: boolean;
contentId: ContentId;
currentCounts: Partial<Record<ReactionType, number>>;
currentType: ReactionType | null;
acceptedActionSeq: ActionSequence;
countsIncludeAcceptedAction: boolean;
}
export interface WebSocketDeltaFrame {
type: "reaction_update";
contentId: ContentId;
windowSeq: number;
deltas: Partial<Record<ReactionType, number>>;
total: number;
timestamp: number;
}
export interface LiveReactionsService {
emitReaction(userId: UserId, request: ReactionRequest): Promise<ReactionResponse>;
removeReaction(userId: UserId, request: Omit<ReactionMutationRequest, "type">): Promise<ReactionResponse>;
getReactionCounts(contentId: ContentId): Promise<ReactionCounts>;
broadcastDeltas(contentId: ContentId, frame: WebSocketDeltaFrame): Promise<void>;
}1. React to Content
POST /api/v1/reactions
Content-Type: application/json
Idempotency-Key: action-uuid-123
{
"content_id": "post-123",
"type": "heart",
"action_id": "action-uuid-123",
"action_seq": 42
}
HTTP/1.1 202 Accepted
Content-Type: application/json
{
"accepted": true,
"current_counts": { "heart": 15235, "laugh": 3421 },
"current_type": "heart",
"accepted_action_seq": 42,
"counts_include_accepted_action": false
}
# current_counts is the latest materialized serving snapshot and may exclude this accepted action.2. Remove Reaction
DELETE /api/v1/reactions?content_id=post-123&action_id=action-uuid-124&action_seq=43
Idempotency-Key: action-uuid-124
HTTP/1.1 202 Accepted
Content-Type: application/json
{
"accepted": true,
"current_counts": { "heart": 15234, "laugh": 3421 },
"current_type": null,
"accepted_action_seq": 43,
"counts_include_accepted_action": false
}3. Retrieve Reaction Counts
GET /api/v1/reactions/counts?content_id=post-123
HTTP/1.1 200 OK
Content-Type: application/json
{
"heart": 15234,
"laugh": 3421,
"wow": 892,
"total": 19547
}4. WebSocket Live Updates
// Server pushes an aggregated frame over WebSocket every 500ms:
{
"type": "reaction_update",
"content_id": "live-456",
"window_seq": 1842,
"deltas": { "heart": 342, "laugh": 89 },
"total": 1523400,
"timestamp": 1760000000123
}Common Error Responses
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 backoffData Model
Redis Structure
# Sharded serving counters
Key: reaction_count:{content_id:<shard>}
Type: Hash
Fields: { "heart": 15234, "laugh": 3421, "angry": 88 }
# Current per user reaction state
Key: reaction_state:{content_id:<shard>}
Type: Hash
Field: user_id
Value: { type, action_seq, action_id }
# Optional advisory membership gate
Key: reacted:{content_id:<shard>}
Type: Set or Bloom Filter
TTL: 24 hours for live streamsCassandra Schema Design
-- Durable current reaction state
CREATE TABLE user_reactions (
content_id UUID,
shard_id SMALLINT,
user_id UUID,
type TEXT, -- heart, laugh, wow, sad, angry, or NULL when removed
action_id UUID,
action_seq BIGINT,
updated_at TIMESTAMP,
PRIMARY KEY ((content_id, shard_id), user_id)
);
-- Persisted aggregate reaction counts. This table is rebuildable from user_reactions.
CREATE TABLE reaction_counts (
content_id UUID PRIMARY KEY,
heart BIGINT NOT NULL,
laugh BIGINT NOT NULL,
wow BIGINT NOT NULL,
sad BIGINT NOT NULL,
angry BIGINT NOT NULL,
total BIGINT NOT NULL,
snapshot_seq BIGINT NOT NULL
-- Monotonic snapshot version. Older snapshots must not overwrite newer ones.
);Fault Tolerance
| Concern Scenario | Mitigation Strategy |
|---|---|
| Redis Cluster Node Outage | Redis failover restores serving state when replicas are healthy. If serving memory is lost, rebuild counters and current reaction state from Cassandra and replay retained Kafka events for the missing interval. |
| Double Counting / Retries | Atomic Redis Lua transitions use action_id and action_seq, while Cassandra ignores stale action_seq values. Aggregate counters remain rebuildable from the durable current state. |
| Kafka Consumer Backlog | The live Redis serving path continues while persistence catches up. Operators alert on lag and scale the db-writer consumer group without blocking reaction acknowledgements. Acknowledged events remain retained in Kafka for replay. |
| Redis Pub/Sub or WebSocket Gap | Clients track window_seq, ignore duplicate frames, detect gaps, and fetch a fresh count snapshot. Pub/Sub is treated as a transport for live deltas rather than durable state. |
1. Toggle Reaction Race Conditions ⭐
If a user rapidly taps Like, Unlike, and Like again, out of order packet delivery could corrupt the count. We use action_id and a monotonic action_seq to preserve user intent instead of relying on packet arrival order. The service treats the accepted sequence at the user and content serialization boundary as authoritative, including across multiple devices. The same sequence also makes retries safe:
# Rapid client action sequence: Like, Unlike, Like within 100ms Accepted chronological intent: - T=0ms: LIKE has action_seq=1 and becomes the current reaction. - T=30ms: UNLIKE has action_seq=2 and removes the current reaction. - T=60ms: LIKE has action_seq=3 and becomes the current reaction again. Result: Net heart count is +1. Network out of order delivery: - Each mutation carries a unique action_id and a monotonic action_seq. - If action_seq=3 arrives before action_seq=1 or action_seq=2, Redis applies 3 and ignores the older actions. - A retry with the same action_id or a stale action_seq becomes a no-op. - Redis set membership is not sufficient for ordering user intent. The action sequence is the authoritative ordering signal. Multiple devices: - The service serializes mutations at the user/content boundary and uses the latest accepted action_seq. - Ties are resolved at the server serialization point, and the accepted sequence is persisted with the current reaction state.
2. Viral Hot Keys (1M+ events on a single stream) ⭐
When 1M users react to a single stream, a single content counter becomes a Redis hot shard. We reduce command volume with batching and distribute the remaining work across sharded sub-counters:
# Scaling hot keys under extreme load (1M reactions/sec on a single content_id)
Solutions:
1. In-Process Aggregation (Reaction Service):
- Buffer reactions in service memory for 100ms.
- Send one batched counter delta with cumulative counts (for example, +342).
- Reduces 1M individual Redis commands to ~200 batched Redis operations/sec across 20 service nodes.
- Batching alone still targets the same Redis shard, so combine it with counter sharding when one content_id is hot.
- The same hot content can also be promoted to a multi partition Kafka key so event ingestion does not serialize on one partition.
2. Sharded Sub-Counters:
- Split the single content counter into N shard keys:
reaction_count:{cid:0}, reaction_count:{cid:1}, ..., reaction_count:{cid:7}
- Choose the shard with hash(user_id) % 8.
- Keep matching reaction state and counter keys for a shard on the same Redis Cluster slot.
- Query: SUM across all 8 sub-counters.
3. Local Redis Sidecars with Write Behind:
- Accept writes into local host Redis instances.
- Background tasks sync accumulated deltas to the central cluster.
- Treat this as an optional optimization because a host failure can lose unflushed state unless another durable path has already acknowledged the mutation.3. Redis Deduplication Set Memory Exhaustion ⭐
Storing 50M user IDs creates large memory pressure, and Redis set overhead makes the raw 800MB UUID payload estimate optimistic. We solve the memory bottleneck with an advisory Bloom filter and exact verification. Durable Cassandra state uses bounded content shards so a mega stream does not create one oversized partition:
For a mega stream, the durable current-state table uses a larger shard count such as N=128, while Redis serving counters may use a smaller operational shard count such as 8. The two shard counts serve different bottlenecks.
# Deduplication memory optimization for 50M unique reactors
Challenge:
- Storing 50M UUID payloads alone at 16 bytes each requires ~800 MB per stream.
- Actual Redis set memory is higher because of object and hash-table overhead.
- Bloom filtering reduces the auxiliary membership footprint, but it does not replace the exact
current reaction state needed to process type changes and removals. That exact state is a
separate memory consideration for mega streams.
Solutions:
1. Bloom Filters:
- Allocate approximately 15 bits per element (~93 MB total).
- This is more conservative than a 1% false positive target when configured with an appropriate number of hash functions.
- Treat the Bloom filter as an advisory front gate. A negative result can skip an exact lookup,
while a positive result triggers an exact current-state check before rejecting a mutation.
- For mega streams, keep exact Redis reaction state for the active-user working set and use
Cassandra for cold-user verification when that state is evicted. If the durable state is still
catching up, defer or retry rather than deciding from the Bloom result alone.
- Never use a Bloom positive as proof that a user has already reacted because false positives are possible.
2. TTL Cleanup:
- Apply 24 hour expiration to live stream deduplication keys after the event ends.
3. Split-Tier Verification:
- Maintain exact in-memory reaction state only for active live users when memory permits.
- Query Cassandra for durable current reaction state when Redis state is unavailable.
- Keep the serving state and Bloom membership key separately so reducing Bloom memory does not
accidentally remove the exact state required for an active user's type change.Additional Considerations
Count Display Formatting & User Perception
Since numbers above 1K are shortened, such as 15.2K or 1.5M, small differences between the Redis serving count and the Cassandra durable count are usually invisible to users. This lets the system prioritize fast in memory aggregation and asynchronous persistence over a heavy synchronous transaction on every reaction.
Abuse Protection and Rate Limiting
Each authenticated user, device, and IP should have a token bucket or equivalent rate limit for reaction mutations so one actor cannot dominate a stream. The gateway rejects excessive requests before they enter Kafka, while the Reaction Service still enforces per user and content state transitions. Rate limits should be higher for ordinary human interaction than for automated traffic, and suspicious bursts can trigger CAPTCHA or temporary backoff without affecting other viewers.
Related Problems
The Like Count for High-Profile Posts architecture focuses on static posts and celebrity-scale counters, emphasizing tiered synchronous and asynchronous write paths, sharded counter keys, and read your own writes consistency. In contrast, this design centers on the live streaming path with WebSocket batched deltas and 500 millisecond tumbling window aggregations. The two designs address fundamentally different concurrency and delivery constraints and should remain distinct during an interview.
Interview Walkthrough
- 25 minute cut
Skip arch50 and arch75 depth unless interviewing for a staff-level role.
- Frame the throughput mismatch with 1M reactions per second (5 min)
- Propose tiered counters with Redis for real time display and Cassandra for persistence (6 min)
- Enforce single reaction deduplication via Redis sets and Bloom filters (5 min)
- Design 500ms batched delta broadcasting over WebSockets (5 min)
- Shard hot counter keys and handle toggle race conditions (4 min)
- State the throughput mismatch: 1M reactions per second arrive during a viral stream, but viewers only need batched counter updates every 500ms.
- Propose tiered counters using patterns from Redis Patterns for System Design for real time display, paired with a Kafka to Cassandra pipeline for durable persistence with 5 second batching.
- Enforce one current reaction per user using atomic Redis state transitions, switching to an advisory Bloom filter with exact verification when deduplication sets exceed memory on mega streams.
- Make the write path safe across crashes by assigning each mutation an action_id and action_seq, publishing the accepted event to Kafka before acknowledgement, and using an atomic Redis transition that is safe to replay after consumer restarts.
- Broadcast delta frames (
{heart: +342, laugh: +89}) over WebSocket connections using protocols detailed in Network Protocols: HTTP, gRPC, WebSocket & DNS, rather than emitting individual reaction events to every viewer. - Apply partitioning techniques from Sharding & Partitioning to split hot counter keys across N sub-counters when one Redis shard becomes the bottleneck for a viral content ID.
- Accept eventual consistency, because formatted counts such as 15.2K naturally absorb small discrepancies between the Redis serving count and Cassandra persistent state while Kafka retains the accepted event stream for replay.
- Handle toggle race conditions with unique action IDs and monotonic action sequences. The atomic state transition applies only the newest action, so out of order Like and Unlike packets cannot revert a newer user intent.
- Highlight the common anti-pattern of attempting to fan out every individual reaction emoji to millions of WebSocket clients, which causes the broadcast gateway layer to collapse.
Engineering Trade-offs
1. Write Path Tiering (Speed vs Durability)
Exact vs approximate counts, push vs poll architectures, and regional fan out strategies represent the key trade offs in high concurrency reaction systems. Writing every reaction synchronously to durable storage at 1M/sec would add unnecessary persistence pressure to the live path. We trade synchronous database durability for highly responsive in memory writes. If a collector or persistence worker crashes, broadcast freshness may lag and Cassandra may trail by up to the planned 5 second batch window, but acknowledged Kafka events remain available for replay so durable counts can catch up.
2. Counter Implementation Architectures
| Approach | Latency | Accuracy | Durability | Best For |
|---|---|---|---|---|
| Redis HINCRBY ⭐ | < 1 ms | Exact in Redis memory | Volatile (syncs periodically) | Real time viewer display updates |
| Cassandra aggregate snapshot | ~5 ms | Exact after reconciliation | Highly durable | Persistent aggregate serving snapshot |
| Flink window aggregation | 1 to 5 sec window | Exact within window when inputs are deduplicated | Highly durable output | Complex analytics by region and time |
| HyperLogLog (HLL) | < 1 ms | ~0.81% error rate | Volatile | Unique reactor count (approximate) |
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.