Design a Real-Time Chat & Instant Messaging System
1. Problem Statement & Scope Clarification
System Mission
Design a globally distributed, highly scalable, real-time instant messaging platform (similar to WhatsApp, Slack, and Discord) capable of supporting ultra-low latency 1-on-1 direct messaging, group conversations (up to 1,000 members), real-time presence tracking (Online/Offline/Last Seen), multi-device synchronization, read receipts, and durable offline catch-up delivery for hundreds of millions of users worldwide.
Functional Requirements
- 1-on-1 Direct Messaging (
SEND_MESSAGE): Bi-directional real-time text and media messaging with instant delivery confirmations (). - Group Chats & Fan-Out: Scalable group channels supporting up to 1,000 participants with persistent membership, unread badge tracking, and read-state management.
- Presence & Heartbeat Engine: Real-time user online/offline status detection with debounce mechanisms to suppress connection flapping on mobile cellular transitions.
- Message Delivery Receipts: Three-stage message lifecycle states:
SENT(Server Ack),DELIVERED(Device Ack), andREAD(User Seen). - Offline Message Synchronization: Durable message storage allowing offline devices to fetch missed messages using sequential sync cursors upon reconnect.
- Push Notifications: Fallback push notification dispatch (Apple APNs / Google FCM via Amazon SNS) when a recipient is disconnected.
Non-Functional Requirements (SLAs & SLOs)
- High Availability: uptime SLA globally across multi-AZ and multi-region AWS deployments.
- Latency:
- 1-on-1 Message Delivery: , .
- Presence State Broadcast: .
- Data Durability: Zero message loss () for all acknowledged messages.
- Connection Scale: Support concurrent active persistent WebSocket connections.
- Consistency: Monotonic sequence order per conversation (-ordered Snowflake sequences); eventual consistency for global presence broadcasts.
2. Capacity & Scale Estimation (Back-of-the-Envelope Math)
User & Traffic Profile
- Daily Active Users (DAU): ().
- Peak Concurrent WebSocket Connections: ().
- Daily Messages Sent: (5 Billion messages/day).
- Average Message Size: (Text content + metadata headers).
Throughput & QPS Derivations
- Average Message Ingest QPS:
- Peak Message Ingest QPS ( multiplier):
- Group Fan-Out Amplification: Assuming of messages are sent to groups with an average size of 8 active online members:
Network Bandwidth
- Ingress Bandwidth (Writes):
- Egress Bandwidth (Delivery to WebSocket Clients):
Memory & Connection Fleet Sizing
- RAM per Idle WebSocket Connection: (TCP socket buffer + SSL session context in kernel/user space).
- Total In-Memory Socket State:
- WebSocket Gateway Cluster Provisioning:
Deploying ECS tasks sized like a
c6g.2xlarge(; on Fargate this is a task size, on EC2 launch type it is the instance type), capped at persistent sockets per task so that socket state stays near , file-descriptor limits are never approached, and a single task failure drops at most of the fleet's connections: Provisioned across 3 AWS Availability Zones ().
Storage Footprint (5-Year Horizon)
- Daily Ingestion Storage:
- 5-Year Storage Footprint (Single Copy):
- DynamoDB Throughput Provisioning:
- Peak Write: .
- Peak Read (Sync + History): .
- DynamoDB partitions older messages () off to Amazon S3 Standard-IA / Glacier via DynamoDB Streams and Kinesis Data Firehose, reducing live transactional table size by .
3. High-Level Architecture & AWS Component Mapping
Synthesizing vector architecture diagram...
Follow Alice's message to Bob. Both phones keep a WebSocket open, entering through Global Accelerator and the NLB to one gateway node. When a client connects, its gateway records it in the "Presence Registry & Pub/Sub Backplane" panel (Redis for presence, a DynamoDB registry with a 120 s TTL that heartbeats keep alive). Alice's gateway writes the message to Kinesis, and in the "Chat Ingestion & Pipeline" panel workers persist it to DynamoDB, then deliver it: if Bob is online, through Redis Pub/Sub to whichever gateway holds his socket; if he is offline, through SNS as a push notification. Writing to the stream before delivering means an acknowledged message is never lost, even if a gateway crashes.
Data Flow Walkthrough
- Connection Establishment:
- Clients connect via AWS Global Accelerator to the nearest edge PoP, traversing AWS dedicated fiber to an internal Network Load Balancer (NLB).
- The NLB forwards TCP handshakes to the WebSocket Gateway Fleet (ECS Fargate running Rust/Netty).
- Upon authentication, the gateway node records the connection mapping (
user_id -> connection_id, node_ip) in ElastiCache Redis and DynamoDB Connection Registry (with a 120-second TTL), subscribing the gateway node to the user's private Redis Pub/Sub channel.
- Message Ingestion & Sequencing:
- Alice dispatches
SEND_MESSAGEover her established WebSocket. The receiving gateway node generates a globally unique 64-bit Snowflake ID (server_msg_id), claims the next per-conversation sequence number with one atomic DynamoDBUpdateItem ... ADD seq :1on theCONV#<id>/METAitem (), emits aMESSAGE_SENT_ACKback to Alice carrying both, and writes the event to Amazon Kinesis Data Streams. - Why two identifiers? Snowflake IDs are unique and roughly time-ordered across the whole platform, but two gateways can never produce a gap-free sequence for one conversation without coordination. The per-conversation counter is gap-free (), which is what clients need for gap detection (Section 8.3) and offline catch-up; the Snowflake ID is what the platform needs for global uniqueness, idempotent persistence, and cross-conversation ordering (search, moderation, analytics).
- Decoupled worker fleets on Amazon EKS consume from Kinesis, persisting the record to the
ChatPlatformTablein Amazon DynamoDB.
- Alice dispatches
- Presence Evaluation & Real-Time Routing:
- The worker inspects Bob's presence key in ElastiCache Redis (
presence:usr_bob). - Case A (Bob Online): Redis returns Bob's active gateway node (
WS_NODE_2). The worker executesPUBLISH channel:WS_NODE_2 payload, andWS_NODE_2pushes the message to Bob's socket in . - Case B (Bob Offline): If no active presence exists, the worker enqueues a push payload to Amazon SNS, which relays to Apple APNs or Google FCM.
- The worker inspects Bob's presence key in ElastiCache Redis (
- Offline Catch-Up Sync:
- When Bob reconnects, his client presents its last confirmed sequence cursor (
last_sync_seq). The gateway executes a range query against DynamoDB (SK > MSG#<last_sync_seq>) to stream missed messages in bulk, in order, with no gaps.
- When Bob reconnects, his client presents its last confirmed sequence cursor (
Concrete Step-by-Step Request Walkthrough: Tracing Direct Message & Delivery State
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| 1 | Alice sends direct message to Bob | Persistent WebSocket active on WS_NODE_1 | WS_NODE_1 parses WebSocket frame; validates auth | Server acknowledges receipt with MESSAGE_SENT_ACK () |
| 2 | Server assigns identifiers | Snowflake generator in gateway process + DynamoDB atomic counter on CONV#conv_991823/META | server_msg_id = 718293847561029384 (global); ADD seq :1 returns sequence_number = 10482 (per-conversation, gap-free) | Message assigned a globally unique ID and a strictly contiguous conversation sequence |
| 3 | Ingestion into Kinesis log | Kinesis producer client | PutRecord executed with partition_key = conv_991823 | Log appended; offset checkpointed |
| 4 | Redis presence check | Redis Hash lookup | HGET presence:usr_bob gateway_id | Bob identified as connected to WS_NODE_2 |
| 5 | Redis Pub/Sub cluster routing | In-memory message routing bus | PUBLISH channel:WS_NODE_2 {conn_id: "conn_bob", ...} | Backplane routes payload across internal VPC () |
| 6 | WS_NODE_2 delivers to Bob | Local socket write buffer | Binary WebSocket frame dispatched over TCP socket | Bob's phone receives message () |
| 7 | Bob emits delivery receipt | Client SDK automatically responds | DELIVERY_ACK sent to WS_NODE_2 | State updated to DELIVERED; Alice's client updated |
| 8 | Durable persistence | EKS Chat Worker consuming Kinesis | PutItem to DynamoDB ChatPlatformTable | Message safely committed to multi-AZ storage |
4. API Interface Design & Wire Protocols
1. Client-to-Gateway WebSocket Protocol (JSON Framing)
json// 1. Client Sending a Message { "action": "SEND_MESSAGE", "client_msg_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d", "conversation_id": "conv_991823", "recipient_id": "usr_bob_402", "content_type": "TEXT", "body": "Hey Bob, let's review the architectural blueprint.", "timestamp_client_ms": 1767225600120 } // 2. Server Acknowledgment to Sender (Status: SENT) { "event": "MESSAGE_SENT_ACK", "client_msg_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d", "server_msg_id": "718293847561029384", "conversation_id": "conv_991823", "sequence_number": 10482, "timestamp_server_ms": 1767225600135 } // 3. Server Delivery to Receiver Device { "event": "RECEIVE_MESSAGE", "server_msg_id": "718293847561029384", "conversation_id": "conv_991823", "sender_id": "usr_alice_101", "content_type": "TEXT", "body": "Hey Bob, let's review the architectural blueprint.", "sequence_number": 10482, "timestamp_server_ms": 1767225600135 } // 4. Client Read Receipt Update { "action": "READ_RECEIPT", "conversation_id": "conv_991823", "last_read_sequence_number": 10482, "timestamp_read_ms": 1767225605000 }
2. REST Management Endpoints
Create / Initialize Conversation
httpPOST /v1/conversations Host: api.chat.aws.internal Authorization: Bearer <jwt_access_token> Content-Type: application/json { "type": "DIRECT", "participants": ["usr_bob_402"] } Response: 201 Created { "conversation_id": "conv_991823", "type": "DIRECT", "created_at": 1767225500000, "participants": ["usr_alice_101", "usr_bob_402"] }
Fetch Historical Messages (Cursor-Based Pagination)
httpGET /v1/conversations/conv_991823/messages?limit=50&before_seq=10482 Host: api.chat.aws.internal Authorization: Bearer <jwt_access_token> Response: 200 OK { "messages": [ { "server_msg_id": "718293847561029384", "sender_id": "usr_alice_101", "body": "Hey Bob, let's review the architectural blueprint.", "sequence_number": 10482, "status": "READ", "timestamp_server_ms": 1767225600135 } ], "has_more": false, "next_cursor": null }
Pre-Signed URL for Media Uploads
httpPOST /v1/media/upload-url Host: api.chat.aws.internal Authorization: Bearer <jwt_access_token> Content-Type: application/json { "conversation_id": "conv_991823", "content_type": "image/jpeg", "byte_size": 2048576 } Response: 200 OK { "upload_url": "https://chat-media-bucket.s3.amazonaws.com/uploads/conv_991823/img_01.jpg?AWSAccessKeyId=...", "cdn_url": "https://cdn.chat.internal/media/conv_991823/img_01.jpg", "expires_in_seconds": 900 }
5. Data Models & DynamoDB Single-Table Schema
To achieve predictable sub-millisecond query performance at PB scale, we utilize a single-table DynamoDB design partitioned by conversation_id and sorted by the gap-free per-conversation sequence number (zero-padded so lexical order equals numeric order). The Snowflake server_msg_id is stored as an attribute and is the idempotency key for the persistence worker.
DynamoDB Table: ChatPlatformTable
- Partition Key (
PK): Entity identifier. - Sort Key (
SK): Specific entity instance or sequence cursor. - Global Secondary Index 1 (
GSI1): Inbox index (GSI1-PK: USER#<user_id>,GSI1-SK: UPDATED#<timestamp>) for listing user chats.
PK (Partition Key) | SK (Sort Key) | GSI1-PK | GSI1-SK | Attributes & Payloads | Access Pattern / Query |
|---|---|---|---|---|---|
USER#<user_id> | CONN#<connection_id> | - | - | gateway_ip, connected_at, ttl_epoch | Active WebSocket socket lookup |
USER#<user_id> | CONV#<conversation_id> | USER#<user_id> | UPDATED#<ts> | unread_count, last_read_seq, joined_at | Fetch user's inbox sorted by activity |
CONV#<conversation_id> | META | - | - | type (DIRECT/GROUP), members_count, name, seq (atomic per-conversation counter) | Conversation metadata; ADD seq :1 allocates the next sequence number |
CONV#<conversation_id> | MEMBER#<user_id> | - | - | role (ADMIN/MEMBER), joined_at | Fetch all members of a group chat |
CONV#<conversation_id> | MSG#<seq:012d> (e.g. MSG#000000010482) | - | - | server_msg_id (Snowflake), sender_id, body, type, status | Fetch paginated chat history / offline catch-up via range query on sequence |
Range Query for Message History & Offline Catch-Up
json{ "TableName": "ChatPlatformTable", "KeyConditionExpression": "PK = :conv_pk AND SK > :last_sync_sk", "ExpressionAttributeValues": { ":conv_pk": {"S": "CONV#conv_991823"}, ":last_sync_sk": {"S": "MSG#000000010481"} }, "Limit": 50, "ScanIndexForward": true }
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~43%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.