Exactly-Once Is a Lie: Engineering Reliable Event-Driven Systems
{"prompt":" \"modern distributed systems engineering control room, sleek server infrastructure with glowing racks in background | large transparent holographic display showing 'Exactly-Once Is a Lie' in bold modern typography, software engineers in discussion pointing at event flow diagrams, floating nodes and message queues visualizing event-driven architecture, subtle duplicate event packet icons overlapping | text elements 'Exactly-Once Is a Lie' rendered in clean sans-serif font, naturally embedded into the holographic interface, clear and highly readable ::7 | cinematic dramatic lighting with cool blue and teal accents, glowing screen illumination, dark professional atmosphere, depth of field blur on background racks ::7 | 8k resolution, hyperrealistic, photorealistic quality, octane render, cinematic composition, sharp focus, high detail, professional photography --ar 16:9 --s 1000 --q 2 --v 5.2\",","originalPrompt":" \"modern distributed systems engineering control room, sleek server infrastructure with glowing racks in background | large transparent holographic display showing 'Exactly-Once Is a Lie' in bold modern typography, software engineers in discussion pointing at event flow diagrams, floating nodes and message queues visualizing event-driven architecture, subtle duplicate event packet icons overlapping | text elements 'Exactly-Once Is a Lie' rendered in clean sans-serif font, naturally embedded into the holographic interface, clear and highly readable ::7 | cinematic dramatic lighting with cool blue and teal accents, glowing screen illumination, dark professional atmosphere, depth of field blur on background racks ::7 | 8k resolution, hyperrealistic, photorealistic quality, octane render, cinematic composition, sharp focus, high detail, professional photography --ar 16:9 --s 1000 --q 2 --v 5.2\",","width":1061,"height":555,"seed":42,"model":"sana","enhance":false,"nologo":true,"negative_prompt":"undefined","nofeed":false,"safe":false,"quality":"medium","image":[],"transparent":false,"isMature":false,"isChild":false,"trackingData":{"actualModel":"sana","usage":{"completionImageTokens":1,"totalTokenCount":1}}}

Exactly-Once Is a Lie: Engineering Reliable Event-Driven Systems

Exactly-Once Is a Lie: Engineering Reliable Event-Driven Systems

Every distributed messaging platform eventually makes the same promise: exactly-once delivery. Kafka sells transactions. SQS offers FIFO queues with deduplication. Pub/Sub ships exactly-once subscriptions. The claims are technically defensible and operationally misleading. The moment a message crosses a process boundary, the producer cannot know whether the consumer processed it — a crash between processing and acknowledgement is indistinguishable from a crash before processing ever started. The only real question is who owns the job of making duplication harmless.

This article covers the engineering that actually makes event-driven systems dependable: the dual-write problem, the transactional outbox, idempotent consumers, ordering and partitioning, sagas, schema evolution, poison-message handling, and the telemetry that tells you when any of it has broken.

Delivery semantics: what is actually on the table

Three semantics exist, and only two are achievable without heroics:

  • At-most-once — fire and forget. The producer publishes and never retries. Messages can be lost, but never duplicated. Acceptable for metrics, clickstream, and telemetry where a lost sample is noise.
  • At-least-once — the producer or broker retries until an acknowledgement arrives. This is the default for almost every real system. Duplicates are guaranteed, not exceptional.
  • Effectively-once — at-least-once transport combined with idempotent consumers, so that duplicates produce no additional observable effect. This is what people mean when they say exactly-once, and it is the only version that survives contact with production.

True exactly-once delivery across arbitrary services is not achievable, because it requires atomic commit across two systems that do not share a transaction coordinator. What is achievable is exactly-once effect: the transport may deliver a message three times, but the business outcome happens once. Engineering effort belongs there, not in chasing a delivery guarantee that cannot exist.

The dual-write problem

Consider the most common event-driven code path in existence: a service receives a request, writes to its database, then publishes an event to a broker. Two independent systems, two independent commitments, no shared transaction. Every failure window is a data-consistency bug waiting to be discovered by an accountant.

  • Commit succeeds, publish fails — the order exists in the database, but no downstream service ever learns about it. Inventory is never reserved, the confirmation email never sends, and the read model silently diverges.
  • Publish succeeds, commit fails — downstream services act on an order that does not exist. This is worse: it is a phantom event that no reconciliation job can easily unwind.
  • Both succeed, client retries — duplicate order, duplicate charge, duplicate shipment.
  • Process dies between the two — nondeterministic, environment-dependent, and typically invisible until volume grows enough for the race to hit.

Wrapping both operations in a try/catch and hoping is not a strategy. Neither is publishing first and writing second — you have simply moved the window. The fix is to make the state change and the intent to publish a single atomic unit.

Pattern 1: the transactional outbox

Instead of writing to the database and then publishing to the broker, write both the domain state and the outgoing event into the same database, inside the same transaction. A separate relay process reads the outbox table and publishes to the broker. Now the atomicity problem is delegated to the one system that already solves it: the relational database.

BEGIN;

INSERT INTO orders (id, customer_id, total_cents, status)
VALUES (:id, :customer_id, :total, 'PLACED');

INSERT INTO outbox (id, aggregate_id, aggregate_type, event_type, payload, created_at)
VALUES (gen_random_uuid(), :id, 'Order', 'OrderPlaced', :payload_json, now());

COMMIT;

The relay can be implemented two ways.

Polling publisher. A worker polls the outbox table, publishes unpublished rows, and marks them published. Use FOR UPDATE SKIP LOCKED so multiple relay instances can run safely in parallel without coordinating.

SELECT id, aggregate_id, event_type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;

Change data capture. A CDC connector such as Debezium tails the database write-ahead log and streams outbox inserts straight to the broker. This removes polling latency and relay infrastructure, at the cost of coupling your deployment to log format versions, slot management, and a new component that must itself be monitored. CDC also surfaces a subtlety people forget: WAL ordering is per-table, and transaction boundaries must be preserved if you need all-or-nothing publication across several events.

Either way, two things become true that were not true before. First, the event cannot be lost — it is durably stored with the state change that caused it. Second, the relay is at-least-once by construction, so duplicates are now a consumer problem rather than a producer problem. Publish lag becomes a first-class metric: outbox_lag_seconds is the age of the oldest unpublished row, and it should page someone when it climbs.

Retention matters too. Prune published rows on a schedule tied to your maximum plausible replay window. Keeping an outbox table forever turns a small coordination table into your largest table.

Pattern 2: idempotent consumers

Once you accept at-least-once transport, every consumer must be idempotent: processing the same message twice must produce the same observable state as processing it once. There are several ways to get there, and they have different costs.

Natural idempotency. Some operations are inherently idempotent. UPDATE orders SET status = 'SHIPPED' WHERE id = :id can run a thousand times with the same result. Prefer these when the domain allows it — they need no extra storage and no coordination. The catch is that anything with a side effect outside the database (sending email, charging a card, calling a partner API) is not naturally idempotent.

Deduplication tables. Record the fact that a message was processed, in the same transaction as the side effect. The unique constraint does the hard work.

CREATE TABLE processed_messages (
  consumer_group TEXT        NOT NULL,
  message_id     UUID        NOT NULL,
  processed_at   TIMESTAMPTZ NOT NULL DEFAULT now(),
  PRIMARY KEY (consumer_group, message_id)
);

The consumer logic then becomes a single atomic unit:

BEGIN;

INSERT INTO processed_messages (consumer_group, message_id)
VALUES (:group, :msg_id)
ON CONFLICT DO NOTHING;

-- if zero rows were inserted, this is a duplicate: skip the work

UPDATE inventory SET reserved = reserved + :qty WHERE sku = :sku;

INSERT INTO order_lines (order_id, sku, qty) VALUES (:order_id, :sku, :qty);

COMMIT;

The critical rule: the dedupe insert and the business write must be in the same transaction. If they are separate, you have recreated the dual-write problem with extra steps.

The critical caveat: this table grows without bound. Real deployments partition it by time and drop partitions older than the maximum possible redelivery window, which is derived from broker retention plus your retry policy — not from optimism. If a message can be redelivered after your dedupe window closes, it will eventually be, at 3 a.m., during an incident.

Idempotency keys at the API edge. The same idea applies to synchronous HTTP. Let the client supply an idempotency key; store the first response against that key; replay the stored response for subsequent requests with the same key. Stripe popularized this pattern for payments, and it is the correct answer for any endpoint that creates a resource with a side effect.

Version-based conditional writes. Where the domain has a natural version or sequence number, use optimistic concurrency: apply the update only if the incoming version is greater than the stored one. This also gives you a sane answer for out-of-order delivery, which dedupe tables alone do not.

Ordering, partitioning, and the hot-partition tax

Kafka and most log-based brokers guarantee order within a partition, not across a topic. Partition by the aggregate identifier — order id, account id, device id — and you get per-entity ordering for free. That is almost always the ordering you actually need.

What you get along with it:

  • Key skew. A single large tenant can dominate one partition while the rest idle. Monitor per-partition throughput, not just topic-level throughput.
  • Rebalance pauses. When a consumer leaves, partitions move. Use cooperative rebalancing and static membership to minimize the disruption, and always commit offsets after processing completes — never before, or a crash silently drops messages.
  • A scalability ceiling. Global ordering requires a single partition, which means a single consumer, which means your throughput ceiling is one machine. If you find yourself needing global order, question it before you build it.

On the producer side, enable idempotent production and keep in-flight requests bounded so retries cannot reorder writes within a partition. If a single logical operation must update two topics atomically, use a transactional producer — otherwise downstream consumers will occasionally see one topic without the other.

Pattern 3: sagas for long-running business transactions

A business transaction that spans several services cannot use a database transaction. The saga pattern replaces atomicity with a sequence of local transactions, each paired with a compensating action that undoes its effect if a later step fails. Compensation is not rollback: a shipped package is not un-shipped, it is recalled. A charged card is not un-charged, it is refunded.

Two coordination styles dominate:

  • Choreography. Each service listens for events and reacts, emitting its own events. Simple for three steps, unmaintainable at ten. Nobody can point to a single place that describes the workflow, and debugging means reconstructing the graph from logs.
  • Orchestration. A coordinator drives the steps and persists the saga state machine. Explicit, observable, testable. More infrastructure, but you can answer the question “where is this order stuck?” with a query rather than an afternoon.

Three design rules separate working sagas from cascading incidents:

  1. Every step needs a compensator, and compensators must be idempotent. A compensation may itself be retried, delayed, or delivered twice.
  2. Use semantic locks to prevent dirty reads. Write intermediate states such as PENDING or RESERVED so other services can see that a transaction is in flight, rather than seeing an inconsistent half-state.
  3. Identify the pivot transaction. After the pivot, the saga can only move forward — retrying until success. Before it, everything is compensable. Drawing this line explicitly prevents the classic failure of compensating a step that other services have already committed against.

Use a durable workflow engine for anything with a timeout longer than a few seconds. Hand-rolled saga state machines built on cron jobs and a status column work until they don’t, and they fail in ways that require manual database surgery at the worst possible moment.

Schema evolution: events are contracts you cannot retract

Once an event is published, it is in someone else’s database, and you cannot redeploy that. Treat every event schema as a public API with a permanent deprecation path.

  • Use a schema registry with Avro, Protobuf, or JSON Schema. Reject incompatible changes at publish time rather than at consume time in a service you cannot see.
  • Only add optional fields with defaults. Never repurpose a field, never change its type, never rename it in place. Add the new field, populate both during a migration window, then remove the old one after every consumer has migrated.
  • Events are facts, not snapshots. Publish past-tense domain events (OrderPlaced, PaymentCaptured) rather than mutable record dumps. A fact that was true at time T stays true forever; a snapshot becomes a lie the moment the row changes.
  • Separate events from commands. A command can be rejected and carries an expectation; an event is a record of something that already happened. Mixing them produces consumers that think they can veto history.

Poison messages, retries, and dead letter queues

Some messages will always fail. The goal is not zero failures but bounded, observable, recoverable failure.

Retry with exponential backoff and jitter, but do not sleep inside the consumer. Blocking a handler stalls the entire partition and turns one bad message into a topic-wide outage. Instead, route failed messages to a retry topic with an increasing delay, or use broker-native retry mechanisms that republish without blocking.

Classify failures before retrying them:

  • Transient — network timeouts, deadlocks, rate limits, downstream 5xx. Retry.
  • Permanent — schema validation errors, missing required fields, business rule violations. Retrying a permanent failure a thousand times is a self-inflicted denial of service. Send it straight to the dead letter queue.

A dead letter queue is not a graveyard. It needs a named owning team, an alert when its depth grows, a dashboard showing failure reasons clustered by type, and a documented replay procedure. Because consumers are idempotent, replaying a DLQ after fixing the bug is safe by construction — which is the entire point of the previous sections.

Observability across asynchronous boundaries

Distributed tracing handles synchronous calls well and asynchronous messaging poorly unless you deliberately propagate context. Inject trace context and correlation identifiers into message headers at publish time, and extract them in the consumer to create a link between spans rather than a false parent-child relationship.

Metrics that actually predict incidents:

  • Consumer lag per partition, not per consumer group. Averages hide the partition that is 40 minutes behind.
  • Outbox lag seconds — the age of the oldest unpublished event.
  • Duplicate suppression ratio — the percentage of messages rejected by the dedupe table. A sudden drop to zero is an alarm, not an improvement: it usually means a constraint was dropped, a retention window shrank, or consumers started running without the dedupe path.
  • Retry rate and DLQ arrival rate, sliced by exception type.
  • End-to-end processing latency from event creation timestamp, not from consumer receipt. The interesting number is how stale the read model is, not how fast the handler runs.

Log the message id, aggregate id, and dedupe key on every consumer invocation. When something goes wrong at 2 a.m., being able to grep one identifier and reconstruct the entire life of a business event is worth more than any dashboard.

Testing for the failures you cannot reproduce

Duplication and reordering bugs rarely appear in development because they require timing. Test them deliberately:

  • Contract tests between producer schemas and consumer expectations, run in CI against the registry.
  • Integration tests against real infrastructure using Testcontainers or ephemeral clusters. Mocks do not model consumer group rebalancing, and rebalancing is where the bugs live.
  • Duplication injection. Publish every message twice in a test environment and assert the final state is identical. This single test catches most idempotency defects before they reach production.
  • Chaos experiments. Kill the consumer between the business write and the offset commit, then verify the retry produces no double effect. Kill the broker mid-publish and verify the outbox relay replays correctly.
  • Reconciliation jobs as a safety net. Periodically compare derived read models against the source of truth and repair drift. Every mature event-driven system has one, because every mature event-driven system has had a divergence nobody predicted.

The operational checklist

Before you call an event-driven system production-ready, confirm each of these:

  • No service writes to a database and publishes to a broker outside a transaction.
  • Every consumer is idempotent, and its dedupe write shares a transaction with its business write.
  • Dedupe retention exceeds the maximum redelivery window implied by broker retention plus retry policy.
  • Messages are keyed by aggregate id, and per-partition lag is monitored.
  • Retries happen on retry topics, never by sleeping inside a handler.
  • Permanent failures skip retries and go directly to a DLQ with an owner.
  • Schemas are registered, and incompatible changes fail at publish time.
  • Trace context propagates through message headers.
  • A reconciliation job exists and runs on a schedule.
  • Someone has rehearsed a DLQ replay in a staging environment within the last quarter.

The takeaway

Exactly-once delivery is not a feature you buy; it is a property you build. Brokers can reduce the frequency of duplicates, and transactions can make multi-topic writes atomic, but no broker can guarantee that a message was processed exactly once by a service it cannot see. The moment you accept that, the architecture becomes clear: make publishing reliable with the outbox, make consuming harmless with idempotency, make workflows recoverable with sagas and compensations, and make the whole thing visible with metrics that measure staleness and duplication rather than throughput. Reliability is not the absence of duplicate messages — it is the presence of a system in which duplicates do not matter.

Comments

No comments yet. Why don’t you start the discussion?

Leave a Reply

Your email address will not be published. Required fields are marked *