Message Ordering Guarantees in Job Queues

Ordering is the guarantee engineers assume they have and most often do not, and this guide treats it as one of the core contracts of Queue Fundamentals & Architecture. A queue that is "FIFO" on the broker side can still execute jobs out of order the moment you run two workers, retry a failed message, or dead-letter one step of a sequence — and the resulting bugs (a refund processed before its charge, an address update overwritten by the stale one it replaced) are among the hardest to reproduce because they depend on timing.

The useful question is never "is this queue ordered?" but "ordered with respect to what, and at what cost to throughput?" Global ordering across every message serialises the whole system through a single consumer. Per-key ordering — every message for one customer, order, or aggregate processed in sequence, while different keys run in parallel — is what nearly every production design actually needs, and it is what SQS FIFO message groups, Kafka partition keys, and RabbitMQ's single active consumer each provide in different shapes.

The Scenario: When Order Starts to Matter

Consider an order service that publishes three events for the same order within 40 milliseconds: OrderCreated, OrderPaid, and OrderShipped. A downstream worker projects them into a read model and triggers emails. On a plain competing-consumers queue with four workers, the three messages are pulled by three different workers. OrderShipped can finish first, the projection writes status = shipped, and then OrderCreated lands and overwrites it with status = created. The customer receives a shipping email for an order the dashboard says has not been paid.

Nothing is broken in isolation. The broker delivered in FIFO order; each worker processed its message correctly. The failure is that delivery order and completion order are different things once there is any parallelism, and the read model was implicitly relying on completion order.

Delivery order is not completion order A queue hands three events for the same order to three workers in sequence. Because the workers run in parallel and take different amounts of time, the shipped event completes before the created event, and the last write to the read model is the oldest state. Delivered 1-2-3, completed 3-2-1 queue (FIFO) created, paid, shipped worker A: created (90 ms) worker B: paid (55 ms) worker C: shipped (20 ms) read model last write wins: "created" Completion order follows processing time, not queue position. The slowest event wins the final write.

Three separate mechanisms reorder work, and a design has to account for all of them:

  • Parallel consumers. Two messages in flight at once can complete in either order. This is the dominant source of reordering and the one people forget because it is invisible in a single-worker development environment.
  • Retries and redelivery. A message that fails and is retried with backoff re-enters the stream behind later messages for the same key. Under at-least-once delivery, a visibility-timeout expiry does the same thing without any explicit failure.
  • Producer-side races. Two API servers publishing events for the same entity have no shared clock. Whichever publish call reaches the broker first is "first", regardless of which business action happened first.

Architectural Overview: Where Ordering Lives

Ordering is a property of a lane: a sequence of messages that is delivered to at most one consumer at a time and advanced only when the head message is acknowledged. Every ordering mechanism in the common brokers is a way of mapping keys onto lanes.

Broker Lane primitive Key that selects the lane Parallelism unit
SQS FIFO Message group MessageGroupId Number of active groups
Kafka Partition Record key (hashed) Partition count
RabbitMQ Queue with single active consumer Routing key to a queue Number of queues
Redis Streams Stream + one consumer Stream name Number of streams
BullMQ Group (Pro) or one queue at concurrency 1 Group id / queue name Groups / queues

The mapping has a direct consequence: throughput scales with the number of independent lanes, not with the number of workers. Add a hundred workers to a Kafka topic with twelve partitions and eighty-eight of them sit idle. The design work is choosing a key fine-grained enough to produce many lanes but coarse enough that everything which must stay in order shares one.

# ordering-design.yaml — record the ordering contract next to the queue definition
queue: order-events
ordering:
  scope: per-key            # none | per-key | global
  key: order_id             # everything for one order shares a lane
  lanes: 64                 # Kafka partitions / expected active SQS groups
  on_failure: block-lane    # block-lane | skip-and-park | reorder-ok
  max_lane_stall: 300s      # alert if a lane head is stuck longer than this
consumers:
  instances: 16             # useful maximum = lanes; more only adds standby capacity
  in_flight_per_lane: 1     # anything above 1 breaks ordering within the key

The on_failure field is the decision most teams never write down. When the head of a lane fails, you can block the lane until it succeeds (strict order, risk of a stalled key), park the failing message and continue (order broken for that key, lane keeps moving), or declare that this consumer tolerates reordering and rely on versioning instead. Each broker has a default, and the defaults differ.

Keys map to lanes; lanes map to consumers Messages for many order ids are hashed onto four lanes. Each lane is consumed by exactly one worker at a time, so messages with the same key stay in order while different lanes proceed in parallel. A fifth worker has no lane and stays idle. Parallelism = number of lanes hash(order_id) mod 4 same key, same lane lane 0: o-17, o-17, o-42 lane 1: o-08, o-51 lane 2: o-23, o-23, o-23 lane 3: o-90 worker 1 worker 2 worker 3 worker 4 worker 5 idle Hot keys (o-23) serialise on their lane; adding workers beyond the lane count adds no throughput.

Implementation 1: Per-Key Ordering on a Managed FIFO Queue

On AWS, SQS FIFO queues implement lanes as message groups. A message group is delivered strictly in order, and SQS will not hand out the next message in a group while an earlier one is in flight. Different groups are delivered in parallel. The full walkthrough is in FIFO ordering with SQS message group IDs; the essential producer and consumer look like this:

# producer.py — publish order events so each order is its own ordered lane
import json
import uuid
import boto3

sqs = boto3.client("sqs")
QUEUE_URL = "https://sqs.eu-west-1.amazonaws.com/123456789012/order-events.fifo"

def publish(order_id: str, event_type: str, payload: dict, event_id: str | None = None) -> None:
    event_id = event_id or str(uuid.uuid4())
    sqs.send_message(
        QueueUrl=QUEUE_URL,
        MessageBody=json.dumps({"type": event_type, "order_id": order_id, **payload}),
        MessageGroupId=order_id,          # the lane: strict order within one order
        MessageDeduplicationId=event_id,  # 5-minute dedup window; make it the event id
    )


# consumer.py — process one group's messages in sequence, never skipping the head
def consume() -> None:
    while True:
        resp = sqs.receive_message(
            QueueUrl=QUEUE_URL,
            MaxNumberOfMessages=10,       # may contain several groups, each in order
            WaitTimeSeconds=20,           # long polling
            VisibilityTimeout=60,
            AttributeNames=["MessageGroupId", "ApproximateReceiveCount"],
        )
        for msg in resp.get("Messages", []):
            try:
                handle(json.loads(msg["Body"]))
            except Exception:
                # Stop processing this batch: later messages in the same group must
                # not run before this one. Leaving them undeleted blocks the group
                # until the visibility timeout expires and SQS redelivers in order.
                break
            sqs.delete_message(QueueUrl=QUEUE_URL, ReceiptHandle=msg["ReceiptHandle"])

Two details matter. First, the break on failure: a receive batch can contain several messages from the same group, and processing message 3 after message 2 failed is exactly the reordering you were paying FIFO to prevent. Second, the dead-letter queue: with a redrive policy, a poison message at the head of a group will eventually move to the DLQ, and the group then continues without it. That is the "skip-and-park" failure policy, applied automatically. If downstream state cannot tolerate a missing step, the consumer must check for gaps — see handling out-of-order events with sequence numbers.

Implementation 2: Partitioned Logs and Single Active Consumers

Kafka provides the same per-key guarantee through partitions. The producer hashes the record key to a partition; within a partition, records have a total order by offset, and a consumer group assigns each partition to exactly one member. The producer settings are what keep that order intact across retries:

# producer.properties — preserve per-partition order under retries
acks=all
enable.idempotence=true                    # broker de-duplicates retried batches by sequence number
max.in.flight.requests.per.connection=5    # safe up to 5 only when idempotence is on
retries=2147483647
delivery.timeout.ms=120000
partitioner.class=org.apache.kafka.clients.producer.internals.DefaultPartitioner

Without idempotence, a retried batch can land after a later batch that succeeded first, silently swapping two records in the partition. The consumer side has its own trap — processing a poll batch concurrently across a thread pool reintroduces the parallel-consumer problem inside a single partition. The detailed setup, including how to fan out safely by key within a partition, is in per-entity ordering with Kafka partition keys.

RabbitMQ approaches ordering differently. A single queue delivers in FIFO order, but with multiple consumers and a prefetch above one, messages are spread across consumers and complete in any order. The x-single-active-consumer argument makes the broker deliver to one consumer at a time and fail over to the next registered consumer if it disconnects:

# rabbit_ordered.py — one queue per shard, one active consumer per queue
import pika

conn = pika.BlockingConnection(pika.ConnectionParameters("rabbitmq"))
ch = conn.channel()

SHARDS = 8
ch.exchange_declare("order-events", exchange_type="x-consistent-hash", durable=True)
for shard in range(SHARDS):
    q = f"order-events.shard-{shard}"
    ch.queue_declare(q, durable=True, arguments={
        "x-single-active-consumer": True,   # standby consumers take over on disconnect
        "x-queue-type": "quorum",           # replicated; SAC is supported on quorum queues
    })
    ch.queue_bind(q, "order-events", routing_key="1")   # equal weight per shard

ch.basic_qos(prefetch_count=1)   # >1 is fine for throughput only if handling stays sequential

The consistent-hash exchange spreads order ids across shards so that eight lanes run in parallel, and single active consumer guarantees each lane has one reader. Single active consumer in RabbitMQ covers failover behaviour and the prefetch subtleties.

Two shapes of the same lane On the left, a Kafka topic with partitions each assigned to one consumer group member. On the right, RabbitMQ shard queues each with one active consumer and a standby that takes over on disconnect. Both provide per-key order with parallelism across lanes. Partition owner vs active consumer Kafka consumer group partition 0 member A owns partition 1 member B owns RabbitMQ shard queues shard-0 (SAC) consumer X active shard-0 (SAC) consumer Y standby

Trade-off Analysis: Strict Order vs Throughput

Every ordering guarantee is paid for in throughput or availability. The table summarises the choices that come up most often:

Approach Ordering scope Throughput ceiling Behaviour when head fails Operational cost
Single queue, one worker Global One worker's rate Whole queue blocks Trivial, but a single point of slowness
SQS FIFO with message groups Per group 300 msg/s per API action, 3,000 with batching; far higher in high-throughput mode Group blocks until DLQ threshold, then skips Managed; 20,000 in-flight messages per queue
Kafka keyed partitions Per partition (hence per key) Scales with partition count Partition blocks unless you park the record Partition count hard to change later
RabbitMQ shards + SAC Per shard queue Scales with shard count Shard blocks; requeue loops possible You manage the shard topology
Unordered + version checks None at transport; enforced at write Unbounded by ordering Nothing blocks Every consumer must compare versions

The last row deserves more attention than it usually gets. For state projections — "the current status of order 17" — you often do not need the transport to preserve order at all. If each event carries a monotonically increasing version and the writer applies UPDATE ... WHERE version < :new_version, stale events become harmless no-ops and the queue can be fully parallel. Ordering at the transport is only mandatory when every intermediate step has a side effect that must happen in sequence, such as a ledger where each entry depends on the balance left by the previous one.

Failure Modes & Recovery

The stuck lane. The head message of a lane fails deterministically, and every later message for that key waits behind it. With per-key lanes the blast radius is one key — which is the point — but a hot key such as a large tenant can make that one key a significant fraction of traffic. Remediation: alert on lane-head age (the SQS ApproximateAgeOfOldestMessage metric is queue-wide, so track per-group age in the consumer), and give the lane a bounded retry count before parking the message to a dead-letter queue with a record of which key it blocked.

Reordering on rebalance. When a Kafka consumer group rebalances, a partition moves from one member to another. If the old member had processed but not committed records, the new owner reprocesses them — in order, but duplicated. If the old member was still processing when the partition was revoked, both can briefly process the same partition. Remediation: commit synchronously in the on_partitions_revoked callback, use the cooperative-sticky assignor to minimise movement, and keep handlers idempotent.

Hidden concurrency inside a consumer. A developer adds asyncio.gather or a thread pool to speed up a slow handler, and ordering silently disappears within each lane. Remediation: make the in-flight-per-lane limit an explicit, tested configuration value, and add a test that feeds a lane of sequenced messages with random handler delays and asserts completion order.

Retries that jump the queue. Framework-level retries in Celery, Sidekiq, and BullMQ re-enqueue a failed job with a delay, which puts it behind later jobs for the same key. The retry strategy that suits an unordered queue is actively harmful on an ordered one; retry in place (sleep and re-attempt inside the handler) or block the lane instead.

When the head of a lane fails If later steps depend on the failed one, block the lane and retry in place with a bound. If the step can be compensated or skipped, park it in a dead-letter queue and continue. If the consumer checks versions, reordering is harmless and the message can simply be retried later. Pick the failure policy per consumer lane head fails after N in-place retries later steps depend on it: block lane, page on-call step is skippable: park in DLQ, continue lane writes are versioned: requeue freely, order is irrelevant

Producer-Side Ordering: Getting the Sequence Right Before the Queue

Transport ordering only preserves the order in which messages arrive at the broker. If two application servers handle two requests for the same order a few milliseconds apart, the one whose publish call wins the network race is first — which may not be the one whose database transaction committed first. The queue then faithfully preserves the wrong order.

The reliable fix is to derive the sequence from the system of record rather than from publish timing. The transactional outbox pattern does this naturally: the business write and the outbox row commit in one transaction, and the outbox row carries a per-aggregate sequence number assigned inside that transaction. A relay publishes outbox rows in sequence order per aggregate, so the broker receives events in commit order even when requests raced.

-- outbox.sql — per-aggregate sequence assigned inside the business transaction
CREATE TABLE outbox (
  id           bigserial PRIMARY KEY,
  aggregate_id text        NOT NULL,
  seq          bigint      NOT NULL,
  event_type   text        NOT NULL,
  payload      jsonb       NOT NULL,
  published_at timestamptz,
  UNIQUE (aggregate_id, seq)              -- two writers cannot claim the same slot
);

-- Inside the same transaction as the order update:
INSERT INTO outbox (aggregate_id, seq, event_type, payload)
SELECT :order_id, COALESCE(MAX(seq), 0) + 1, 'OrderPaid', :payload
FROM outbox WHERE aggregate_id = :order_id;
-- The UNIQUE constraint turns a concurrent race into a retryable conflict
-- instead of two events with the same sequence number.

The relay then publishes with the aggregate id as the ordering key and the sequence number in the message body. Downstream consumers get two protections at once: the transport keeps events for one aggregate in order, and the explicit sequence lets them detect a gap if a message was ever parked or lost. A consumer that trusts the sequence number rather than arrival order also survives a future broker migration, which is the moment many implicit ordering assumptions break.

Sequence assigned at commit, not at publish The business update and an outbox row with a per-order sequence number commit in one database transaction. A relay reads unpublished outbox rows in sequence order and publishes them with the order id as the ordering key, so the queue receives events in commit order. Commit order becomes queue order one transaction UPDATE orders INSERT outbox seq=3 relay ORDER BY seq per order ordered lane key = order_id, body has seq Two racing API servers cannot swap events: the UNIQUE (aggregate_id, seq) constraint serialises them.

Performance Tuning

The throughput of an ordered system is governed by three numbers: the number of active lanes, the service time of the head message, and how evenly keys spread across lanes.

  • Lane count. For Kafka, choose partitions for the peak parallelism you will need in two years, because increasing the count remaps keys and breaks ordering for in-flight keys during the change. Sixty-four or 128 partitions is a common starting point for a busy event topic; rebalancing consumer groups gets slower as the count rises, so do not go to thousands without reason.
  • Head service time. Because a lane processes one message at a time, per-lane throughput is 1 / service_time. A handler that takes 50 ms caps each lane at 20 messages per second no matter how much hardware you add. Move slow, order-independent side effects (sending an email, calling an analytics API) onto a separate unordered queue triggered by the ordered handler.
  • Key skew. Measure the distribution of messages per key. If the top 1% of keys carry 40% of traffic, those keys set your tail latency. Options include splitting a hot key into sub-keys where business rules allow (per-account ordering instead of per-tenant), or giving hot tenants a dedicated partition.
  • Batching within a lane. Pulling ten messages from one SQS group and processing them sequentially in one receive cycle amortises the API round trip without breaking order; it is the single cheapest throughput win for FIFO queues.
# Per-lane stall detection: age of the oldest unprocessed message per partition
max by (topic, partition) (
  kafka_consumergroup_lag_seconds{consumergroup="order-projector"}
) > 300

# Skew: share of traffic carried by the busiest partition
max(rate(kafka_topic_partition_current_offset{topic="order-events"}[10m]))
  / sum(rate(kafka_topic_partition_current_offset{topic="order-events"}[10m]))

A busiest-partition share far above 1 / partition_count means the key choice, not the hardware, is the bottleneck.

Testing That Ordering Actually Holds

Ordering bugs hide in development because a single local worker processes everything sequentially. A test that proves ordering must introduce the conditions that break it: several consumers, random handler latency, injected failures, and redelivery. The harness below runs against a real broker in CI (a containerised RabbitMQ, LocalStack SQS, or a single-node Kafka) and asserts that every key's events were applied in sequence.

# test_ordering.py — property-style check: per-key order survives chaos
import random
import threading
from collections import defaultdict

def test_per_key_order_under_parallel_consumers(queue, start_consumers):
    applied: dict[str, list[int]] = defaultdict(list)
    lock = threading.Lock()

    def handler(msg):
        time.sleep(random.uniform(0, 0.05))          # random service time
        if random.random() < 0.05:
            raise RuntimeError("injected transient failure")   # forces redelivery
        with lock:
            applied[msg["key"]].append(msg["seq"])

    for key in (f"order-{n}" for n in range(200)):
        for seq in range(1, 21):
            queue.publish(key=key, body={"key": key, "seq": seq})

    start_consumers(handler, count=8)
    queue.wait_until_drained(timeout=120)

    for key, seqs in applied.items():
        deduped = [s for i, s in enumerate(seqs) if i == 0 or s != seqs[i - 1]]
        assert deduped == sorted(deduped), f"{key} applied out of order: {seqs}"

The deduplication step matters: under at-least-once delivery a retried message can be applied twice in a row, which is a separate property (idempotency) from ordering. Run the test with the production failure policy switched on, and run it again after any change to consumer concurrency settings — those changes are where ordering most often regresses.

FAQ

Does a FIFO queue guarantee my jobs run in order? Only if one consumer processes one message at a time for each ordering key. A FIFO queue guarantees delivery order; with parallel workers, retries, or batches processed concurrently, completion order can differ. SQS FIFO enforces the one-in-flight-per-group rule for you; most other queues leave it to the consumer.

Should I use global ordering or per-key ordering? Per-key ordering almost always. Global ordering limits the whole system to a single consumer's throughput and makes one failing message block everything. Identify the entity whose history must be sequential — an order, an account, a document — and use its id as the ordering key.

Can I increase Kafka partitions without breaking ordering? Not transparently. Adding partitions changes which partition a key hashes to, so for a short period new records for a key go to a new partition while older ones are still being consumed from the old one. Drain or pause producers for affected keys, or over-provision partitions up front.

How do I keep order when a message must be retried? Retry in place inside the handler (with a bounded number of attempts and backoff) rather than re-enqueueing, or leave the message unacknowledged so the broker redelivers it at the head of the lane. Re-enqueueing with a delay moves the message behind later ones for the same key.

Related