- Published on
Where a Transaction Ends in a Message Consumer
- Authors

- Name
- Mehdi Akiki
Reference
A message handler can fit on one screen and still cross three transaction systems:
consume broker record
update application database
call external API
publish result event
acknowledge broker record
Putting these lines inside one function does not make them atomic. I draw each durable boundary before I decide retry behaviour.
The first unavoidable gap
Take the smallest useful consumer:
read message
write database
acknowledge message
The database and broker usually cannot commit as one ordinary local transaction. I must choose the order.
If I acknowledge first:
ack broker
crash
database write never happens
The message is lost to this consumer.
If I commit the database first:
commit database
crash
broker redelivers
database write runs again
The delivery is duplicated. I normally prefer this direction because a duplicate can be detected, while an acknowledged but unapplied record may be unrecoverable.
Put identity inside the database transaction
The inbox pattern couples duplicate detection and the domain write:
begin;
insert into consumer_inbox (consumer, message_id, payload_hash)
values (:consumer, :message_id, :payload_hash)
on conflict do nothing
returning message_id;
-- Only when the insert returned a row:
update account
set balance = balance + :delta
where account_id = :account_id;
commit;
After commit, the consumer acknowledges the broker record.
If it crashes before acknowledgement, redelivery finds the inbox identity and skips the domain mutation. For the complete pattern, see The Inbox Pattern: Deduplicate Before Business Logic Runs.
Write the crash matrix
For the database-and-ack path, I expect:
| Crash point | Durable state | Redelivery outcome |
|---|---|---|
| before transaction | nothing | process normally |
| after inbox insert, before commit | transaction rolls back | process normally |
| after domain write, before commit | transaction rolls back | process normally |
| after commit, before broker ack | inbox and domain state present | detect duplicate, then ack |
| after broker ack | completed | no normal redelivery |
This table is more valuable than saying the consumer is “exactly once.” It identifies the stored evidence that makes each retry safe.
An external API is outside this transaction
Now add:
commit database
charge payment provider
acknowledge message
The database transaction cannot roll back a payment already accepted by another service. Holding the database lock while making the HTTP call does not make the provider part of PostgreSQL.
There are two dangerous crash directions:
provider charged, local completion not stored -> possible duplicate charge
local completion stored, provider not called -> missing charge
I give the external effect a stable idempotency key derived from the intended business operation. After an ambiguous timeout, I query or retry with the same identity instead of creating a new charge.
If the provider does not support idempotency or reconciliation, the guarantee is weaker and must be stated as such.
Publish through an outbox
Suppose the consumer must update the database and emit AccountCredited.
Publishing directly creates another gap:
database commit succeeds
event publish fails
I write the domain mutation and an outbox row in one local transaction:
begin;
insert into consumer_inbox (...);
update account set ...;
insert into outbox (
event_id,
event_type,
aggregate_id,
payload,
created_at
) values (
:event_id,
'AccountCredited',
:account_id,
:payload,
now()
);
commit;
A separate relay publishes the outbox event and records delivery progress. This turns “database plus broker” into two recoverable local steps rather than one imaginary distributed transaction.
Broker transactions have a specific boundary
Some brokers support transactions that couple consumed positions and produced records inside that broker. Kafka's producer API documents transactions that can atomically send output records and commit offsets for consumed records.
That is useful for a Kafka-to-Kafka topology:
consume Kafka A
produce Kafka B
commit consumed offsets in the Kafka transaction
It does not automatically include an unrelated SQL database or SaaS HTTP call. I name the participants before using the word “transactional.”
Long work makes ownership another boundary
A consumer can lose its lease or partition assignment while processing. The old worker may continue and commit after a new worker has received the same record.
For a short database transaction, inbox uniqueness protects the business mutation. For long-running work, I may also need:
- a claim with owner and expiry;
- periodic lease renewal;
- a fencing token checked at commit;
- cancellation when partition ownership changes;
- idempotent external-effect identity.
A mutex inside one process cannot fence another worker.
Do not make the transaction wider than necessary
It is tempting to hold a database transaction open while parsing a large file, calling an API, or waiting for a model. This increases lock time and still does not make remote effects atomic.
I separate phases:
1. validate immutable input
2. short transaction: claim identity and record intended work
3. perform bounded external work with stable effect identity
4. short transaction: record observed outcome and enqueue next event
5. acknowledge when the recovery contract permits it
The exact sequence depends on whether failure may leave work pending, but every durable transition has an owner and a retry rule.
Test boundaries by crashing, not by reading code
I add fault injection after each important operation:
after source receive
after inbox claim
after domain mutation
after database commit
after external request reaches provider
after outbox publish
before and after source acknowledgement
For each point, I restart the worker and assert the final business effect count, not only that the handler returned success.
A transaction ends at the system that commits it. Once a handler touches another database, broker, or API, there is a gap. Reliable consumers do not hide these gaps; they fill them with identity, durable intent, reconciliation, and tested recovery.