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.
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.
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.
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
- Exactly-Once vs At-Least-Once Delivery — why consumers must be idempotent.
- Preventing Duplicate Job Execution with Idempotency — patterns across frameworks.
- Transactional Outbox Pattern for Job Enqueue — the producer-side counterpart.
- Replaying Dead-Letter Messages in RabbitMQ — replays these claims protect against.