The Dual-Write Problem and Its Solutions
Summary
The two actions of writing to the DB and publishing to the broker are not atomic. The outbox is a technique that folds those two into one DB transaction.
Why this was needed
Let us start with the most common code.
def create_order(req):
order = db.save(Order(req)) # 1. DB 저장
broker.publish("OrderCreated", order) # 2. 이벤트 발행
return order
These five lines have a serious problem. If 1 succeeded but 2 failed, the order exists but nobody knows. Conversely, if the transaction rolls back after 2 succeeded, the event for an order that does not exist spreads into the world. This is the dual-write problem.
The idea "can't we just bind the two in one distributed transaction?" is natural, but most brokers cannot take part in XA transactions, and 2PC itself creates an availability problem of holding participants hostage when the coordinator fails. So in practice people take a different route.
How it works
The idea of the outbox is simple. Instead of sending the event directly to the broker, you insert one row into the outbox table in the same transaction as the business data. The DB's ACID guarantees the atomicity of the two writes. Then a separate process (the relay) reads that table and moves it to the broker.
There are four important columns in the table design. aggregate_id (used as the partition key), event_type, payload, and the status or the publication time. There are two relay approaches. Polling, which periodically queries PENDING, and CDC, which reads the DB's transaction log (WAL, binlog). CDC has low latency and little DB load, but adds one more operational component.
There is one thing you must understand here. The outbox eliminates loss but cannot eliminate duplication. If the relay dies right after publishing to the broker and before changing the status to PUBLISHED, it publishes the same event again after restarting. This is at-least-once, and at this point the consumer's idempotency becomes essential. Exactly-once is not a property created by delivery but a property created by the receiving side.
A saga solves a different problem. It composes one business transaction spanning several services as a chain of local transactions and, on failure, compensating transactions. If payment succeeds and shipping fails, it executes a compensating step that cancels the payment. The key point is that a saga is not a rollback but a cancellation that goes forward — an email already sent cannot be undone, and you only send one more "cancellation notice email".
A saga has two forms: choreography (each service listens to events and takes the next action) and orchestration (a central coordinator directs the order). If there are three steps or fewer, choreography is light, and beyond four nobody knows where the flow is going, so orchestration is better.
What you meet in the field
An outbox table grows bloated if left alone. You must periodically delete published rows or drop partitions. And unless you have a partial index on the condition status='PENDING', the relay query scans the whole table.
A mistake that often appears in compensating transactions is ignoring the fact that the compensation itself can fail. Compensations must also be retried, and therefore compensations must also be idempotent.
How to design events
For the outbox and the saga to work, the events flowing over them must themselves be well designed. If this part is weak, however precisely you implement the technique, it collapses a few months later.
An event is a fact, not a command. OrderCreated is something that has already happened and cannot be undone, and what the listener does is decided by the listener. On the other hand, if you send a command like SendEmail as an event, the sender must know the receiver's circumstances, and the point of splitting the services is lost.
Decide whether to include the needed information or pass only a reference. If you put the whole order in the event, the listener does not need to ask back, but the event gets large and the values in it can go stale. Conversely, if you pass only the ID, the listener must look it up every time, so load concentrates on the original service, and if that service dies, the event cannot be processed and the benefit of splitting asynchronously is lost. The practical compromise is to include values that must be fixed as facts at that time, and to leave as references those for which you need to see the current state. It is the same criterion as the earlier normalization discussion.
Assume the format will change. Once an event is published, several services read it at their own pace, so you cannot deploy the sender and the receiver at the same time. So allow only changes that add fields, and issue changes that delete or change the meaning as an event with a new name. If you register the schema somewhere and have a mechanism that checks compatibility, this rule is followed automatically.
Include the time and the order. If you put the occurrence time and a sequence number in the event, the receiving side can recognize and ignore an old event that arrives late. This is the practical solution to the problem we mentioned earlier that retries break ordering, and without these values the receiving side comes to mistake the arrival order for the occurrence order.
What you will do in the next lab
You create the orders and outbox tables with SQLite, first reproduce the inconsistency of the dual write, then fold it into a single transaction, make a relay that publishes to Redis, confirm that duplicates arise when the relay dies midway, and then add deduplication to the consumer.