TOPIC #110Beginner 8 min read

Message Queues: Producer-Consumer Model & Competing Consumers

CSD
CompleteSystemDesign Editorial
Report an issue
Key takeawayCore Architecture Summary

Master point-to-point worker queues: Enqueue, dequeue, visibility timeouts, heartbeat lease extensions, acknowledgments (ACK/NACK), and horizontal worker scaling.

Key Glossary Concepts in this TopicAll Glossary Terms
Interactive Lab · 📥 Competing Consumer Lease LabFull lab guide

Competing Consumers & Visibility Timeouts

Workers dequeue messages, hold invisibility leases, and crash mid-flight — watch un-ACKed messages resurface.

Available

8

ApproxNumberOfMessages

ACK'd

0

DeleteMessage issued

Redeliveries

0

duplicate handlers required

Autoscaler Target

2

⌈depth × task-time ÷ 60s SLA⌉

QUEUE: worker fleet lease state

W1IDLE

ReceiveMessage (long poll 20s)

W2IDLE

ReceiveMessage (long poll 20s)

W3IDLE

ReceiveMessage (long poll 20s)

SIZING ADVISOR:

Tasks run 4–21s. Rule of thumb: timeout ≈ 3× P99. At 12s the timeout is SHORTER than long tasks → healthy workers see duplicate deliveries.

SQS hides a message for the visibility timeout the moment a worker polls it. An ACK (DeleteMessage) removes it forever; a crash leaves no ACK, so the lease expires and a surviving worker retries with ReceiveCount + 1. That is at-least-once delivery — handlers must be idempotent, and the autoscaler should read queue depth, not CPU.

Competing Consumers & Visibility Timeout Lifecycle

How visibility timeouts and ACKs prevent lost messages when background workers crash.

Competing Consumers & Visibility Timeout Lifecycle
100%
Touchpad: Pinch to zoom • Drag to pan
Rendering visual architecture flowchart...

01.The Producer-Consumer Architecture

The Producer-Consumer pattern decouples work generation from work execution through an intermediate buffer:

  • Producer: Any application service that constructs a structured message (JSON, Protocol Buffers, or Avro) and enqueues it via an API call (SendMessage in SQS or basic.publish in RabbitMQ).
  • Message Queue: A distributed, durable store that holds messages in sequential order, managing delivery state, timeouts, and replication across multiple Availability Zones.
  • Consumer (Worker): A background process that polls or receives messages, executes business logic, and notifies the queue upon completion.

In point-to-point queues, each message is delivered to and processed by exactly one consumer, distinguishing it from broadcast publish-subscribe models.

02.The Competing Consumers Pattern & Auto-Scaling

To process heavy workloads, multiple identical worker instances subscribe to the same queue. This is known as the Competing Consumers Pattern:

  • Dynamic Load Distribution: Workers pull messages from the queue independently. Faster workers naturally pull more messages, while slower workers handle fewer, creating automatic load balancing without a central coordinator.
  • Queue-Depth-Driven Auto-Scaling: Instead of scaling worker fleets based on CPU utilization (which can be misleading for I/O-bound tasks), systems scale on Backlog Depth per Worker:

Target Workers = ≤ft\lceil \frac{Queue Depth (ApproximateNumberOfMessages)}{Target Processing Rate per Worker × Acceptable Latency SLA} \right\rceil

For example, if the queue holds 12,000 messages, each worker handles 10 msgs/sec, and messages must clear within 60 seconds, the auto-scaler provisions \lceil 12,000 / (10 × 60) \rceil = 20 workers.

03.Acknowledgments (ACK/NACK) & Visibility Timeouts

When a worker retrieves a message, the queue cannot immediately delete it—if the worker crashes during execution, the message would be permanently lost. Instead, distributed queues use Visibility Timeouts:

  1. State Transition to In-Flight: When Worker A dequeues Message 1, the queue hides Message 1 from all other workers for a configured Visibility Timeout window (e.g., 30 seconds).
  2. Positive Acknowledgment (ACK): Worker A finishes processing in 4 seconds and issues a DeleteMessage call with the unique ReceiptHandle. The queue permanently purges the message.
  3. Negative Acknowledgment (NACK / Reject): If Worker A detects a transient failure (e.g., downstream database temporarily locked), it sends a NACK, causing the queue to immediately reset visibility so another worker can retry without waiting for the timeout.
  4. Worker Crash & Automatic Failover: If Worker A dies from an Out-Of-Memory (OOM) error or hardware failure, no ACK is received. Once the 30s timeout elapses, Message 1 automatically transitions back to visible state. Worker B polls the queue, receives Message 1 with an incremented ReceiveCount = 2, and completes the task.

04.Heartbeating (Lease Renewal) & Polling Optimizations

Visibility Timeout Sizing & Heartbeating

Setting the visibility timeout too short causes duplicate processing (Worker B starts processing Message 1 while Worker A is still working on it). Setting it too long delays crash recovery.

  • Rule of Thumb: Configure default visibility timeout to 3× your P_{99} task execution time.
  • Heartbeat Lease Renewal: For tasks with variable execution times (e.g., video encoding taking anywhere from 10 seconds to 10 minutes), the worker runs a background heartbeat thread. Every 15 seconds, it calls ChangeMessageVisibility(visibilityTimeout = 30s) to extend its lease until completion.

Short Polling vs Long Polling

  • Short Polling: The queue samples a subset of storage servers and returns immediately, even if no messages are found. This generates millions of empty HTTP 200 OK calls, driving up cloud bills.
  • Long Polling (WaitTimeSeconds = 20): The queue connection remains open for up to 20 seconds until a message arrives, reducing empty responses by over 90\% and reducing delivery latency from seconds to milliseconds.

Architectural Trade-offs & Production Realities

Architectural Advantages

  • Horizontal elasticity: Add or remove worker nodes with zero reconfiguration of the queue or producers
  • Fault resilience: Unacknowledged messages automatically reappear upon worker process failure
  • Natural load balancing: High-speed workers process more tasks without centralized scheduling overhead

Trade-offs & Constraints

  • At-least-once delivery requires workers to be strictly idempotent to avoid duplicate mutations
  • Strict FIFO ordering across multiple competing consumers cannot be guaranteed without partition locks
Production Implementation in Big Tech
Amazon Web Services• Amazon SQS (Simple Queue Service)

AWS SQS powers distributed asynchronous compute for millions of architectures. SQS distributes incoming messages across a redundant cluster of storage servers. When consumers invoke ReceiveMessage with a 20-second long-poll, SQS locks the message with a visibility timeout, tracks receipt handles, and automatically re-surfaces unacknowledged items if worker containers terminate unexpectedly.

Staff+ Engineering Takeaways

  • Competing consumers pull from a shared queue to balance computational workloads dynamically.
  • Visibility timeouts hide in-flight messages from other workers while execution is underway.
  • If a worker crashes without sending an ACK, the visibility timeout expires, enabling automatic failover.
  • Long polling (`WaitTimeSeconds=20`) eliminates empty API responses and cuts cloud billing costs.

Topic Knowledge Check

Exercise 1 of 3 • Test your architectural comprehension.

Exercise 1 of 30 answered
1

What occurs if a worker process crashes due to an Out-Of-Memory (OOM) error while processing a message from AWS SQS?

Rate This Architecture ChapterFeedback & Rating

How clear and actionable was this distributed systems breakdown?

Interactive Engineering Workbenches: