Guaranteeing Event Publication With the Outbox Pattern
Goal
After creating the dual-write inconsistency yourself, eliminate loss with the outbox pattern and absorb the remaining duplicates with consumer idempotency, completing the whole path by hand.
Why it matters
Code that calls DB saving and event publishing side by side is everywhere, and works fine normally. The problem appears only at the moment the broker wobbles for 3 seconds, and the inconsistency created then quietly remains. You discover a state where the order exists but the inventory did not decrease days later in settlement. The outbox solves this problem by flipping it from "let's make the broker take part in the transaction" to "let's write the intent to publish into the DB as well". In exchange, a new property arises — loss disappears but duplication appears. If the relay dies right after publishing and before updating the status, it sends the same event again after restarting. So this lab treats duplication not as a bug but as a design premise, and finishes within one lab even the absorbing of it on the consumer side.
Steps
- With
/root/outbox/init.py, create/root/outbox/app.db. There must be two tables,orders(id, sku, qty)andoutbox(event_id, aggregate_id, seq, event_type, payload, status), and the default ofstatusisPENDING. /root/outbox/dualwrite.pysaves one order and then fails to publish. After running it, on the first line of/root/outbox/dualwrite.out, writeINCONSISTENT orders=<n> published=<m>. n and m must differ./root/outbox/place_order.pywrites to orders and outbox in the same transaction. Running it 3 times results in 3 rows in orders and 3 rows in outbox./root/outbox/relay.pyRPUSHes the rows withstatus='PENDING'to the Redis listoutbox.eventsand changes those rows toPUBLISHED. The list length must not grow even if you run it twice.- Events with the same
aggregate_idmust enter the queue in ascendingseqorder. In/root/outbox/order_check.out, leaveORDER OK. /root/outbox/relay_crash.pyonly publishes and does not update the status. If you run a normal relay again afterward, the sameevent_identers the queue twice. In/root/outbox/atleastonce.out, writeDUPLICATE event_id=<id> count=2./root/outbox/consumer.pydrains the queue while filtering duplicates byevent_id. Leave the processing result in the Redis hashprocessed. Even with duplicates, the size ofprocessedmust equal the number of unique events.- In
/root/outbox/report.txt, write four lines:orders=<n>,outbox=<n>,enqueued=<n>, andprocessed_unique=<n>.enqueuedmust be greater than or equal toprocessed_unique.
Notes
- A partial index for looking up
PENDING:CREATE INDEX ... ON outbox(status) WHERE status='PENDING' - Checking Redis:
redis-cli LLEN outbox.events,redis-cli HLEN processed - Common mistake 1: the relay trying hard to make publishing and the status update a single atomic unit — it is impossible. Accepting duplicates and filtering them at the consumer is the answer.
- Common mistake 2: sorting by timestamp without
seq— if two events occur in the same millisecond, the order flips.
Create the orders and outbox tables
With /root/outbox/init.py, create /root/outbox/app.db. There must be two tables, orders(id, sku, qty) and outbox(event_id, aggregate_id, seq, event_type, payload, status), and the default of status is PENDING.
Python's sqlite3 module is enough. An outbox row needs an event identifier, an aggregate identifier, a type, a body, and a status.
Reproduce the dual-write inconsistency
/root/outbox/dualwrite.py saves one order and then fails to publish. After running it, on the first line of /root/outbox/dualwrite.out, write INCONSISTENT orders=<n> published=<m>. n and m must differ.
Just make the DB save succeed and only the broker publication fail. Count the items in the two stores and save in a file that they differ.
Fold it into a single transaction
/root/outbox/place_order.py writes to orders and outbox in the same transaction. Running it 3 times results in 3 rows in orders and 3 rows in outbox.
Put the two INSERTs in the same connection and the same commit. If an exception occurs before the commit, neither should exist.
Move it to the broker with a relay
/root/outbox/relay.py RPUSHes the rows with status='PENDING' to the Redis list outbox.events and changes those rows to PUBLISHED. The list length must not grow even if you run it twice.
Read the PENDING rows, push them into the Redis list, and change the status. Even if you run it several times, it must not move again what was already moved.
Preserve order per aggregate
Events with the same aggregate_id must enter the queue in ascending seq order. In /root/outbox/order_check.out, leave ORDER OK.
Events of the same order must go out in the order they occurred. Think about what to use as the sort criterion — timestamps can be equal.
Create duplicates with a relay crash
/root/outbox/relay_crash.py only publishes and does not update the status. If you run a normal relay again afterward, the same event_id enters the queue twice. In /root/outbox/atleastonce.out, write DUPLICATE event_id=<id> count=2.
Imitate the situation where it publishes but dies without being able to update the status. If you run it again, the same event enters the queue twice.
Add deduplication to the consumer
/root/outbox/consumer.py drains the queue while filtering duplicates by event_id. Leave the processing result in the Redis hash processed. Even with duplicates, the size of processed must equal the number of unique events.
You just need to remember the identifiers of events already processed. A Redis set data structure or SET NX fits.
Report the whole path in numbers
In /root/outbox/report.txt, write four lines: orders=<n>, outbox=<n>, enqueued=<n>, and processed_unique=<n>. enqueued must be greater than or equal to processed_unique.
Gather into one file the number of orders, the number of outbox rows, the number put in the queue, and the number of unique items actually processed. The relationship between the first three values and the last one is the whole of this pattern.