Skip to content

09 — Message Queues

Message queues enable asynchronous communication between services. A producer sends a message to a queue, and a consumer processes it when ready — decoupling the two services in time and space.

Analogy: A message queue is like an email inbox. You send an email (produce a message), and the recipient reads and acts on it when they have time (consume). You don’t need to wait on the phone for them to respond.


Without message queues:

  • Synchronous coupling — a slow consumer blocks the producer
  • Traffic spikes — burst traffic crashes downstream services
  • Retry complexity — failed operations require complex retry logic
  • Fan-out — sending the same data to multiple consumers is complex

flowchart LR
subgraph Producers["Producers"]
P1["Web App"]
P2["Mobile App"]
P3["Cron Job"]
end
subgraph Queue["Message Queue"]
Q["Queue / Topic / Stream"]
end
subgraph Consumers["Consumers"]
C1["Email Service"]
C2["Notification Service"]
C3["Analytics Pipeline"]
C4["Search Index"]
end
P1 & P2 & P3 --> Queue
Queue --> C1 & C2 & C3 & C4
style Producers fill:#3b82f6,color:#fff
style Queue fill:#7c3aed,color:#fff
style Consumers fill:#059669,color:#fff

ModelDescriptionExample
Point-to-PointOne message → one consumerSQS, RabbitMQ work queues
Pub/Sub (Fan-out)One message → all subscribersSNS, Kafka consumer groups
StreamingOrdered log of events, replayableKafka, Kinesis
Dead Letter Queue (DLQ)Failed messages after retriesAll queue systems

flowchart TB
MSG["Message Systems"] --> Queue["Queue<br/>SQS, RabbitMQ<br/>Point-to-point<br/>Message deleted after consume"]
MSG --> PubSub["Pub/Sub<br/>SNS, RabbitMQ<br/>Broadcast to all<br/>subscribers"]
MSG --> Stream["Stream<br/>Kafka, Kinesis<br/>Ordered log<br/>Replayable"]
MSG --> DLQ["Dead Letter Queue<br/>Failed messages<br/>Retry policies"]
style MSG fill:#7c3aed,color:#fff
style Queue fill:#3b82f6,color:#fff
style PubSub fill:#059669,color:#fff
style Stream fill:#f59e0b,color:#fff
style DLQ fill:#ef4444,color:#fff

sequenceDiagram
participant User as User
participant App as Web App
participant Queue as Message Queue
participant Worker as Worker (Consumer)
participant DB as Database
participant Email as Email Service
User->>App: Place order
App->>DB: Save order (status: processing)
App->>Queue: SendMessage({ orderId: 123 })
App-->>User: 202 Accepted (order received)
Note over Worker: Worker polls queue
Worker->>Queue: ReceiveMessage
Queue-->>Worker: { orderId: 123 }
Worker->>DB: Update order (processing)
Worker->>Worker: Validate payment
Worker->>Email: Send confirmation email
Worker->>DB: Update order (confirmed)
Worker->>Queue: DeleteMessage (done)
alt Failed after retries
Queue->>DLQ["Dead Letter Queue"]: Move message
end

TechnologyTypePersistenceOrderingBest For
Apache KafkaStreamDisk (configurable)Partition-orderEvent streaming, logs, metrics
RabbitMQQueueDisk or MemoryOptionalTask queues, pub/sub, RPC
AWS SQSQueueDisk (AWS)Best-effort (Standard)Decoupling microservices
AWS SNSPub/SubNone (push)N/AFan-out notifications
Redis StreamsStreamMemory/DiskYesReal-time, lightweight
Google Pub/SubPub/SubDisk (GCP)Yes (ordered)Event-driven, serverless

flowchart TB
subgraph Producers["Producers"]
Prod1["App Events"]
Prod2["DB Changes"]
Prod3["Metrics"]
end
subgraph Kafka["Apache Kafka Cluster"]
Topic1["Topic: user-events<br/>Partition 0 | 1 | 2"]
Topic2["Topic: page-views<br/>Partition 0 | 1"]
end
subgraph Consumers["Consumer Groups"]
CG1["Analytics Group"]
CG2["Notification Group"]
end
Prod1 & Prod2 --> Topic1
Prod3 --> Topic2
Topic1 --> CG1 & CG2
Topic2 --> CG1
style Producers fill:#3b82f6,color:#fff
style Kafka fill:#7c3aed,color:#fff
style Consumers fill:#059669,color:#fff

ConceptDescription
ProducerSends messages to a queue/topic
ConsumerReceives and processes messages
TopicNamed channel (Kafka, SNS)
PartitionSub-division of a topic (parallelism in Kafka)
Consumer GroupGroup of consumers sharing work (Kafka)
BrokerServer hosting queue/topic partitions
Retry PolicyHow many times to retry failed messages
DLQ (Dead Letter Queue)Final destination for messages that can’t be processed
Message TTLHow long a message lives before auto-deletion

PatternDescriptionUse Case
Work QueueMultiple workers compete for messagesImage processing, email sending
Pub/SubEach subscriber gets every messageNotifications, audit logs
Event SourcingStore all events, rebuild stateAudit trails, financial systems
Transactional OutboxDB transaction → message to queueReliable cross-service communication
SagaDistributed transaction via messagesMicroservice transactions

ChoiceProsCons
SynchronousSimple, consistentCoupled, slower, no retry
Async with queueDecoupled, resilientEventual consistency, complexity
KafkaHigh throughput, replayableHeavy operations, higher latency
SQSFully managed, simpleLimited ordering (Standard), vendor lock
At-least-once deliveryNo message lossDuplicate processing (idempotency needed)
Exactly-onceNo duplicatesLower throughput, complex

StrategyDescription
Increase partitionsMore parallelism (Kafka)
More consumersScale consumer count within group
Batch processingConsume messages in batches
Rate limitingControl how fast consumers process
Priority queuesCritical messages processed first

  1. What’s the difference between a queue and a pub/sub system?
  2. How does Kafka achieve high throughput?
  3. How do you handle duplicate messages in a queue?
  4. What is the transactional outbox pattern?
  5. How would you design a reliable notification system with message queues?

SystemMessage Queue Usage
UberKafka for ride events, trip tracking, analytics
NetflixKafka for event streaming (2M+ events/sec)
LinkedInKafka invented at LinkedIn for activity streams
AmazonSQS for order processing, inventory updates

  • Message queues decouple producers (senders) from consumers (processors)
  • SQS/RabbitMQ = simple queues — one message consumed by one worker
  • Kafka = event stream — ordered, replayable, high throughput
  • SNS/Pub-Sub = one message broadcast to many subscribers
  • DLQ (Dead Letter Queue) catches messages that failed after all retries
  • At-least-once delivery is the default — design your consumers to be idempotent
  • Use queues for async processing, load smoothing, and decoupling services