Skip to content
Rafe Uddaraj

Engineering Roadmap: Why 'Exactly Once' Delivery Is Mostly a Lie

11 min readEnglishRead in Bangla
On this page

Imagine a real-world production incident. You built a payment API for an e-commerce platform. When a user clicks the checkout button, your payment service processes the request and calls an external payment gateway like Stripe or PayPal. During system design review, the architecture flow looked clean and predictable.

Then, at 3 AM, your PagerDuty alert fires. Customer support reports that hundreds of users were charged twice for the exact same order on their credit cards.

You open your production logs to trace what happened behind the scenes:

  1. Your Payment Service dispatched an HTTP POST request to the payment gateway: Charge $100.
  2. The payment gateway processed the charge successfully and debited $100 from the user's card.
  3. Just one millisecond before the 200 OK response packet reached your server, an intermediate network link dropped, causing a TCP Socket Timeout on your application server.
  4. From the perspective of your Payment Service, it cannot determine whether the money was charged or if the request dropped before reaching the gateway.
  5. To keep the service fault-tolerant, an automatic retry mechanism kicked in and dispatched an identical retry request.
  6. The payment gateway received the second request as a new transaction and charged the user another $100.
Double Charge Production Incident
Double Charge Production Incident

A user was billed $200 instead of $100. This illustrates the classic crisis of Exactly-Once Delivery in distributed systems.

While an "Exactly-Once" guarantee sounds appealing, achieving it across unreliable networks is mathematically impossible. When engineers migrate from monoliths to microservices, the most common pitfall is treating the network as fully reliable.

In this deep dive, we explore the low-level mechanics of distributed systems. We examine why network-level exactly-once delivery cannot exist and how senior engineers construct Effectively-Once architectures in production using Idempotency, Transactional Outbox, and Inbox Patterns.


Network Uncertainty and The Lost ACK Problem

To understand fundamental uncertainty in distributed architectures, we must examine the physical realities of the network layer. Whenever two independent nodes communicate over a network, the intermediate links, switches, routers, and middleboxes are inherently fallible.

A foundational theorem in computer science, the Two Generals' Problem, mathematically proves that two separate nodes communicating over an unreliable channel can never reach 100% consensus on state using finite messages.

When a client transmits a request to a server, three distinct failure scenarios can occur from the client's perspective:

  1. Request Lost in Transit: The request packet drops due to network congestion or routing failure before reaching the server.
  2. Server Crashes Pre-Execution: The server receives the packet but crashes (power loss, kernel panic, OOM) before execution completes.
  3. Execution Succeeded but ACK Lost: The server receives the request, mutates state in its database, but the acknowledgment (ACK / 200 OK) packet drops on the return path.
The Lost ACK Problem
The Lost ACK Problem

The critical challenge is that when the sender encounters a timeout error, it receives identical symptoms (ETIMEDOUT / ECONNRESET) across all three distinct scenarios.

If the client assumes failure and retransmits, scenario 3 creates a duplicate mutation. If the client assumes success, scenarios 1 and 2 result in silent, permanent data loss.

The Post Office Analogy

If you send an important letter through the mail and receive no response, you cannot determine if the postal carrier lost the envelope or if the recipient received it and simply forgot to reply. To guarantee receipt, you are forced to send a duplicate letter.


Three Delivery Semantics: At-Most-Once, At-Least-Once, Exactly-Once

To navigate network uncertainty, distributed systems classify message transport into three distinct operational models:

Delivery Semantics Comparison
Delivery Semantics Comparison

1. At-Most-Once (0 or 1 delivery)

The system transmits a message once. If an error or timeout occurs, it never retries (Fire and Forget).

  • Advantage: Guarantees zero duplicate processing.
  • Risk: Any network hiccup leads to permanent data loss.
  • Use Case: High-throughput telemetry, real-time metrics, video streaming packets, and high-frequency IoT sensors where dropped samples are acceptable.

2. At-Least-Once (1 or more deliveries)

The sender repeatedly retransmits the message until it receives explicit acknowledgment from the receiver.

  • Advantage: Message loss rate drops to near zero.
  • Risk: Whenever an ACK packet drops, the receiver receives duplicate messages.
  • Use Case: Financial ledgers, order workflows, and messaging systems like Apache Kafka, RabbitMQ, and AWS SQS by default.

3. Exactly-Once (Exactly 1 delivery)

A theoretical model where a message travels across the network exactly once and executes business logic exactly once. At the pure network transport layer, this is physically impossible. However, we can simulate an Effectively-Once model by combining At-Least-Once transport with application and storage-level deduplication.


The Core Misconception: Delivery vs. Processing

A common architectural mistake is conflating message delivery with message execution:

  • Delivery (Transport Layer): Transferring bytes from the broker or sender across a socket to the consumer.
  • Processing (Application Layer): Parsing the payload and executing domain logic in memory.
  • Side Effect (Persistence / External Layer): Committing state changes, such as mutating database rows, creating charges on Stripe, or dispatching emails.

You cannot prevent duplicate packet delivery across an unreliable network. However, if you design your execution layer so that receiving the same payload ten times produces exactly one side effect, the system behaves as Effectively-Once to your users.


The Real Enemy is Duplicate Side Effects, Not Duplicate Messages

Receiving a duplicate network frame causes minimal harm. Catastrophic failures occur when duplicate payloads trigger duplicate side effects:

  • Charging a user's credit card twice
  • Decrementing warehouse inventory twice
  • Inserting duplicate database records
  • Dispatching multiple confirmation notifications
The Real Enemy: Duplicate Side Effects
The Real Enemy: Duplicate Side Effects

Production architectures must embrace duplicate transport while building strict guards around every persistent side effect.


Idempotency: The Practical Escape Hatch

Idempotency is the primary defensive mechanism against duplicate execution.

In mathematics and computer science, an operation is defined as idempotent if applying it multiple times yields the exact same result as applying it once:

Idempotency Flow Architecture
Idempotency Flow Architecture

How Idempotency Keys Work in Practice:

  1. The client generates a unique UUID for the intent and attaches it to the HTTP header: Idempotency-Key: 9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d.
  2. Before running business logic, the server checks the key in a shared persistence layer (e.g., PostgreSQL or Redis).
  3. If the key exists with status COMPLETED, the server bypasses execution and immediately returns the cached response.
  4. If the key is new, the server acquires a lock (IN_PROGRESS), executes domain logic, persists the result, and transitions status to COMPLETED.

REST Semantics and Idempotency

In REST API design, GET, PUT, and DELETE verbs are naturally idempotent by specification. Repeating a DELETE /items/42 multiple times leaves the system in the same target state. In contrast, POST is non-idempotent by default. Production POST endpoints handling billing, fund transfers, or order submission should mandate an Idempotency-Key header.


Database Unique Constraints and the Check-Then-Insert Trap

A naive implementation of idempotency relies on application memory checks:

JavaScript
// DANGEROUS ANTI-PATTERN: Check-Then-Insert Race Condition
const existing = await db.query('SELECT * FROM payments WHERE idempotency_key = $1', [key]);
if (!existing.rows.length) {
// If concurrent requests arrive simultaneously, both evaluate to true!
await chargeCard(amount);
await db.query('INSERT INTO payments (idempotency_key, status) VALUES ($1, $2)', [key, 'DONE']);
}

If two duplicate requests hit separate node instances concurrently, both queries return empty results, and both threads proceed to charge the user. This is the classic Check-Then-Insert race condition.

Application-level checks cannot guarantee safety in concurrent systems. Database unique constraints provide the only authoritative synchronization boundary.

Database Level Unique Constraint
Database Level Unique Constraint
SQL
CREATE TABLE idempotency_keys (
key VARCHAR(255) PRIMARY KEY,
status VARCHAR(50) NOT NULL, -- 'IN_PROGRESS', 'COMPLETED', 'FAILED'
response_body JSONB,
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);

Database B-Tree indexes enforce row uniqueness at the storage engine level. When multiple workers attempt concurrent inserts for the same key, exactly one transaction succeeds while all others fail with a unique violation error.


Consumer Crashes and the Inbox Pattern

When processing asynchronous message queues (Kafka, RabbitMQ, SQS), worker crashes create duplicate processing risks:

  1. A consumer retrieves a message from the queue.
  2. The consumer credits $100 to the user's account balance.
  3. Just before dispatching the ACK back to the broker, the consumer process suffers an Out of Memory (OOM) crash.

Because the broker never receives an ACK, it assumes delivery failed. Upon consumer reboot, the broker redelivers the message, causing another $100 credit.

Inbox Pattern for Consumer Crashes
Inbox Pattern for Consumer Crashes

To prevent this vulnerability, engineers deploy the Inbox Pattern.

The core invariant of the Inbox Pattern: Persisting the processed message identifier and updating domain data must execute within a single atomic database transaction.

SQL
BEGIN;
-- 1. Insert message ID into inbox (message_id holds a UNIQUE constraint)
-- On duplicate delivery, the insert fails and aborts the entire transaction block
INSERT INTO message_inbox (message_id, received_at)
VALUES ('msg-ord-98765', NOW());
-- 2. Mutate business state within the same transaction scope
UPDATE user_wallets
SET balance = balance + 100
WHERE user_id = 'usr-001';
-- 3. Record audit trail
INSERT INTO wallet_transactions (user_id, amount, type)
VALUES ('usr-001', 100, 'CREDIT');
-- 4. Commit atomically
COMMIT;

If the server crashes prior to COMMIT, the database rolls back all mutations. If it crashes after COMMIT, the message ID is securely persisted, causing future redelivery attempts to trigger a unique constraint violation and exit cleanly.


The Dual-Write Problem and Transactional Outbox Pattern

The Inbox Pattern protects message consumption. But how do we guarantee reliable message publishing?

Consider an order service that writes an order to the database and then publishes an OrderCreated event to Kafka:

JavaScript
// DANGEROUS ANTI-PATTERN: The Dual-Write Problem
await db.query('INSERT INTO orders ...'); // Step 1: Database Commit
await kafkaProducer.send('OrderCreated', ...); // Step 2: Process crash or network loss!

If the process crashes between Step 1 and Step 2, the order exists in the database, but downstream inventory services never receive notification. This is the Dual-Write Problem: two distinct distributed storage engines cannot participate in an atomic commit without complex coordination protocols.

The standard solution is the Transactional Outbox Pattern:

  1. Domain mutations and outbound message payloads are written to an outbox_events table inside the same local database transaction.
  2. Because both writes occur in a single database transaction, ACID properties guarantee both succeed or both roll back.
  3. An asynchronous relay process (or Change Data Capture engine like Debezium) tails the outbox table and streams events to the message broker.
Transactional Outbox Architecture
Transactional Outbox Architecture

Outbox Guarantees

The Transactional Outbox pattern provides At-Least-Once delivery. If the relay worker publishes an event to Kafka and crashes before updating the outbox record, the event will be republished. Consequently, downstream consumers must implement the Inbox Pattern.


Retry Storms, Exponential Backoff, and Jitter

When a downstream service experiences partial degradation, upstream clients operating under At-Least-Once semantics begin retrying aggressively.

When the downstream service attempts recovery, it is hit with an overwhelming surge of queued retries, triggering an immediate relapse. This self-inflicted outage is known as a Retry Storm (or Thundering Herd problem).

Retry Storms vs Exponential Backoff with Jitter
Retry Storms vs Exponential Backoff with Jitter

To mitigate retry storms, production systems combine two strategies:

  1. Exponential Backoff: Increasing retry intervals exponentially ( seconds).
  2. Full Jitter: Introducing random variation into retry delays to distribute incoming traffic evenly across time.
JavaScript
// Production-grade Exponential Backoff with Full Jitter
async function fetchWithResilience(url, options, maxRetries = 5) {
for (let attempt = 0; attempt < maxRetries; attempt++) {
try {
const response = await fetch(url, options);
if (response.ok) return response;
// Do not retry 4xx client errors (excluding 429 Rate Limit)
if (response.status !== 429 && response.status < 500) {
throw new Error(`Client Error: ${response.status}`);
}
} catch (error) {
if (attempt === maxRetries - 1) throw error;
// 1. Calculate Exponential Base Delay: 1s, 2s, 4s, 8s, 16s...
const baseDelay = Math.pow(2, attempt) * 1000;
// 2. Full Jitter: Uniform random distribution in [0, baseDelay]
const delay = Math.random() * baseDelay;
console.warn(`Attempt ${attempt + 1} failed. Retrying in ${Math.round(delay)}ms...`);
await new Promise(resolve => setTimeout(resolve, delay));
}
}
}

The Myth vs. Reality of Apache Kafka's Exactly-Once

A prevalent industry misunderstanding is that enabling Kafka's Exactly-Once Semantics (processing.guarantee=exactly_once_v2) guarantees end-to-end exactly-once delivery across your entire stack.

Kafka's transactional guarantees apply strictly within closed Kafka boundaries:

  • Where Kafka EOS Applies: Consuming from Kafka Topic A, performing in-memory stream transformations (Kafka Streams), and writing to Kafka Topic B within a single transactional coordinator boundary.
  • Where Kafka EOS Fails: When a consumer calls an external HTTP endpoint (e.g., Stripe, SendGrid, or a third-party microservice). If the consumer crashes during the external network call, Kafka rolls back the read offset, triggering re-execution and duplicate outbound API calls.
Kafka EOS Boundary Limitations
Kafka EOS Boundary Limitations

End-to-End Limitations: Database Transactions vs. External APIs

You cannot atomically bind an external HTTP network request inside a relational database transaction.

End-to-End Failure Matrix
End-to-End Failure Matrix
JavaScript
// HIGHLY DANGEROUS ANTI-PATTERN: Network calls inside DB Transactions
await db.query('BEGIN');
await db.query('UPDATE accounts SET balance = balance - 100 WHERE id = $1', [userId]);
// If this network request hangs for 15s, the database row lock remains open!
// If DB COMMIT fails after Stripe succeeds, the customer is charged without balance deduction!
const paymentRes = await stripe.charges.create({ amount: 10000 });
await db.query('COMMIT');

Golden Rule of Database Transactions

Never execute external HTTP calls or network I/O inside an open database transaction block. Latency spikes will exhaust your connection pool, and partial failures leave distributed state unrecoverable.


Real-World Architecture: Payment Processing Pipeline

Below is a robust, production-grade payment architecture designed to handle network failures and duplicate events across every tier:

Production Payment System Architecture
Production Payment System Architecture

End-to-End Lifecycle:

  1. Client Layer: The client application submits a payment request with a unique client request ID.
  2. Ingress & DB Guard: The API service writes an outbox record with an idempotency key and returns a 202 Accepted response.
  3. Async Processing: A background worker picks up the outbox event and initiates the external gateway call with a dedicated provider key.
  4. Gateway Interaction: The worker communicates with Stripe using the idempotency key. Network timeouts trigger safe retries using the same key.
  5. Webhook Deduplication: Incoming gateway webhooks are validated against an inbox table by event ID, preventing duplicate order completion logic.

Effectively-Once System Architecture

High availability and consistency require defense in depth across all architectural tiers:

Distributed Reliability Layers
Distributed Reliability Layers
  1. Transport Layer (At-Least-Once): Automatic retry mechanisms with timeout handling.
  2. Traffic Smoothing Layer (Backoff & Jitter): Mitigating thundering herds during system degradation.
  3. Consumer Guard (Inbox Pattern & Unique Constraints): Eliminating race conditions and duplicate writes at the persistence layer.
  4. Publisher Guard (Transactional Outbox): Eliminating dual-write inconsistencies between database records and message queues.
  5. Workflow Orchestration (Saga & Compensations): Executing compensating transactions to gracefully roll back distributed state upon terminal errors.

Recap and Production Readiness Checklist

Keep this distributed systems failure matrix in mind when designing resilient backends:

Distributed System Failure Matrix
Distributed System Failure Matrix

Production Readiness Checklist:

  • Clarify Terminology: Distinguish transport semantics from processing semantics in architecture reviews.
  • Enforce DB Uniqueness: Eliminate in-memory checks in favor of database unique constraints.
  • Mandate Idempotency Keys: Require idempotency headers across all state-mutating POST endpoints.
  • Implement the Inbox Pattern: Bind message ID persistence and domain mutations in a single atomic transaction.
  • Implement the Outbox Pattern: Decouple event publishing from transaction boundaries using transactional outbox tables.
  • Isolate DB Transactions: Keep all network calls and third-party APIs outside database transaction blocks.
  • Apply Backoff with Jitter: Configure exponential backoff with full jitter across all retry policies.
  • Deduplicate Webhooks: Validate incoming webhook event IDs against an inbox table before processing.
  • Structured Distributed Tracing: Propagate trace_id, idempotency_key, and message_id across all application logs.

Network partitions in distributed systems are inevitable operational realities. By pairing At-Least-Once delivery with database-enforced idempotency, your system maintains bulletproof consistency regardless of network volatility.

All articles

Get in touch

Questions about a video, an article, or working together.