Design StreamChat — StreamHub's messaging layer supporting one-to-one and group chats with delivery receipts, online presence, and message history. This is a classic product design: scope first, estimate scale, then draw boxes.
Requirements
Functional requirements
- One-to-one and group messaging (up to 500 members per group).
- Real-time delivery to online users; push notification to offline users.
- Message history: fetch last N messages and paginate backward.
- Delivery and read receipts (sent ✓, delivered ✓✓, read ✓✓ blue).
- Online/offline presence indicators and typing indicators.
- Share images and links (media stored in object storage).
Non-functional requirements
- Latency: message delivery under 200 ms for online users.
- Scale: 50M DAU; 10M concurrent WebSocket connections.
- Availability: 99.9% — chat is core engagement.
- Ordering: Messages in a conversation appear in send order.
- Durability: Messages never lost once server acknowledges send.
Out of scope: Voice/video calls, end-to-end encryption, message editing (mention as future work).
Estimation
- 50M DAU × 40 messages/day = 2B messages/day ≈ 23K messages/sec average, ~70K peak.
- Avg message size: 200 bytes text + metadata ≈ 500 bytes → 2B × 500 B = 1 TB/day storage.
- 5-year retention for premium users: ~1.8 PB — requires tiered storage and archival.
- 10M concurrent WebSocket connections → ~10K connections per server × 1,000 connection servers.
Chat system (AWS)
API design
# REST — history and metadata
GET /v1/conversations/{id}/messages?cursor=xxx&limit=50
POST /v1/conversations # create 1:1 or group
GET /v1/conversations/{id}/members
# WebSocket — real-time channel
WS /v1/chat/connect
→ client sends: { "type": "message.send", "conversation_id": "c_123", "body": "gg wp" }
← server pushes: { "type": "message.new", "message": { "id": "m_456", ... } }
← server pushes: { "type": "receipt.delivered", "message_id": "m_456" }
← server pushes: { "type": "presence.update", "user_id": "u_789", "status": "online" }
CREATE TABLE messages (
id BIGINT PRIMARY KEY,
conversation_id BIGINT NOT NULL,
sender_id BIGINT NOT NULL,
body TEXT,
media_url TEXT,
created_at TIMESTAMPTZ NOT NULL,
status VARCHAR(16) DEFAULT 'sent'
);
CREATE INDEX idx_messages_conv ON messages (conversation_id, created_at DESC);
CREATE TABLE conversations (
id BIGINT PRIMARY KEY,
type VARCHAR(8) NOT NULL, -- direct, group
created_at TIMESTAMPTZ DEFAULT now()
);
CREATE TABLE conversation_members (
conversation_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
joined_at TIMESTAMPTZ DEFAULT now(),
PRIMARY KEY (conversation_id, user_id)
);
High-level architecture
Core services
- Chat service: Manages WebSocket connections; routes messages to connected recipients.
- Connection registry: Redis hash
user_id → server_id— which server holds each user's socket. - Message store: Cassandra partitioned by
conversation_id— write-heavy, time-ordered reads. - Presence service: Redis sorted sets track online users with heartbeat TTL.
- Notification service: Push to offline users via FCM/APNs (see notification-system article).
- Media service: Presigned S3 upload; store URL in message record.
StreamHub production architecture (AWS)
Deep dive: cross-server message routing
When User A (on Server 1) sends to User B (on Server 3):
Routing steps
- Server 1 persists message to Cassandra; returns ack to User A.
- Server 1 looks up User B in Redis:
user:B → server_3. - Server 1 publishes to Redis pub/sub channel
server_3with the message payload. - Server 3 receives pub/sub event; pushes to User B's WebSocket.
- Server 3 sends delivery receipt back through the same path.
def route_message(msg: Message, recipient_id: int):
target_server = redis.get(f"presence:{recipient_id}")
if target_server:
redis.publish(f"chat:server:{target_server}", msg.to_json())
else:
notification_service.send_push(recipient_id, msg)
Deep dive: group message fan-out
Group with 500 members: don't synchronously push to 500 sockets. Persist once, then async fan-out via internal queue — each recipient resolved through connection registry independently.
Quick recall
Everything you need if you only revisit this box.
- WebSocket for online real-time; push notifications for offline users.
- Redis connection registry routes messages across chat server pods via pub/sub.
- Cassandra (partition by conversation_id) for message storage at billions of rows.
- 10M concurrent connections ≈ 1,000 servers at 10K connections each.
- 2B messages/day ≈ 23K msg/sec — partition DB by conversation, async group fan-out.
- Delivery receipts via separate lightweight events; presence via heartbeat + TTL.
Test yourself
Answer these before moving on — recall is what makes it stick.