Lesson 06

Message Queues

Kafka, RabbitMQ, pub/sub, delivery guarantees (at-most, at-least, exactly-once).

Plays in the sticky player at the bottom of the page

Transcript

Lesson 6: Message Queues

A reference-style deep dive into Kafka, RabbitMQ, pub/sub, delivery guarantees (at-most, at-least, exactly-once). Read this as an article, not a transcript. The accompanying audio is the spoken companion; the article below is the canonical written reference.

Audience. Engineers designing, building, or operating distributed systems who want a clear mental model rather than a checklist of tools.

Prerequisites. Working knowledge of HTTP, basic SQL, and the idea of running more than one server behind a load balancer.

Table of contents

  1. Why message queues matters
  2. The core mental model and where to start 3-12. In-depth sections below (see "Lesson body")

Lesson diagram

Lesson 6 diagram — Message Queues

Figure 1. The canonical message queues topology and control flow covered in this lesson.


1. Why queues matter

A message queue decouples the producer of work from the consumer. The producer writes once; the queue durably stores the message; one or more consumers process at their own pace. Queues turn synchronous coupling into asynchronous coupling. The cost is latency and complexity; the benefit is elasticity, durability, and retry.

Three motivations drive queue adoption:

  • Smooth bursts. The producer writes at peak rate; consumers process at average rate. The queue absorbs the difference.
  • Decouple systems. The producer does not know or care which consumers exist.
  • Durability. Messages survive consumer crashes; processing resumes when the consumer is back.

2. Delivery guarantees — the central trade-off

A queue makes three promises:

At-most-once

Messages may be lost but never duplicated. The broker delivers once and forgets.

queue.publish('orders', msg)
# If the broker crashes before persisting, the message is lost.
# If the consumer crashes, the message is lost.

Fast, lossy. Used for: telemetry, metrics, anything where the cost of a duplicate or a missed message is low.

At-least-once

Messages are never lost but may be duplicated. The broker retries until the consumer acks.

queue.publish('orders', msg)
# Broker retries until consumer acks.
# If the consumer crashes after processing but before ack, the message is redelivered.

The default for most production queues. Safe if the consumer is idempotent.

Exactly-once

Messages are never lost and never duplicated. Requires idempotent processing AND broker-side dedup.

# Producer: idempotent (no duplicate publishes from retries)
producer.send('orders', key=key, value=value)  # broker dedups by key
# Consumer: idempotent (safe to process the same message twice)
def process(msg):
    if db.exists(msg.id):  # already processed
        return
    db.apply(msg)
    db.mark(msg.id)

The hardest to achieve. Most systems that claim "exactly-once" mean "at-least-once + idempotent processing", which is effectively exactly-once when the consumer is correctly idempotent.

3. Kafka — the log-based broker

Kafka's core abstraction is a partitioned, append-only log.

topic: orders
  partition 0: [msg_0, msg_1, msg_2, msg_3, ...]  (offset 0, 1, 2, 3, ...)
  partition 1: [msg_0, msg_1, msg_2, ...]
  partition 2: [msg_0, msg_1, ...]

Producers append to the log; consumers track their own offset.

# Producer
producer.send('orders', key=b'order-42', value=b'...')
producer.flush()

# Consumer
consumer = KafkaConsumer('orders', group_id='billing')
for msg in consumer:
    process(msg)
    consumer.commit()  # advance offset

Properties:

  • High throughput. Millions of messages per second per cluster.
  • Durable. Messages persist on disk; replayable by resetting offset.
  • Partitioned. Order is preserved within a partition; cross-partitioned only by key.
  • Consumer groups. Multiple consumers share partitions; each partition is processed by one consumer in the group.

4. RabbitMQ — the queue-based broker

RabbitMQ's core abstraction is a queue + exchange routing.

producer → exchange → queue (bound by routing key) → consumer

Three exchange types:

  • Direct — exact routing key match.
  • Topic — pattern match (e.g. orders.*.us).
  • Fanout — broadcast to all bound queues.
channel.exchange_declare('orders', exchange_type='topic')
channel.queue_declare('billing', durable=True)
channel.queue_bind('billing', 'orders', routing_key='order.placed.*')

channel.basic_publish('orders', 'order.placed.us', body)

Properties:

  • Rich routing. Topic exchanges let you route by any string pattern.
  • Per-message ack. Consumers ack individual messages; no offset tracking.
  • No replay. Once a message is acked, it's gone.
  • Lower throughput than Kafka; great for transactional workloads.

5. NATS — the lightweight pub/sub

NATS is built for speed and simplicity.

nc = nats.connect('nats://localhost:4222')
nc.publish('orders.new', json.dumps(order).encode())

async def handler(msg):
    await process(msg)
await nc.subscribe('orders.new', cb=handler)

Properties:

  • At-most-once by default. No persistence; messages are fire-and-forget.
  • Extremely fast. Sub-millisecond latency; millions of messages per second.
  • Subject-based routing. Hierarchical subjects with wildcard subscriptions.
  • Optional persistence via JetStream (adds durability and replay).

6. Pub/sub vs work queue

Two fan-out patterns, with different semantics:

Pub/sub (fan-out)

Each subscriber gets a copy of each message. Use when multiple independent consumers need the same event.

# Pub/sub
nc.publish('user.signup', msg)         # every subscriber gets it
# Subscribers: email_service, analytics_service, billing_service

Work queue

Each message is delivered to exactly one consumer. Use when the work is parallelizable.

# Work queue
channel.basic_publish('image.process', body)
# Only one of N consumers gets this message; they share the load.

Both are useful. Most production systems use both: pub/sub for events, work queue for jobs.

7. Ordering guarantees

Most queues make no global ordering promise. Within a partition (Kafka) or queue (RabbitMQ), order is preserved.

For ordered processing:

  • Use a single partition per ordering key. All messages for a given user go to the same partition.
  • Use a single consumer per partition. Otherwise the consumer reorders.
# Kafka: hash key ensures same partition
producer.send('orders', key=b'user-42', value=msg)   # always partition 0
producer.send('orders', key=b'user-99', value=msg)   # always partition 2

The trade-off: per-key ordering means per-key throughput is limited by one partition.

8. Dead letter queues

A message that fails processing repeatedly is a poison message. It must be quarantined, not retried forever.

# RabbitMQ: DLX configuration
channel.queue_declare(
    'orders.process',
    arguments={
        'x-dead-letter-exchange': 'orders.dlx',
        'x-message-ttl': 60000,   # 60 seconds
    }
)
# After 60s of retries, the message goes to the dead-letter queue.
# Operators inspect it later.

A dead letter queue is a separate queue that catches the poison messages. Operators alert on its depth.

9. Backpressure and queue depth

When consumers fall behind, the queue grows. Without backpressure, the queue fills memory or disk and the broker crashes.

# Prometheus: queue depth alert
alert: QueueDepthHigh
expr: max(queue_depth_messages) > 100000
for: 5m

Three responses to backpressure:

  1. Alert + scale consumers. The cleanest response.
  2. Drop messages. For low-priority queues where loss is acceptable.
  3. Block producers. The producer's send call blocks until the queue has space.

10. Idempotent processing

The single most important property of a consumer. Idempotent means the same message processed twice has the same effect as processing it once.

# Idempotent consumer
def process(msg):
    if redis.exists(f"processed:{msg.id}"):
        return   # already processed
    db.apply(msg)
    redis.set(f"processed:{msg.id}", 1, ex=86400)  # 24h TTL

Idempotency makes at-least-once safe. Without it, every retry is a potential duplicate.

Patterns for idempotency:

  • Idempotency key — every message has a unique id; consumer dedups by it.
  • Database constraints — unique indexes prevent duplicate writes.
  • Compare-and-set — only apply if the current state matches the expected prior state.

11. Schema management

Messages are contracts between producers and consumers. Schema changes must be coordinated.

# Avro with schema registry
producer.send('orders', value=avro_encode({'order_id': 42, 'total': 99.99}, schema))
consumer.decode(msg.value, schema)

Without a schema registry, producers and consumers can drift: the producer adds a field, the consumer breaks. With a schema registry, schema versions are enforced; consumers see only what they understand.

12. Key takeaways

  • At-least-once + idempotent consumer = effectively exactly-once.
  • Kafka is the log; RabbitMQ is the queue; NATS is the bus. Pick the right primitive.
  • Pub/sub for events; work queue for jobs.
  • Per-key ordering requires per-key partitioning.
  • Dead letter queues catch poison messages; monitor their depth.
  • Backpressure must be explicit; unbounded queues are bugs.

Appendix: terms

  • At-most-once / at-least-once / exactly-once — the three delivery guarantees.
  • Partition — an ordered log within a Kafka topic.
  • Consumer group — a set of consumers that share a topic's partitions.
  • Dead letter queue (DLQ) — a queue that catches poison messages.
  • Idempotent — produces the same effect whether processed once or many times.
  • Schema registry — a service that enforces message schemas.

Appendix: source dialogue excerpt

The audio for this lesson was synthesized from the following Cantonese dialogue (verbatim, not translated):

  • M: 各位同學早晨, 我係子謙。歡迎收聽系統架構課程第六課。今日嘅主題係 Message Queues。…
  • F: 大家好, 我係曉晴。Message queue 係 asynchronous communication 嘅 infrastructure, 將 producer 同 consumer 解耦。今日我哋會拆解 Kafka 同 RabbitMQ 嘅 design 同 trade-off, 同埋 delivery seman…
  • M: 首先講解基本概念。Message queue 嘅核心係 producer 將 message send 去 queue, consumer 從 queue receive message。Producer 同 consumer 唔需要同時 online, 即係 producer 可以 send message 之後即時…
  • F: Message queue 嘅 fundamental 用途有三個。第一個係 decoupling, 即係 producer 同 consumer 唔需要知道對方嘅存在, 只需要知道 message 格式。第二個係 buffering, 即係 producer 嘅 burst 可以 queue 喺 broker, co…
  • M: 好, 第一個 important concept 係 delivery guarantee。常見嘅 guarantee 有三種, at most once, at least once, exactly once。At most once 即係 message 可能 lost 但係永遠唔會 duplicate, 即係 …
  • F: Exactly once 即係 message 唔會 lost 亦唔會 duplicate, 即係每個 message 只會 process 一次。Exactly once 係理論上 impossible 喺 asynchronous network 入面, 因為 consumer process 嘅 failure …

Full dialogue contains 28 segments; see script_raw.json in the source folder.

Lesson quiz · 30 questions

Question 1 of 30Answered 0 / 30
Question 1 of 30

A message queue provides which decoupling?

Pick an answer to lock it in. We'll tell you immediately whether you got it right and show an explanation. Then press Enter or click Next to continue.

Shortcuts:ABCDpick answer on current questionEntergo to next unanswered
30 unanswered