Idempotent Consumers with Postgres Unique Constraints

When a consumer's side effect is a database write, the database itself can guarantee that the effect happens once, however many times the message is delivered. This guide builds that guarantee with unique constraints as part of Exactly-Once vs At-Least-Once Delivery in Queue Fundamentals & Architecture. Unlike TTL-based keys in a cache, a unique constraint never expires, survives failovers with the data it protects, and is checked atomically with the write.

Problem Statement

A loyalty service consumes OrderCompleted events from RabbitMQ and credits points to customer balances. Redeliveries — from consumer crashes, network blips, and a replay of a dead-letter queue last month — have credited some orders two or three times; finance found 11,000 excess points in a quarterly audit. An earlier fix checked "has this order been credited?" before crediting, but two consumers processing the same redelivered message concurrently both passed the check. You want each order credited exactly once regardless of redelivery timing, concurrency, or replays weeks later, with a check that costs one extra row write.

Prerequisites

  • Postgres 12+ (any version with INSERT ... ON CONFLICT).
  • A consumer that performs its side effect as database writes in one transaction.
  • A stable identity per message effect: a producer-assigned event id, or a natural key such as the order id.
  • Late acknowledgement: ack the message after the database transaction commits.

Step 1 — Create a Processed-Messages Table

A small table with a primary key on the effect identity is the ledger of what has been applied. Keep it narrow; it will receive one row per processed message.

CREATE TABLE processed_messages (
  consumer     text        NOT NULL,             -- which consumer applied it
  message_key  text        NOT NULL,             -- event id or natural key
  processed_at timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (consumer, message_key)
);
-- Supports time-based cleanup (Step 5)
CREATE INDEX processed_messages_time ON processed_messages (processed_at);

Including consumer in the key lets several services share the table while each keeps its own record — the loyalty consumer and the analytics consumer may both process OrderCompleted:9912, and neither should block the other.

Claim and effect in one transaction For each delivery, the consumer opens a transaction, inserts a row keyed by consumer and message key with ON CONFLICT DO NOTHING, and applies the points credit only if the insert affected a row. Then it commits and acknowledges. A duplicate delivery hits the existing primary key, the insert affects zero rows, and the transaction commits without crediting again. One transaction, two writes BEGIN ... COMMIT INSERT processed_messages ON CONFLICT DO NOTHING UPDATE balances only if a row was inserted Duplicate: insert affects 0 rows, no credit, commit, ack. The record never expires.

Step 2 — Insert the Claim and Apply the Effect Together

The whole trick is doing both writes in one transaction and branching on whether the claim inserted.

# consumer.py
import psycopg

def handle_order_completed(conn: psycopg.Connection, msg: dict) -> str:
    with conn.transaction():
        claimed = conn.execute(
            """INSERT INTO processed_messages (consumer, message_key)
               VALUES ('loyalty', %s) ON CONFLICT DO NOTHING""",
            (f"order-completed:{msg['order_id']}",),
        ).rowcount
        if claimed == 0:
            return "duplicate"                                  # already applied: no-op
        conn.execute(
            "UPDATE balances SET points = points + %s WHERE customer_id = %s",
            (points_for(msg), msg["customer_id"]))
        conn.execute(
            "INSERT INTO point_ledger (customer_id, order_id, points) VALUES (%s, %s, %s)",
            (msg["customer_id"], msg["order_id"], points_for(msg)))
    return "applied"

def on_message(channel, method, properties, body):
    result = handle_order_completed(conn, json.loads(body))
    channel.basic_ack(method.delivery_tag)                      # ack AFTER commit, either way

If the process crashes after commit but before the ack, the redelivered message finds the row and returns duplicate. If it crashes before commit, the transaction rolls back — claim and credit together — and the redelivery applies it normally. There is no window in which the effect happened but the claim did not.

Step 3 — Understand the Concurrency Behaviour

The failed check-then-act approach let two consumers pass the check simultaneously. With INSERT ... ON CONFLICT, the unique index serialises them: the second inserter waits on the first transaction's uncommitted row, then sees the conflict once it commits.

T1: INSERT (loyalty, order-completed:9912)         -> inserts, holds row lock
T2: INSERT (loyalty, order-completed:9912)         -> blocks on T1's uncommitted key
T1: UPDATE balances ...; COMMIT
T2: unblocks -> conflict -> DO NOTHING, rowcount 0 -> returns "duplicate"
(if T1 had rolled back instead, T2's insert would succeed and T2 would apply the credit)

That last line is the property that makes this safe in both directions: a crash in the first consumer does not lose the effect, because the waiting duplicate takes over. Keep transactions short so the waiting is brief; do not call slow external services inside them.

The unique index serialises duplicates Transaction one inserts the message key and holds its lock while crediting points. Transaction two tries to insert the same key and blocks. When transaction one commits, transaction two sees the conflict, inserts nothing, and skips the credit. If transaction one had rolled back, transaction two's insert would succeed and it would apply the credit instead. Concurrent duplicates, one credit txn 1 INSERT key: ok credit points COMMIT txn 2 INSERT same key: waits on txn 1 conflict: 0 rows, no credit No check-then-act gap: the check and the claim are the same statement.

Step 4 — Prefer Natural Keys or Conditional Upserts Where They Exist

A separate processed-messages table is the general solution. Often the effect's own table already has a natural unique key, and the dedup can live there with no extra table.

-- The ledger row itself is the idempotency record: one ledger entry per order
ALTER TABLE point_ledger ADD CONSTRAINT one_credit_per_order UNIQUE (order_id);

WITH credit AS (
  INSERT INTO point_ledger (customer_id, order_id, points)
  VALUES ($1, $2, $3)
  ON CONFLICT (order_id) DO NOTHING
  RETURNING customer_id, points
)
UPDATE balances b SET points = b.points + c.points
FROM credit c WHERE b.customer_id = c.customer_id;      -- runs only if the insert happened

For state updates (last-writer-wins projections), a version-guarded upsert is idempotent without any claim at all: ON CONFLICT (id) DO UPDATE ... WHERE target.version < EXCLUDED.version. That approach, and why it also tolerates reordering, is covered in handling out-of-order events with sequence numbers.

Step 5 — Clean Up Old Claims Without Weakening the Guarantee

A processed-messages table grows by one row per message forever unless cleaned. Deleting rows re-opens the door for very late duplicates, so delete only beyond the longest plausible redelivery or replay window — and document that window.

-- Nightly: keep 30 days of claims (covers redelivery, DLQ triage, and most replays)
DELETE FROM processed_messages
WHERE processed_at < now() - interval '30 days'
  AND ctid IN (SELECT ctid FROM processed_messages
               WHERE processed_at < now() - interval '30 days' LIMIT 50000);
-- repeat until 0 rows; small batches keep locks and WAL bursts short

At high volume, partition the table by day and drop old partitions instead of deleting rows — faster and free of vacuum load, as discussed in Database-Backed Job Queues. For effects where no window is safe (payments), keep the natural-key constraint from Step 4 on the business table permanently; it costs nothing extra.

How long must a claim live? Redeliveries arrive within minutes, dead-letter triage and replays within days to a couple of weeks. A 30-day retention on the processed-messages table covers all of these. Natural unique keys on business tables, such as one ledger entry per order, never expire and cover replays of any age. Duplicate arrival windows redelivery: minutes DLQ triage and replay: days 30-day claim retention Natural keys on business tables have no window at all: use them for money.

Step 6 — Handle Effects That Leave the Database

When the consumer must also call an external API (send an email, charge a card), the database transaction cannot roll it back. Order the steps so the external call is protected by the claim and by the provider's own idempotency key.

def handle_with_external_call(conn, msg):
    key = f"order-completed:{msg['order_id']}"
    with conn.transaction():
        claimed = conn.execute("INSERT INTO processed_messages (consumer, message_key) "
                               "VALUES ('receipts', %s) ON CONFLICT DO NOTHING", (key,)).rowcount
        if not claimed:
            return "duplicate"
        conn.execute("INSERT INTO outbox (kind, payload) VALUES ('send_receipt', %s)",
                     (json.dumps(msg),))                 # external call happens later, once

Writing to an outbox in the same transaction, then having a separate relay perform the call with an idempotency key, keeps the database as the single source of truth — the approach in the transactional outbox pattern.

Verification

def test_redelivery_applies_once(conn):
    msg = {"order_id": 9912, "customer_id": 7, "total_cents": 5000}
    assert handle_order_completed(conn, msg) == "applied"
    assert handle_order_completed(conn, msg) == "duplicate"
    assert balance(conn, 7) == points_for(msg)

def test_concurrent_duplicates_apply_once(pool):
    msg = {"order_id": 9913, "customer_id": 8, "total_cents": 1000}
    with ThreadPoolExecutor(8) as ex:
        results = list(ex.map(lambda _: run_in_new_conn(pool, handle_order_completed, msg), range(8)))
    assert results.count("applied") == 1 and results.count("duplicate") == 7

In production, emit a counter of duplicate outcomes. It shows the real redelivery rate of your broker and consumer setup — useful context the next time someone proposes removing the dedup "because duplicates never happen".

Gotchas & Edge Cases

Autocommit mode. If the claim insert and the effect run in separate autocommit statements, a crash between them leaves a claim with no effect — the message is then skipped forever. Always use one explicit transaction.

Keys from message ids that change. Some brokers assign new message ids on republish or DLQ replay. Key on a producer-assigned event id or a natural business key, not the broker's delivery id.

Serialization failures. Under SERIALIZABLE isolation, concurrent duplicates may fail with serialization errors instead of waiting. Retry the transaction on 40001; the retry will see the conflict and return duplicate.

Hot keys. A natural key that many messages share (a per-customer "daily summary" key) serialises all of them. Use the finest key that still identifies one effect.

FAQ

Is this better than Redis SET NX? For effects that live in the database, yes: the claim commits atomically with the effect and never expires. Redis keys are a good fast pre-check and fit effects outside the database; see deduplicating jobs with Redis SET NX keys.

Does the extra insert hurt throughput? It adds one small index insert per message — typically well under a millisecond. The table's growth and cleanup are the real cost; partitioning handles both.

What if the consumer and the effect use different databases? Put the claim in the same database as the effect — that is what makes them atomic. A claim in database A and a write in database B brings back the dual-write problem. If effects span databases, use the outbox approach from Step 6 so each database changes in its own local transaction, driven by idempotent steps.

Can I use this with Kafka? Yes. It is the standard way to get exactly-once effects when a Kafka consumer writes to a database, where Kafka transactions do not reach. See Kafka exactly-once semantics with transactions.

Related