PrepZone Logo
PrepZone

Message Queues Fundamentals

Decouple producers and consumers with durable queues, acknowledgements and back-pressure.

Read these first

When StreamHub finishes transcoding a VOD upload, the API cannot synchronously notify every follower, update search indexes, and send push alerts in one HTTP response. It publishes a video.ready event to a queue and returns immediately while workers handle the rest.

Why queues exist

Problems queues solve

  • Decoupling: Producer and consumer scale, deploy, and fail independently.
  • Spike absorption: A flash sale or viral stream queues work instead of overwhelming downstream services.
  • Reliability: Messages persist on disk until acknowledged — survive consumer crashes.
  • Async workflows: One event triggers multiple consumers (notifications, analytics, search indexing).

Amazon MSK event pipeline

produceCOMPUTE
EKS APIorder placed
INTEGRATION
Amazon MSKorders.placed.v1
COMPUTE
InventoryEKS worker
COMPUTE
Email svcEKS worker
ANALYTICS
AnalyticsFlink / EMR
API publishes events; worker fleets scale independently on consumer lag.

Core concepts

ConceptMeaningStreamHub example
ProducerService that publishes messagesUpload service after transcode completes
ConsumerService that reads and processes messagesNotification worker, search indexer
Queue / TopicDurable buffer between themSQS queue or Kafka topic `video.events`
AcknowledgementConsumer confirms successful processingDelete from SQS or commit Kafka offset
Dead letter queueHolds messages that failed max retriesMalformed transcode events for manual review
  • Producer

    MeaningService that publishes messages
    StreamHub exampleUpload service after transcode completes
  • Consumer

    MeaningService that reads and processes messages
    StreamHub exampleNotification worker, search indexer
  • Queue / Topic

    MeaningDurable buffer between them
    StreamHub exampleSQS queue or Kafka topic `video.events`
  • Acknowledgement

    MeaningConsumer confirms successful processing
    StreamHub exampleDelete from SQS or commit Kafka offset
  • Dead letter queue

    MeaningHolds messages that failed max retries
    StreamHub exampleMalformed transcode events for manual review

Delivery guarantees

AspectAt-most-onceAt-least-once
BehaviourMessage may be lost; never duplicatedMessage always delivered; may arrive twice
MechanismAck before process (fire-and-forget)Ack after process; retry on failure
Consumer requirementNone — loss is acceptableMust be idempotent
StreamHub useLive viewer count heartbeatsPayment events, push notifications
  • Behaviour

    At-most-onceMessage may be lost; never duplicated
    At-least-onceMessage always delivered; may arrive twice
  • Mechanism

    At-most-onceAck before process (fire-and-forget)
    At-least-onceAck after process; retry on failure
  • Consumer requirement

    At-most-onceNone — loss is acceptable
    At-least-onceMust be idempotent
  • StreamHub use

    At-most-onceLive viewer count heartbeats
    At-least-oncePayment events, push notifications

Exactly-once delivery is the holy grail but practically achieved only via idempotent consumers plus transactional outbox — not by the broker alone.

Message anatomy

Java
{
  "message_id": "msg_8f3a2b1c",
  "event_type": "video.ready",
  "timestamp": "2026-04-01T14:22:00Z",
  "payload": {
    "video_id": "vid_9912",
    "streamer_id": "sh_4420",
    "duration_sec": 3600
  },
  "metadata": {
    "trace_id": "abc123",
    "retry_count": 0
  }
}

Back-pressure and flow control

Preventing consumer overload

  • Prefetch limit: Consumer fetches N messages at a time, not the entire queue.
  • Visibility timeout: Message hidden from other consumers while one processes it (SQS pattern).
  • Consumer lag monitoring: Alert when lag exceeds SLA — scale consumers or fix slow processing.
  • Rate limiting downstream: Throttle consumer calls to external APIs (push providers, email).

StreamHub upload pipeline

StreamHub production architecture (AWS)

HTTPSstaticmissAPICLIENT
Mobile / WebStreamHub cli…
NETWORK
Route 53GeoDNS routing
NETWORK
CloudFrontCDN + WAF edge
NETWORK
AWS ALBTLS terminati…
NETWORK
API GatewayJWT · rate li…
STORAGE
Amazon S3media origin
COMPUTE
Amazon EKSAPI · auth · …
DATABASE
ElastiCachesessions · ho…
DATABASE
RDS Postgresprimary + rep…
INTEGRATION
Amazon MSKdomain events
ANALYTICS
OpenSearchstream discov…
OPS
CloudWatchmetrics · X-R…
End-to-end path from user to data — reference this when placing any new service.

Upload API → SQS transcode-jobs → worker fleet → Kafka video.events → notification, search, and analytics consumers.

Quick recall

Everything you need if you only revisit this box.

  • Queues decouple producers from consumers and absorb traffic spikes.
  • At-least-once is the practical default — design idempotent consumers.
  • Dead letter queues capture poison messages after max retries.
  • Prefetch limits and visibility timeouts prevent consumer overload.
  • One event can fan out to multiple consumers via pub/sub topics.
  • Monitor consumer lag — it is the primary queue health signal.

Test yourself

Answer these before moving on — recall is what makes it stick.