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
Core concepts
| Concept | Meaning | StreamHub example |
|---|---|---|
| Producer | Service that publishes messages | Upload service after transcode completes |
| Consumer | Service that reads and processes messages | Notification worker, search indexer |
| Queue / Topic | Durable buffer between them | SQS queue or Kafka topic `video.events` |
| Acknowledgement | Consumer confirms successful processing | Delete from SQS or commit Kafka offset |
| Dead letter queue | Holds messages that failed max retries | Malformed transcode events for manual review |
Producer
MeaningService that publishes messagesStreamHub exampleUpload service after transcode completesConsumer
MeaningService that reads and processes messagesStreamHub exampleNotification worker, search indexerQueue / Topic
MeaningDurable buffer between themStreamHub exampleSQS queue or Kafka topic `video.events`Acknowledgement
MeaningConsumer confirms successful processingStreamHub exampleDelete from SQS or commit Kafka offsetDead letter queue
MeaningHolds messages that failed max retriesStreamHub exampleMalformed transcode events for manual review
Delivery guarantees
| Aspect | At-most-once | At-least-once |
|---|---|---|
| Behaviour | Message may be lost; never duplicated | Message always delivered; may arrive twice |
| Mechanism | Ack before process (fire-and-forget) | Ack after process; retry on failure |
| Consumer requirement | None — loss is acceptable | Must be idempotent |
| StreamHub use | Live viewer count heartbeats | Payment events, push notifications |
Behaviour
At-most-onceMessage may be lost; never duplicatedAt-least-onceMessage always delivered; may arrive twiceMechanism
At-most-onceAck before process (fire-and-forget)At-least-onceAck after process; retry on failureConsumer requirement
At-most-onceNone — loss is acceptableAt-least-onceMust be idempotentStreamHub use
At-most-onceLive viewer count heartbeatsAt-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
{
"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)
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.