Skip to content

VOL 05 / CH 17 / LESSON 02

17.2 Acknowledgment Boundaries, Idempotent Consumption, and Backlog Recovery ​

"The message is processed exactly once" sounds like a broker configuration, but in reality, it spans producers, brokers, consumers, and the business database. As long as there's any possibility that an acknowledgment could be lost, the sender must make a decision between "possibly succeeded" and "possibly failed."

Three Delivery Semantics: First Define the Scope ​

  • at-most-once: Messages may be lost, but no active retry is performed;
  • at-least-once: At least one delivery within the specified fault model and retention conditions, with possible redelivery during recovery;
  • exactly-once: Within a clearly defined boundary, the final outcome is observationally equivalent to a single delivery.

"Exactly-once" must explicitly define its scope. Kafka transactions can atomically write across multiple Kafka partitions and allow read_committed consumers to exclude uncommitted and aborted transaction records. However, this does not automatically include a single email, an HTTP call, or any external database within the same transaction boundary.

Producer Confirmation Still Introduces Uncertainty ​

If a producer times out after sending a message:

text
broker does not receive -> retry is necessary
broker has persisted the message but confirmation is lost -> retry may result in duplicates

Kafka's idempotent producer uses mechanisms like producer identity and partition sequence to eliminate duplicate retries within a session scope. Similarly, in RabbitMQ, if the broker confirms a publication but the connection fails, the producer may not know whether it received the confirmation. As a result, business events should still rely on stable event_id or idempotent keys.

Consumer Ack Must Follow Business Processing ​

text
receive
validate schema
perform idempotent business transaction
record event_id / resulting state
ack or commit offset

Acknowledging the message before writing to the database risks permanently losing business processing if the process crashes between those two steps. Writing to the database before acknowledging introduces the possibility of message redelivery, but an idempotent business table can detect and ignore duplicates.

This standalone PostgreSQL 16+ lab uses one fixed event ID and an increment of 5. Create the tables once, then replay BEGIN through SELECT to simulate redelivery. The anonymous block commits deduplication and the business update together.

sql
CREATE TABLE consumed_event (
  consumer_name text NOT NULL,
  event_id uuid NOT NULL,
  consumed_at timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (consumer_name, event_id)
);
CREATE TABLE inventory_projection (
  item_id bigint PRIMARY KEY,
  received bigint NOT NULL
);
INSERT INTO inventory_projection VALUES (42, 0);

BEGIN;
DO $$
DECLARE affected bigint;
BEGIN
  INSERT INTO consumed_event (consumer_name, event_id)
  VALUES ('inventory-projector', '00000000-0000-0000-0000-000000000001')
  ON CONFLICT (consumer_name, event_id) DO NOTHING;
  GET DIAGNOSTICS affected = ROW_COUNT;
  IF affected = 1 THEN
    UPDATE inventory_projection SET received = received + 5
    WHERE item_id = 42;
    GET DIAGNOSTICS affected = ROW_COUNT;
    IF affected <> 1 THEN
      RAISE EXCEPTION 'projection target missing';
    END IF;
  END IF;
END $$;
COMMIT;
SELECT received FROM inventory_projection WHERE item_id = 42;

The first and repeated executions both return received=5. A missing target raises an error and rolls back the deduplication insert too, allowing a retry after repair. Real applications must bind validated event fields rather than reuse these constants. The retention period for idempotent records must cover both the longest possible replay window of the broker and the manual recovery window.

Outbox Solves the Dual Write Problem of Database Commit and Message Publication ​

If a business transaction first modifies the database and then publishes a message, the process might crash between these two steps. If it first publishes a message and then updates the database, consumers might receive facts that haven't yet been validated or committed.

Create these order and outbox tables in the practice database. The data-modifying CTE creates an event only for an order that actually transitions from pending to paid, preventing a zero-row UPDATE from publishing a false payment fact:

sql
CREATE TABLE orders (
  order_id bigint PRIMARY KEY,
  status text NOT NULL CHECK (status IN ('pending', 'paid'))
);
CREATE TABLE outbox_event (
  event_id uuid PRIMARY KEY,
  aggregate_type text NOT NULL,
  aggregate_id bigint NOT NULL,
  event_type text NOT NULL,
  payload jsonb NOT NULL
);
INSERT INTO orders VALUES (9182, 'pending');

BEGIN;
WITH changed AS (
  UPDATE orders SET status = 'paid'
  WHERE order_id = 9182 AND status = 'pending'
  RETURNING order_id, status
)
INSERT INTO outbox_event (
  event_id, aggregate_type, aggregate_id, event_type, payload
)
SELECT '00000000-0000-0000-0000-000000000002'::uuid,
       'order', order_id, 'OrderPaid',
       jsonb_build_object('order_id', order_id, 'status', status)
FROM changed;
COMMIT;
SELECT status FROM orders WHERE order_id = 9182;
SELECT count(*) FROM outbox_event;

The first execution returns paid and event count 1. Replay only BEGIN through the two SELECT statements: the count remains 1. A nonexistent order ID must also produce no event. If the outbox insert fails, the order update in the same statement rolls back.

A separate relay or CDC component reads from the outbox to publish events. Relays may re-publish messages, so consumers must still be idempotent. While the Outbox ensures atomic creation of facts, it does not guarantee end-to-end exactly-once semantics across the entire pipeline.

Order Holds Only Within Selected Keys and Channels ​

If an order event must maintain sequence, it should use the same partition key/queue and carry the aggregate version:

json
{
  "event_id": "...",
  "aggregate_id": "order-9182",
  "aggregate_version": 7,
  "event_type": "OrderPaid"
}

Consumers of full-state snapshots can reject older versions. With delta events, skipping a version and accepting a later one may lose a decrement; hold the gap or rebuild from the authoritative source. There is no inherent global order across partitions; parallel processing, retry queues, and dead-letter replay can alter the final sequence.

Poison Message and Retry Budget ​

Deserialization failures, permanent business constraint errors, and transient dependency faults cannot be handled with the same infinite retry loop.

Recommended categorization:

  • transient: exponential backoff, jitter, and finite retry counts;
  • rate-limited: respect the service's retry window and enforce rate limiting;
  • permanent/schema: route to an isolation queue, preserve the original message, error details, and version;
  • operator action: replay after manual correction, maintaining original key order.

Dead-letter queues are not trash bins. They must include alerts, ownership assignment, remediation procedures, replay tools, and a strategy for handling subsequent failures.

Backpressure and Backlog Recovery ​

Queues can turn burst traffic into backlog, but they can't magically increase downstream capacity. You should also monitor:

  • ingress rate versus sustainable processing rate;
  • consumer lag or backlog age, rather than message count alone;
  • Processing latency and failure rate per message;
  • skew between partition and queue;
  • Broker disk, replication, and retention quotas;
  • Time required to recover from peak backlog to normal levels.

Assume 6 million queued events, ongoing ingress of 8,000/s, and sustainable processing of 10,000/s. Net drainage is 2,000/s, so ideal recovery takes 3,000 seconds, or 50 minutes; retries and hotspots extend it. If ingress is instead 10,000/s with capacity 8,000/s, backlog keeps growing until retention or storage limits intervene. Peak shaving works only if the long-term average consumption rate exceeds the long-term average production rate, or if the system can discard or aggregate some events.

RabbitMQ consumer prefetch controls the number of unacknowledged deliveries; in Kafka, max.poll.interval.ms, batch size, and processing thread model collectively influence rebalance and throughput. Parameters should be tuned around the cost of processing a single message and memory budget through benchmarking, not copied from fixed recommendations.

Schema Evolution ​

Once messages enter the long-retention log, both old consumers and historical replay processes will encounter them. Events must include versioning and compatibility strategies:

  • Adding optional fields is generally backward compatible;
  • Changing the meaning of a field is more dangerous than changing its name;
  • Before removing a field, verify that all consumers and historical replay mechanisms are unaffected;
  • Events should represent facts that have occurred, not be reused as remote procedure call command packages;
  • In CI pipelines, validate schema compatibility and representativeness of older message samples.

References ​

The next chapter will no longer compare product marketing claims, but instead establish a selection process that records facts, constraints, costs, and exit strategies.

Built with VitePress | Software Systems Atlas