skip to content

A team's outbox relay publishes events out of order across different aggregates, and occasionally a consumer receives the same OrderPlaced event twice. What guarantees does the transactional outbox pattern actually provide here, and what must consumers/producers do to handle it correctly?

level: seniorimportance: must knowfreq 55%

answer

  1. at-least-once, never exactly-once
  2. idempotent consumers, dedupe by event ID
  3. ordering only guaranteed per-aggregate via partition key
  4. concurrent relays break ordering unless partitioned consistently
  5. poller lag under load = growing backlog, not lost events

basics

~20 s

The outbox pattern guarantees an event eventually gets published if the database write succeeded, but not that it's published exactly once or perfectly ordered across everything - so consumers must be built to safely handle duplicate or occasionally reordered messages.

solid answer

~40 s

Transactional outbox gives at-least-once delivery, not exactly-once: a relay can crash after publishing but before marking a row 'sent,' causing a duplicate on retry, so consumers must be idempotent (dedupe by event id or use upsert-style handlers). Ordering is only reliably preserved within a single aggregate's stream if the relay processes and publishes rows in insertion order and routes by a partition key (e.g., order id) that keeps that aggregate's events on one partition; across different aggregates, or with a concurrent multi-threaded relay, global ordering is not guaranteed and usually isn't needed. Consumers requiring strict ordering must key their processing (or rely on the broker's partitioning) by the same aggregate id the producer used.

go deeper

for a junior

Should understand in plain terms that the same event might sometimes arrive twice and that's expected, not a bug to panic over.

for a middle

Should state 'at-least-once' explicitly and know consumers need some form of deduplication.

for a senior

Should explain per-aggregate ordering via partition keys, and how concurrent relay instances can break ordering if not partitioned consistently.

for a principal

Should discuss monitoring/alerting for relay lag and backlog growth, capacity planning for the outbox table, and how to design the event schema (event ids, aggregate ids) up front to make idempotency and ordering tractable at scale.

## One strong promise, several weaker ones The transactional outbox pattern makes one strong promise and several weaker ones, and conflating them is the most common source of production incidents with this pattern. - **The strong promise:** if the business transaction committed, the corresponding outbox row is durably recorded and will eventually be published - no event is silently lost because the state change happened. - **What the pattern does not promise, on its own:** exactly-once delivery or global ordering. Both of those gaps have specific, well-understood causes. ## Duplication Start with duplication. The relay's job is: 1. read an unpublished row, 2. publish it to the broker, 3. mark it published. Those are two separate operations (publish, then mark) with no shared transaction between the broker and the outbox table's status update - the same fundamental atomicity gap that motivated the outbox pattern in the first place, just one level further down the pipeline. If the relay crashes, or the network drops, between a successful publish and the status update committing, the restarted relay sees the row as still unpublished and republishes it. This is why the pattern is correctly described as **at-least-once** rather than exactly-once: duplicates are not a bug, they are an accepted, structural consequence of prioritizing 'never lose an event' over 'never duplicate an event,' since guaranteeing the latter would require its own distributed transaction. The fix lives entirely on the consumer side. Every published event needs a stable, unique identifier (the outbox row's primary key, or a UUID assigned at insert time), and consumers must either - maintain a dedupe store (a processed-event-ids table with a unique constraint, checked before applying an event), or - make the downstream effect naturally idempotent - for instance, an upsert keyed by order id rather than an append-only insert, so processing the same event twice produces the same end state as processing it once. ## Ordering Ordering has a different, more subtle cause. A single-threaded relay that reads rows strictly in primary-key (insertion) order and publishes them one at a time preserves global insertion order by construction - row 5 is always sent before row 6. But that guarantee evaporates the moment the relay is scaled out for throughput: if three relay instances each poll and publish a disjoint slice of rows concurrently, their publishes can interleave in the broker in essentially any order, because nothing coordinates the timing between three independent processes. - **Across different aggregates, it usually doesn't matter.** A consumer rarely cares whether `OrderPlaced` for order #500 was published before or after `ShipmentCreated` for order #900, since they're unrelated. - **Within a single aggregate's stream, it matters a great deal.** If `OrderPlaced` and `OrderCancelled` for the same order #500 arrive out of order, a naive consumer could process the cancellation and then 'un-cancel' it by processing the placement afterward. The standard fix is to route by a partition key equal to the aggregate id (e.g., Kafka's producer partition key set to `order_id`), which guarantees that Kafka preserves order within that partition even though different aggregates' events land on different partitions and have no ordering relationship to each other. A relay that wants to scale out safely therefore either shards its polling by aggregate id (so all of one aggregate's rows are always handled by the same relay instance) or relies on the broker's partitioning to reimpose per-aggregate order downstream, rather than trying to enforce a single global publish order. ## Relay lag under load The third failure mode worth naming explicitly is relay lag under load: the backlog of unpublished outbox rows growing because the relay can't keep up with the insert rate. This is not data loss - the rows are safely sitting in the table - but it does mean the gap between 'business transaction committed' and 'downstream systems find out' widens, which can violate freshness expectations for things like a search index or a cache that assumes near-real-time updates. Production systems typically monitor an 'oldest unpublished row age' metric specifically to catch this before it becomes user-visible staleness, and treat a growing backlog as a capacity or relay-scaling signal rather than a correctness bug. ## Where it shows up A concrete example of these guarantees in practice: Debezium's own outbox event router documentation is explicit that consumers must handle duplicate delivery, and recommends including an aggregate-id-derived key on the Kafka message specifically so Kafka's per-partition ordering keeps each aggregate's events in commit order even under a scaled-out capture pipeline - it does not claim, and explicitly disclaims, exactly-once or global ordering as a property of the pattern itself.

  • How would you make a consumer of these events idempotent in practice?
    Attach a stable, unique event id (often the outbox row's primary key or a UUID generated at insert time) to every published message, and have the consumer record processed ids in its own store - either via a dedupe table with a unique constraint, or by making the downstream write itself naturally idempotent (e.g., an upsert keyed by order id rather than an insert-only append).
  • How do you preserve per-aggregate ordering when scaling the relay to multiple concurrent instances?
    Partition work by aggregate id - either by sharding which rows each relay instance polls (e.g., hash of order id mod N instances) or by publishing to a Kafka topic partitioned by order id, so all events for a given order land on the same partition and are consumed in commit order. A single unpartitioned relay naturally preserves global insertion order but doesn't scale horizontally without care.
  • What does it look like in production when the relay falls behind (poller lag)?
    The outbox table's backlog of unpublished rows grows, and the time between a business transaction committing and its event reaching consumers increases - visible as a widening gap on an 'oldest unpublished row age' metric. It's not data loss, but it can violate downstream freshness expectations (e.g., a search index or cache staying stale) and, if unmonitored, can eventually exhaust storage or slow the poll query itself as the table grows.

It's like a mail-sorting facility that guarantees every letter you drop in the same mailbox eventually gets delivered and that letters to the same address arrive in the order they were dropped - but it makes no promise about the relative order of letters going to different addresses, and if a truck breaks down mid-route, some letters get redelivered rather than lost.

saying these in an interview costs you the question

  • Claims the pattern gives exactly-once delivery
  • Assumes global cross-aggregate ordering is guaranteed
  • Has no answer for how consumers detect/handle a duplicate event
  • Doesn't connect relay concurrency to ordering risk

context