Design a Distributed Message Queue
1. Problem Statement & Scope Clarification
System Mission
Design a highly scalable, fault-tolerant, multi-tenant distributed message queue service (combining semantics of Amazon SQS and Apache Kafka/Amazon MSK) supporting asynchronous decoupling, durable at-least-once and strictly FIFO delivery semantics, configurable retention, dynamic visibility timeouts, and automated Dead-Letter Queue (DLQ) redrive capabilities.
Functional Requirements
- Producer Ingestion (
SendMessage/SendMessageBatch): Allow upstream microservices to produce single or batched messages with custom delay and deduplication keys. - Consumer Ingestion (
ReceiveMessage): Allow downstream consumers to poll messages with configurable visibility timeouts and long polling (up to 20 seconds). - Acknowledgment & Deletion (
DeleteMessage/DeleteMessageBatch): Explicitly acknowledge successful processing via cryptographic receipt handles. - FIFO & Partition Ordering: Guarantee strictly in-order delivery and exactly-once ingestion (5-minute content-based or token-based deduplication) per
MessageGroupId, while scaling horizontally across millions of groups. Processing remains at-least-once: a consumer that crashes after side effects but beforeDeleteMessagewill see the message again (see the idempotency invariant in Section 3). - Visibility Timeout Extension (
ChangeMessageVisibility): Allow long-running worker tasks to heartbeat and extend execution leases to prevent double processing. - Dead-Letter Queue (DLQ) & Redrive: Automatically divert poison pills exceeding max retry limits (
maxReceiveCount) to an isolated DLQ without halting partition progress.
Non-Functional Requirements (SLAs & SLOs)
- High Availability: uptime across multi-AZ AWS deployments; zero single point of failure.
- Durability: Zero acknowledged message loss () achieved via synchronous multi-AZ replica quorum.
- Throughput & Scalability:
- Standard Queues: Horizontally unlimited ().
- FIFO Queues: per API action baseline, with 10-message batching, and up to in high-throughput mode (per-
MessageGroupIdpartitioning; exact ceiling is region-dependent).
- End-to-End Latency: Publish-to-receive P99 Latency .
2. Capacity & Scale Estimation (Back-of-the-Envelope Math)
Ingestion & Egress Throughput
- Daily Ingestion Volume: (2 Billion messages/day).
- Average Message Size: (Max payload: ).
- Average Ingest QPS:
- Peak Ingest QPS ( peak-to-average ratio):
- Ingest Bandwidth:
- Fan-Out Egress ( average consumer groups):
Storage Footprint & Replication Capacity
- Raw Daily Storage Ingested:
- 7-Day Retention Storage:
- 3-Way Multi-AZ Storage Replication + Index Overhead:
3. High-Level Architecture & AWS Component Mapping
Synthesizing vector architecture diagram...
Follow a message from the top. Producers connect over mutual TLS through an NLB to brokers in three AZs. In the "Consensus & Metadata Coordinator" panel, brokers use Raft or ZooKeeper to agree on which broker leads each partition, so every write has exactly one owner. In the "Durable Storage Engine" panel, each broker appends messages to its own disk log before acknowledging, and moves sealed segments older than 24 hours to S3, so cheap object storage holds long retention without large disks. Consumers long-poll the brokers for new messages. In the "Dead-Letter Queue Subsystem" panel, a message that fails more than 3 times is moved to the DLQ, where an inspector worker examines it instead of letting it block the queue.
Data Flow Walkthrough
- Producer Ingestion & Raft Quorum Commit: Producer dispatches
SendMessagewith payload and optional deduplication key. Broker Leader maps message to partition viaHash(GroupId) % N, checks in-memory sliding deduplication cache, appends to local io2 EBS Write-Ahead Log (WAL), and replicates to followers via Raft, committing once majority quorum () acknowledges. - Consumer Long Polling & Visibility Lease: Consumer issues
ReceiveMessagewithWaitTimeSeconds=20. Broker parks the socket without CPU spin. Upon message arrival, broker sets an ephemeralVisibilityDeadline = Now + 30s, generates a cryptographicReceiptHandleembedding a monotonicLeaseEpoch, and hides the message from concurrent consumers. - Acknowledgment & Poison Pill Mitigation: The worker finishes processing and calls
DeleteMessage(ReceiptHandle). If the worker crashes or times out, the lease expires and the message becomes visible again. Ifreceive_count > maxReceiveCount(e.g. 3), the broker atomically routes the message to the Dead-Letter Queue (DLQ), unblocking the partition.
At-Least-Once Delivery & Consumer Idempotency Invariant: Because network partitions or worker crashes during acknowledgment can cause message re-delivery after visibility timeout expiry, distributed queues guarantee at-least-once delivery. Exactly-once processing semantics require downstream consumers to implement idempotent processing using database unique constraints on (producer_id, deduplication_id) or transactional outbox tables.
Concrete Step-by-Step Request Walkthrough: Tracing Message Ingestion, Lease, and ACK
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| 1 | Producer issues SendMessageGroupId="grp_1", Token="tok_99" | NLB routes to Broker 1 (AZ-1 Leader); broker checks 5-min sliding deduplication cache | Leader computes partition index Hash("grp_1") % N;initiates append to local io2 EBS WAL | Message buffered in memory; Raft replication initiated |
| 2 | Broker Leader coordinates Raft consensus replication | Followers in AZ-2 and AZ-3 receive AppendEntries RPC | Followers write sequential entry to local WAL; Follower 1 ACKs; Follower 2 ACKs | 2 of 3 quorum ACKs received; Leader advances High Watermark (HW) |
| 3 | Leader emits response to producer | Monotonic sequence number committed; message enters active in-memory ring buffer | Leader updates cluster commit index; returns message ID and sequence to producer | 200 OK(P99 latency: ) |
| 4 | Consumer worker issues ReceiveMessageWaitTime=20s, Visibility=30s | Broker matches message in ring buffer; generates ephemeral ReceiptHandle | Broker sets VisibilityDeadline = Now + 30s;increments receive_count = 1 | Message payload + ReceiptHandledispatched to consumer |
| 5 | Worker completes business logic; invokes DeleteMessage(ReceiptHandle) | Broker validates embedded cryptographic signature and checks LeaseEpoch == CurrentEpoch | Broker marks message permanently deleted; offset reclaimed without consumer re-delivery | DeleteMessageResponse(success=true)emitted to worker |
4. API Interface Design & Wire Protocol
Core HTTP/gRPC API Specifications
protobufsyntax = "proto3"; package hispeeddesign.queue.v1; service MessageQueueService { rpc SendMessage (SendMessageRequest) returns (SendMessageResponse); rpc SendMessageBatch (SendMessageBatchRequest) returns (SendMessageBatchResponse); rpc ReceiveMessage (ReceiveMessageRequest) returns (ReceiveMessageResponse); rpc DeleteMessage (DeleteMessageRequest) returns (DeleteMessageResponse); rpc ChangeMessageVisibility (ChangeMessageVisibilityRequest) returns (ChangeMessageVisibilityResponse); } message SendMessageRequest { string queue_name = 1; string message_body = 2; int32 delay_seconds = 3; string message_group_id = 4; // Mandatory for FIFO queues string message_deduplication_id = 5; // SHA-256 or custom deduplication token map<string, string> message_attributes = 6; } message SendMessageResponse { string message_id = 1; string sequence_number = 2; // Strict monotonic 64-bit sequence in FIFO string md5_of_body = 3; } message ReceiveMessageRequest { string queue_name = 1; int32 max_number_of_messages = 2; // 1 to 10 int32 visibility_timeout_seconds = 3; // Default: 30s, Max: 43,200s (12 hours) int32 wait_time_seconds = 4; // Long polling timeout: 0 to 20s } message ReceiveMessageResponse { repeated QueueMessage messages = 1; } message QueueMessage { string message_id = 1; string receipt_handle = 2; // Ephemeral cryptographic lease token string message_body = 3; int32 receive_count = 4; // Incremented on each consume attempt int64 publish_timestamp_ms = 5; string message_group_id = 6; } message DeleteMessageRequest { string queue_name = 1; string receipt_handle = 2; } message DeleteMessageResponse { bool success = 1; } message ChangeMessageVisibilityRequest { string queue_name = 1; string receipt_handle = 2; int32 visibility_timeout_seconds = 3; } message ChangeMessageVisibilityResponse { bool success = 1; int64 new_visibility_deadline_ms = 2; }
5. Storage Engine Internals & Zero-Copy Architecture
1. Append-Only Segment Log Structure
Every topic partition stores messages sequentially inside fixed-size segment files ( each) on disk:
.log(Segment Data): Sequential binary stream of messages[Length | CRC32 | MagicByte | Attributes | Timestamp | KeyLength | Key | ValLength | Value]..index(Sparse Offset Index): Maps logical message offset to physical byte offset in.log(1 entry every of log data)..timeindex(Sparse Timestamp Index): Maps epoch millisecond timestamp to logical message offset for time-based seeks and TTL expiry.
Synthesizing vector architecture diagram...
What to notice: the index does not list every message, only one entry for about every 4 KB of log data. The panel "Disk Segment 00000000000000000000 (.log)" shows the messages stored one after another: offset 0 holds order_1, offset 1 holds order_2 and offset 2 holds order_3. The panel "Sparse Index (.index)" holds only three entries: logical offset 0 is at byte 0, logical offset 100 is at byte 4096, and logical offset 200 is at byte 8192. To find a message, the broker looks up the largest index entry that is not greater than the wanted offset, jumps to that byte in the .log file, and reads forward from there. For example, a read of offset 150 starts at byte 4096 (offset 100) and scans forward at most about 4 KB.
2. Zero-Copy Kernel Data Path (sendfile)
To eliminate CPU and memory bus bottlenecks during high-throughput consumer egress, the broker relies on the Linux sendfile(2) system call:
Synthesizing vector architecture diagram...
Both panels send the same file data to a consumer's socket. In the "Traditional Context Switching" panel, the data is copied four times: the disk controller copies it into kernel memory (DMA), the CPU copies it into the application's buffer, the CPU copies it back into the kernel's socket buffer, and the network card copies it out (DMA), with a switch between kernel and user mode around each system call. In the "Zero-Copy sendfile()" panel, the application asks the kernel to send the file directly, so the page-cache data never enters user space and the network card reads it straight from kernel memory. The two CPU copies disappear, which is why Kafka can serve consumers at close to network line rate.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~45%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.