Competing Consumers vs Pub/Sub Fan-Out

Two consumption patterns look similar on a diagram and behave completely differently: competing consumers, where several workers share one queue and each message is processed once, and fan-out, where every subscriber gets its own copy. This guide shows when each applies and how to build them on common brokers, as part of Producer-Consumer Pattern Design in Queue Fundamentals & Architecture.

Problem Statement

An order service publishes OrderPlaced to a single RabbitMQ queue. Three teams consume it: fulfilment reserves stock, email sends confirmations, and analytics records the sale. All three subscribed to the same queue — and each order was handled by only one of them, whichever grabbed the message first. Two-thirds of orders were missing confirmations or analytics. After a hotfix gave each team its own queue, a new problem appeared: when the analytics consumer was down for a day, its queue grew to 2 million messages and the broker hit its memory alarm, blocking publishing for everyone. You want each team to receive every event, each team to scale its own workers, and one team's outage to never affect the others.

Prerequisites

  • A broker that supports either exchanges/topics with multiple bound queues (RabbitMQ, SNS+SQS, Google Pub/Sub) or a log with independent consumer groups (Kafka, Redis Streams, NATS JetStream).
  • A list of consumers per event type and whether each needs every event.
  • Queue depth limits or storage budgets per subscriber.

Step 1 — Distinguish "Do This Work" from "This Happened"

The pattern follows from what the message means.

  • A command ("resize image 42", "send email 918") is work to be done once. Many workers share it: competing consumers on one queue.
  • An event ("order 77 was placed") is a fact that several independent parties care about. Each party needs its own copy: fan-out, with each subscriber then using competing consumers within its own queue to scale.
OrderPlaced (event)
  -> fulfilment queue  -> 6 fulfilment workers compete
  -> email queue       -> 2 email workers compete
  -> analytics queue   -> 1 analytics worker

Mixing the two — several teams on one queue — is the bug in the problem statement: it treats an event as a command, so each event is "done" by whoever takes it first.

Share the work or copy the fact On the left, one queue feeds three workers of the same service; each message goes to exactly one of them. On the right, an OrderPlaced event is published to an exchange or topic that copies it into three queues, one each for fulfilment, email, and analytics; each subscriber receives every event and scales its own workers. Competing consumers vs fan-out competing: each message once resize queue worker 1 worker 2 worker 3 fan-out: every subscriber, every event OrderPlaced fulfilment queue email queue analytics queue

Step 2 — Build Fan-Out on RabbitMQ with an Exchange per Event Stream

Publish to an exchange, not a queue. Each subscriber declares and owns its queue and binds it to the exchange; the exchange copies each message into every bound queue.

# publisher: knows only the exchange
ch.exchange_declare("orders", exchange_type="topic", durable=True)
ch.basic_publish("orders", routing_key="order.placed", body=json.dumps(evt),
                 properties=pika.BasicProperties(delivery_mode=2, message_id=evt["event_id"]))

# each subscriber owns its queue and its limits
def declare_subscriber(name: str, max_len: int):
    ch.queue_declare(f"orders.{name}", durable=True, arguments={
        "x-queue-type": "quorum",
        "x-max-length": max_len,                       # this subscriber's own ceiling
        "x-overflow": "reject-publish-dlx",            # overflow goes to its DLX, not everyone's alarm
        "x-dead-letter-exchange": f"orders.{name}.dlx",
    })
    ch.queue_bind(f"orders.{name}", "orders", routing_key="order.*")

declare_subscriber("fulfilment", max_len=500_000)
declare_subscriber("email", max_len=500_000)
declare_subscriber("analytics", max_len=2_000_000)

x-max-length with an overflow behaviour is the isolation mechanism the hotfix lacked: a stalled analytics consumer fills its queue to its limit, and overflow goes to its dead-letter exchange, instead of consuming broker memory until the alarm blocks every publisher.

Step 3 — Build Fan-Out on AWS with SNS and SQS

On AWS, SNS is the exchange and each subscriber has its own SQS queue. Filter policies replace routing keys.

resource "aws_sns_topic" "orders" { name = "orders-events" }

module "subscriber" {
  for_each = { fulfilment = ["OrderPlaced", "OrderCancelled"], email = ["OrderPlaced"], analytics = [] }
  source   = "./sqs-subscriber"          # creates queue + DLQ + redrive policy
  name     = each.key
}

resource "aws_sns_topic_subscription" "sub" {
  for_each             = module.subscriber
  topic_arn            = aws_sns_topic.orders.arn
  protocol             = "sqs"
  endpoint             = each.value.queue_arn
  raw_message_delivery = true
  filter_policy        = length(module.subscriber[each.key].event_types) > 0 ? jsonencode({
    event_type = module.subscriber[each.key].event_types }) : null
}

Each SQS queue scales independently and holds up to 14 days of backlog without affecting the others — isolation is built in, which is one reason SNS+SQS is the default fan-out on AWS. Per-subscriber DLQs follow configuring an SQS redrive policy.

One subscriber's outage stays its own With a queue per subscriber and a length limit on each, the analytics consumer's outage grows only the analytics queue up to its limit. The fulfilment and email queues keep draining at their normal depth, and publishers are never blocked. With a shared broker and no limits, the analytics backlog would consume broker memory and eventually block every publisher. Queue depth during an analytics outage fulfilment ~200, draining email ~150, draining analytics growing toward its own limit x-max-length Publishers never block.

Step 4 — Build Fan-Out on a Log with Consumer Groups

Kafka, Redis Streams, and NATS JetStream store each event once and let every consumer group keep its own position. Fan-out is simply one consumer group per subscriber; competing consumers are several members within a group.

# Kafka: each team is a consumer group on the same topic
fulfilment = Consumer({"bootstrap.servers": B, "group.id": "fulfilment", "enable.auto.commit": False})
email      = Consumer({"bootstrap.servers": B, "group.id": "email",      "enable.auto.commit": False})
analytics  = Consumer({"bootstrap.servers": B, "group.id": "analytics",  "enable.auto.commit": False})
for c in (fulfilment, email, analytics):
    c.subscribe(["orders"])
# Scaling a team = running more members of its group, up to the partition count.

Isolation works differently here: a stalled group does not add storage (the event is stored once), it just falls behind in lag. The risk is falling behind the topic's retention — if analytics is down longer than retention.ms, events it never read are deleted. Alert on consumer-group lag in time against retention.

One log, many positions The orders topic stores each event once. The fulfilment group's committed offset is near the head, the email group is slightly behind, and the analytics group is far behind after an outage. Storage does not grow with the lag, but if analytics falls behind the retention boundary, the oldest events it has not read are deleted. Consumer groups on one topic orders topic: retention 7 days, stored once fulfilment email analytics: 6 days behind retention boundary: unread events deleted past here

The two families therefore fail differently under a stalled subscriber. Queue-based fan-out fails loudly and locally: the stalled subscriber's queue fills to its limit and overflows to its DLQ, while everyone else carries on. Log-based fan-out fails quietly and late: nothing looks wrong for days, and then the lagging group silently loses the events that aged out. Both are manageable, but they need different alerts — queue depth against the length limit for the first, lag in time against retention for the second. Log semantics versus queue semantics are compared in Kafka vs RabbitMQ for task queues.

Step 5 — Scale Each Side Independently

Within each subscriber, competing consumers scale throughput. Scale on that subscriber's own backlog, never on the aggregate.

# KEDA: fulfilment scales on its own queue, independent of analytics
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata: { name: fulfilment-workers }
spec:
  scaleTargetRef: { name: fulfilment-worker }
  minReplicaCount: 2
  maxReplicaCount: 30
  triggers:
    - type: rabbitmq
      metadata: { queueName: orders.fulfilment, mode: QueueLength, value: "200" }

For the log-based design, the scaling ceiling per group is the partition count; for queue-based designs it is limited only by the broker and downstream capacity. The mechanics are in scaling workers with KEDA on queue length.

Step 6 — Add a New Subscriber Without Touching Publishers

The test of a good fan-out design is adding a fourth team. With exchanges or SNS, the new team declares its queue and binding; the publisher does not change. With a log, the new team starts a consumer group, optionally from the earliest retained offset to backfill history.

# RabbitMQ: new subscriber, zero publisher changes
rabbitmqadmin declare queue name=orders.loyalty durable=true arguments='{"x-queue-type":"quorum","x-max-length":500000}'
rabbitmqadmin declare binding source=orders destination=orders.loyalty routing_key='order.placed'

# Kafka: new group, backfill from the start of retention
kafka-consumer-groups.sh --bootstrap-server $B --group loyalty --topic orders --reset-offsets --to-earliest --execute

A queue-based subscriber only sees events published after its binding exists; if it needs history, it must be backfilled from another source. That difference often decides between the two families.

Verification

# Every subscriber received every event over a test window
for q in fulfilment email analytics; do
  rabbitmqctl list_queues name messages_ready message_stats.publish_details.rate | grep "orders.$q"
done
# Publish 10,000 test events; each subscriber's processed count must be 10,000

Then stop one consumer and confirm the others' queue depths, latency, and the publisher's publish rate are unaffected, and that the stopped subscriber's queue stops growing at its length limit.

Gotchas & Edge Cases

Shared queue by accident. Two teams using the same queue name in configuration silently turn fan-out into competing consumers. Name queues by subscriber (orders.email), and have each team own its declaration.

Unbounded subscriber queues. Without per-queue limits, the slowest subscriber sets the broker's memory usage. Limit every queue.

Ordering across subscribers. Each subscriber sees events in publish order per queue or partition, but subscribers progress independently; never assume email has run before fulfilment.

Duplicate delivery per subscriber. Each subscriber must be idempotent on its own; fan-out multiplies deliveries, and each copy can be redelivered.

FAQ

Can one queue serve both patterns? No. A queue gives each message to one consumer. Fan-out needs a queue (or consumer group) per subscriber.

Pub/sub without a broker? Redis PUBLISH/SUBSCRIBE delivers only to connected subscribers and keeps nothing — fine for live notifications, not for work that must not be lost. Use streams or a broker for durable fan-out.

Should the publisher know who its subscribers are? No — that is the point of fan-out. The publisher writes to an exchange, topic, or log and knows nothing about who consumes it; subscribers own their queues and bindings. If the publisher enqueues directly into each team's queue, every new consumer needs a publisher change, and the publisher becomes coupled to every downstream team's availability.

How many subscribers is too many? Queue-based fan-out copies each message per subscriber, so storage and broker work grow linearly. Beyond a few dozen subscribers on a high-volume stream, a log with consumer groups is usually more efficient.

Related