You update a row in your database, then publish an event to Kafka so other services can react. The update commits. The publish fails — the broker is down, the network blips, the process crashes between the two operations. Now the order table says one thing and every downstream consumer believes another. You have just built a dual-write bug, and it will happen on the worst possible day.
The transactional outbox pattern removes this failure mode by making the event part of the same transaction as the state change. Instead of publishing to the broker, you insert a row into an outbox table inside your business transaction. Either both the state change and the event commit, or neither does. A separate relay then reads the outbox and publishes the events. Atomicity is no longer your problem — it is the database’s problem, and databases are very good at it.
This post walks through building the pattern end to end in Go with PostgreSQL: the outbox table, the transactional writer, a polling relay with the concurrency problems that come with it, and a CDC-based alternative for when polling is not enough.
The Dual-Write Problem, Precisely
Here is the naive version that looks fine in code review and fails in production:
func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
if err := s.repo.InsertOrder(ctx, o); err != nil {
return err
}
// If this fails, the order exists but nobody was told.
return s.events.Publish(ctx, "order.created", o)
}
Retrying the publish helps but cannot fix it: if the process dies after InsertOrder commits and before Publish runs, the event is simply gone. Retrying the whole operation creates a different problem — duplicate orders. Two systems, two commits, no shared transaction. That is the whole story.
The outbox pattern collapses the two commits into one. The event becomes data, and data gets transactions for free.
The Outbox Table
CREATE TABLE outbox (
id UUID PRIMARY KEY,
aggregate_type TEXT NOT NULL,
aggregate_id TEXT NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
published_at TIMESTAMPTZ
);
CREATE INDEX idx_outbox_unpublished
ON outbox (created_at)
WHERE published_at IS NULL;
A few deliberate choices here. The payload is JSONB so consumers get a self-contained event, not a change log row. The published_at column marks completion instead of deleting rows — keeping published events around for a while gives you a free audit trail and makes "did we lose anything?" a query instead of an archaeology project. The partial index keeps the relay's scan cheap: it only ever looks at unpublished rows. A nightly cleanup job can prune old published rows so the table does not grow forever.
Writing Events Transactionally
The writer does exactly two things in one transaction — mutate business state and append to the outbox:
type OutboxEvent struct {
ID uuid.UUID
AggregateType string
AggregateID string
EventType string
Payload []byte
}
func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
tx, err := s.db.Begin(ctx)
if err != nil {
return err
}
defer tx.Rollback(ctx)
if _, err := tx.Exec(ctx,
`INSERT INTO orders (id, customer_id, total, status)
VALUES ($1, $2, $3, $4)`,
o.ID, o.CustomerID, o.Total, "created"); err != nil {
return err
}
payload, err := json.Marshal(o)
if err != nil {
return err
}
if _, err := tx.Exec(ctx,
`INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload)
VALUES ($1, $2, $3, $4, $5)`,
uuid.New(), "Order", o.ID.String(), "OrderCreated", payload); err != nil {
return err
}
return tx.Commit(ctx)
}
If the commit succeeds, both the order and the event exist. If anything fails before the commit, neither does. There is no window in between. One important discipline: the payload should be a fact ("order 42 was created with these items"), not an instruction ("send a confirmation email"). Consumers decide what a fact means. That keeps the schema of your events stable even as downstream workflows change.
The Relay: Polling Without Shooting Yourself
The relay reads unpublished rows, publishes them to the broker, and marks them published. The obvious implementation has two classic bugs, so let us build it correctly the first time.
func (r *Relay) Run(ctx context.Context) {
ticker := time.NewTicker(200 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := r.relayBatch(ctx); err != nil {
log.Printf("relay batch: %v", err)
}
}
}
}
func (r *Relay) relayBatch(ctx context.Context) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback(ctx)
// FOR UPDATE SKIP LOCKED: concurrent relays never grab the same row,
// and nobody blocks waiting for a row another relay holds.
rows, err := tx.Query(ctx, `
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED`)
if err != nil {
return err
}
defer rows.Close()
var batch []OutboxEvent
for rows.Next() {
var e OutboxEvent
if err := rows.Scan(&e.ID, &e.AggregateType, &e.AggregateID,
&e.EventType, &e.Payload); err != nil {
return err
}
batch = append(batch, e)
}
if err := rows.Err(); err != nil {
return err
}
if len(batch) == 0 {
return nil
}
for _, e := range batch {
if err := r.publish(ctx, e); err != nil {
// Stop the batch at the first failure. Rows already published
// in this transaction stay locked until rollback; they will be
// retried in order on the next tick.
return fmt.Errorf("publish %s: %w", e.ID, err)
}
}
_, err = tx.Exec(ctx,
`UPDATE outbox SET published_at = now()
WHERE id = ANY($1)`,
idsOf(batch))
if err != nil {
return err
}
return tx.Commit(ctx)
}
The two bugs this avoids:
- Publishing twice. If the process crashes after
publishbut beforeCommit, the rows roll back to unpublished and the next tick publishes them again. The outbox guarantees at-least-once delivery, never exactly-once — so every consumer must be idempotent, typically by tracking event IDs it has already processed. This is the same idempotency discipline you already apply to write APIs, moved one layer down. - Duplicate work under concurrency. Running two relay replicas with plain
FOR UPDATEmakes the second one block on the first's locks. WithSKIP LOCKED, each replica claims a disjoint slice of the queue and both make progress. It is the same primitive that powers Postgres-based job queues.
One subtlety worth naming: holding the transaction open while publishing to the broker means broker latency turns into lock duration. With SKIP LOCKED other relays just route around the slow one, so the system degrades gracefully rather than deadlocking. Keep batches small and broker timeouts tight.
Also note the ordering guarantee. Within one aggregate (one order, one account), events must usually be processed in order. ORDER BY created_at gives you that per relay run, but if you scale to many relay replicas, rows can be published by different replicas and interleave. If strict per-aggregate ordering matters, partition the work: have each relay claim only rows whose aggregate_id hashes to it, or accept that consumers must tolerate occasional reordering within an aggregate.
Polling Versus CDC
Polling is simple and good enough for most systems: latency equals your tick interval, and a couple of hundred milliseconds is rarely a problem. When you need lower latency or want to skip the relay process entirely, change data capture is the alternative. A tool like Debezium streams the database's write-ahead log, reads your outbox inserts as they commit, and applies an outbox event router transformation that converts each row into a well-formed broker message — routing by aggregate type, using the aggregate ID as the message key so per-aggregate order is preserved in partitioned topics.
The trade-offs are real. CDC gives you sub-second latency, no polling load, and no extra relay code — but you now operate Kafka Connect (or a similar pipeline) and own its configuration. Polling gives you a 200-line Go program and one more table. Start with polling. Move to CDC when latency requirements or message volume justify it, not because it is more interesting.
What the Outbox Does Not Fix
Be honest about the remaining gaps:
- Consumers still need idempotency. At-least-once delivery means duplicates will arrive. Event IDs exist precisely so consumers can dedupe.
- The outbox does not sequence across aggregates. Two different orders may be relayed in any relative order. Design consumers around per-aggregate invariants, not global ones.
- Read-your-own-writes across services stays eventual. The ordering service sees the order immediately; the notification service sees it one tick later. If a flow needs synchronous confirmation, that is an API call, not an event.
- Schema evolution is your job. JSONB payloads are flexible but untyped. Version your event types (
OrderCreatedV2) or embed a version field before the first consumer ships, not after.
The pattern pairs naturally with sagas for multi-service workflows: the outbox guarantees a service's decisions reliably become events, and a saga choreographs what happens next. If you already use Watermill or a similar Go event library, check whether it ships an outbox component before writing your own — and even then, the table schema and relay loop above are small enough that owning them is often simpler than a dependency.
Wrapping Up
Dual writes are one of those bugs that pass every test and fail every audit. The transactional outbox replaces the problem with something databases already solve: atomic commits. The implementation is deliberately boring — one table, one insert inside your existing transaction, one polling loop with FOR UPDATE SKIP LOCKED. Boring here is the feature. Start with the polling relay, make consumers idempotent from day one, and reach for CDC only when your latency budget demands it.