Handling Out-of-Order Events with Sequence Numbers

Transport ordering is fragile, so this guide takes the other route described in Message Ordering Guarantees, part of Queue Fundamentals & Architecture: stamp each event with a per-entity sequence number and make the consumer correct no matter what order events arrive in. The technique lets you keep a fully parallel, unordered queue and still produce a correct result.

Problem Statement

A subscription service publishes PlanChanged events to a standard queue consumed by twelve workers. A customer who upgrades and then quickly downgrades sometimes ends up on the upgraded plan, because the downgrade event completed first and the upgrade overwrote it. Switching to an ordered transport would cap throughput and has its own failure modes during rebalances and dead-lettering. You want the consumer to reach the correct final state regardless of delivery order, to never apply the same event twice, and to notice when an event is missing entirely.

Prerequisites

  • The ability to add a field to the event payload at the producer, ideally assigned inside the same transaction as the state change (an outbox or a version column).
  • A consumer datastore that supports conditional writes: UPDATE ... WHERE, INSERT ... ON CONFLICT, DynamoDB condition expressions, or Redis Lua.
  • Agreement on what kind of consumer this is — a state consumer (only the latest value matters) or a step consumer (every event has a side effect). The two need different handling below.

Step 1 — Assign the Sequence at the Source of Truth

The sequence must reflect the order in which changes were committed, not the order in which they were published. The simplest source is a version column on the entity, incremented by the same UPDATE that changes the state.

-- Increment the version in the same statement that changes the plan
UPDATE subscriptions
SET plan = :new_plan,
    version = version + 1,
    updated_at = now()
WHERE id = :subscription_id
RETURNING version;           -- publish this value as the event's seq
# producer.py — the event carries entity id + version
def change_plan(sub_id: str, new_plan: str) -> None:
    with db.transaction() as tx:
        version = tx.scalar(UPDATE_SQL, subscription_id=sub_id, new_plan=new_plan)
        tx.execute(INSERT_OUTBOX, payload={
            "type": "PlanChanged",
            "subscription_id": sub_id,
            "plan": new_plan,
            "seq": version,              # strictly increasing per subscription
            "event_id": f"{sub_id}:{version}",
        })

Because the database serialises updates to one row, two concurrent plan changes receive versions 7 and 8 in commit order even if their publish calls race. Wall-clock timestamps are not a substitute: clocks on different hosts disagree by milliseconds, and two events in the same millisecond have no order at all.

Arrival order does not matter to a version check A consumer stores version 6 for a subscription. Event 8 arrives first and is applied, moving the stored version to 8. Event 7 arrives late and is ignored because 7 is not greater than 8. Event 9 arrives and is applied. The final state equals the state after event 9. Arrivals: 8, 7, 9. Final state: 9. stored v6 plan = basic seq 8 arrives 8 > 6: applied seq 7 arrives 7 <= 8: no-op seq 9 arrives 9 > 8: applied Works for state consumers, where event 8 already contains everything event 7 would have set.

Step 2 — State Consumers: Apply Only Newer Versions

If each event carries the full resulting state of the fields it touches, an older event can simply be ignored once a newer one has been applied. The whole check is one conditional write:

# state_consumer.py — last-writer-wins by version, atomically
def handle_plan_changed(ev: dict) -> None:
    rows = db.execute(
        """
        INSERT INTO billing_view (subscription_id, plan, version)
        VALUES (:sid, :plan, :seq)
        ON CONFLICT (subscription_id) DO UPDATE
          SET plan = EXCLUDED.plan, version = EXCLUDED.version
          WHERE billing_view.version < EXCLUDED.version    -- stale events are no-ops
        """,
        sid=ev["subscription_id"], plan=ev["plan"], seq=ev["seq"],
    ).rowcount
    if rows == 0:
        metrics.stale_events.inc()      # useful signal: how often reordering happens

This also makes the consumer idempotent: a duplicate of event 8 fails the version < 8 check and changes nothing. It is the cheapest correct design for projections, caches, search indexes, and any read model, and it is why many teams never need an ordered transport. The broader duplicate-handling picture is in preventing duplicate job execution with idempotency.

Step 3 — Step Consumers: Buffer Early Arrivals

Some consumers cannot skip an event. A ledger, an email sequence, or a workflow step must process 7 before 8 because each has its own side effect. For these, an event that arrives ahead of its predecessor is parked until the gap fills.

# step_consumer.py — hold early arrivals, drain in order
import json

PARK_KEY = "parked:{sid}"           # Redis sorted set of early events, scored by seq

def handle_step(ev: dict) -> None:
    sid, seq = ev["subscription_id"], ev["seq"]
    last = store.last_applied(sid)                       # from the consumer's own DB
    if seq <= last:
        return                                           # duplicate or already covered
    if seq > last + 1:
        redis.zadd(PARK_KEY.format(sid=sid), {json.dumps(ev): seq})
        redis.expire(PARK_KEY.format(sid=sid), 86400)
        return                                           # ack: it is safely parked
    apply_step(ev)                                       # side effect + last_applied=seq, atomically
    drain_parked(sid, seq)

def drain_parked(sid: str, last: int) -> None:
    key = PARK_KEY.format(sid=sid)
    while True:
        nxt = redis.zrangebyscore(key, last + 1, last + 1)
        if not nxt:
            return
        ev = json.loads(nxt[0])
        apply_step(ev)
        redis.zrem(key, nxt[0])
        last += 1

Acknowledging a parked event is safe only because the park is durable. If the parking store can lose data, nack the early event instead and let the broker redeliver it later — slower, but nothing is lost. Two consumers handling events for the same subscription can race in drain_parked; guard the apply with a per-entity lock or a conditional update on last_applied so exactly one wins each sequence.

Park early, drain when the gap fills The consumer has applied up to sequence 6. Event 8 arrives and is parked because 7 is missing. When event 7 arrives it is applied, and the drain step then finds 8 in the parking buffer and applies it, leaving the consumer at sequence 8 with every step executed in order. Every step runs, and runs in order seq 8 arrives parked: 8 (waiting for 7) seq 7 arrives apply 7, then drain 8 last_applied = 8

Step 4 — Put a Deadline on Gaps

A parked event whose predecessor never arrives would wait forever. The predecessor may be in a dead-letter queue, lost to a producer bug, or simply slow. Give gaps a deadline and escalate when it passes.

# gap_sweeper.py — run every minute
GAP_DEADLINE_SECONDS = 300

def sweep() -> None:
    for key in redis.scan_iter("parked:*"):
        sid = key.decode().split(":", 1)[1]
        oldest = redis.zrange(key, 0, 0, withscores=True)
        if not oldest:
            continue
        ev = json.loads(oldest[0][0])
        age = time.time() - ev["published_at"]
        if age > GAP_DEADLINE_SECONDS:
            missing = store.last_applied(sid) + 1
            alert("sequence gap", subscription=sid, missing_seq=missing, waiting=len_parked(key))
            request_republish(sid, from_seq=missing)     # ask the producer's outbox to resend

request_republish works when the producer keeps its outbox rows for a while: the consumer can ask for "subscription 42 from sequence 7" and the relay re-sends them. That turns a missing message from a silent data bug into a self-healing delay.

Step 5 — Instrument Reordering, Parking, and Gaps

Once consumers tolerate reordering, reordering stops being visible as a bug — which is good for users and bad for operators, because a slowly growing problem (a producer that has started losing events, a consumer lagging on one shard) no longer shows up anywhere unless you measure it. Three counters and one gauge cover it:

# metrics.py — prometheus_client
from prometheus_client import Counter, Gauge, Histogram

STALE = Counter("events_stale_total", "Events ignored because a newer version was applied",
                ["event_type"])
PARKED = Counter("events_parked_total", "Early events parked waiting for a predecessor",
                 ["event_type"])
GAP_ESCALATIONS = Counter("event_gaps_escalated_total", "Gaps that exceeded the deadline",
                          ["event_type"])
PARKED_NOW = Gauge("events_parked_current", "Events currently parked across all entities")
REORDER_DISTANCE = Histogram("event_reorder_distance", "seq - last_applied for early arrivals",
                             buckets=(1, 2, 3, 5, 10, 25, 100))

events_stale_total divided by events consumed gives the reorder rate of your transport; on a busy parallel queue a few percent is normal and harmless. event_reorder_distance tells you how far ahead early events tend to be — mostly 1 means ordinary races between workers; large values mean something upstream is holding events back. events_parked_current is the one to alert on: it should hover near zero and drain on its own. A value that climbs steadily means gaps are opening faster than they close, and event_gaps_escalated_total tells you when the deadline in Step 4 has started firing.

# Alert: parked events are accumulating rather than draining
deriv(events_parked_current[15m]) > 0 and events_parked_current > 500

# Dashboard: share of events that arrived stale
sum(rate(events_stale_total[5m])) / sum(rate(events_consumed_total[5m]))
Healthy parking drains; a stuck gap climbs The first trace rises briefly after a burst and falls back to near zero within a minute as predecessors arrive. The second trace climbs steadily because one missing event holds back every later event for its entity, which is the signal to alert on. events_parked_current over 30 minutes bursts that drain: normal steady climb: a gap is not closing 0 min 30 min

Verification

Write a test that shuffles a sequence of events, delivers them with duplicates, and checks both consumer types:

import random

def test_state_consumer_converges_under_shuffle():
    events = [{"subscription_id": "s1", "plan": f"p{i}", "seq": i} for i in range(1, 51)]
    delivered = events + random.sample(events, 10)       # 10 duplicates
    random.shuffle(delivered)
    for ev in delivered:
        handle_plan_changed(ev)
    assert db.scalar("select plan from billing_view where subscription_id='s1'") == "p50"

def test_step_consumer_applies_each_step_once_in_order():
    events = [{"subscription_id": "s2", "seq": i} for i in range(1, 31)]
    random.shuffle(events)
    for ev in events:
        handle_step(ev)
    assert applied_steps("s2") == list(range(1, 31))

In production, graph stale_events_total and the number of parked events. A steady low rate of stale events is normal on a parallel queue; a growing parked count means a gap is not closing.

Gotchas & Edge Cases

Partial-state events. Last-writer-wins only works when each event carries the full value of every field it affects. A PlanChanged event that sends only a delta ("add 5 seats") is a step event; ignoring an older one loses data.

Sequences that restart. Deleting and recreating an entity with the same id resets its version to 1, and every new event looks stale. Include a generation or creation timestamp in the comparison, or never reuse ids.

Sequences shared across entity types. One global counter across all entities is not a per-entity sequence: gaps are normal (other entities consumed the numbers), so gap detection cannot work. Use a counter per entity.

Memory in the parking store. A consumer outage upstream can leave millions of parked events. Put a TTL on parked keys and a size cap per entity, and alert well before either is reached.

FAQ

Should I use sequence numbers or an ordered queue? For state projections, sequence numbers with conditional writes are simpler and faster. Use an ordered transport when every event has a side effect that must run in sequence and parking would be complicated — and even then keep the sequence number for gap detection.

Can I use timestamps instead of sequence numbers? Only if a single writer produces them from one monotonic clock. Across multiple hosts, clock skew and same-millisecond events make timestamps unreliable for ordering; a version incremented by the database is exact.

What if events come from several independent producers? Each producer can only sequence what it owns. If two services both change a subscription, route changes through one owner, or have each event carry a vector of per-source versions and merge by field.

Related