TT Lab
Get started
Learn Learning paths Courses

Microservice Architecture

Guaranteeing Event Publication With the Outbox Pattern

Continue in TT Lab

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

  1. 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.
  2. /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.
  3. /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.
  4. /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.
  5. Events with the same aggregate_id must enter the queue in ascending seq order. In /root/outbox/order_check.out, leave ORDER OK.
  6. /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.
  7. /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.
  8. 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.

Notes

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.