The Outbox Pattern for Reliable Event Publishing

Learn how the transactional outbox prevents lost events, how relays publish them, and how idempotent consumers handle duplicates in distributed systems.

published: reading time: 34 min read author: GeekWorkBench updated: June 17, 2026
Quick Summary

The transactional outbox records a business change and its event in one database transaction, then a polling relay or CDC connector publishes committed rows. This guide compares relay designs, schema and retention choices, high availability, and failure recovery, while explaining why delivery is at least once and how consumer idempotency handles duplicates. It also covers backlog monitoring, payload security, and the operational trade-offs that help teams decide when an outbox is worth running.

The Outbox Pattern for Reliable Event Publishing

Introduction

The dual-write problem appears when an application updates its database and publishes a broker message as separate operations. Either operation can succeed while the other fails, leaving services with inconsistent state and requiring reconciliation.

# A crash after the database write can lose the notification.
db.update_order(order)
broker.publish("order.updated", order)

# Store the event with the business change in one transaction.
with db.transaction():
    db.update_order(order)
    db.insert("outbox", {"event_type": "order.updated", "payload": order})

The relay publishes the committed outbox row later. This closes the gap between the database write and event recording, though delivery can be delayed and repeated.

The outbox pattern writes business data and an event record in the same database transaction, then uses a separate relay to publish the event. This guide covers that flow, the duplicate-delivery trade-off, implementation choices, and the monitoring needed to operate it safely.

Outbox Design and Relay Operations

Core Concepts

The outbox uses the database transaction as the shared boundary. Instead of publishing directly to the broker, the application writes to two tables in the same transaction: the business table and an outbox table.

def place_order(order):
    with db.transaction():
        db.insert("orders", order)
        db.insert("outbox", {
            "event_type": "order_created",
            "payload": json.dumps(order),
            "created_at": now()
        })

Both writes happen atomically. Transaction commits, both order and outbox entry exist. Transaction rolls back, neither exists.

A separate process, often called the Relay or Publisher, reads the outbox table and publishes events to the broker.

class OutboxRelay:
    def process_outbox(self):
        events = db.fetch_unpublished(limit=100)
        for event in events:
            try:
                message_broker.publish(event["event_type"], event["payload"])
                db.mark_published(event["id"])
            except Exception as error:
                log_publish_failure(event["id"], error)
                # Leave the row unpublished so the next poll can retry it.

After a successful publish, the relay marks the outbox entry as published. If the broker accepts an event but the relay does not record that success, the relay may publish it again after restarting. Consumers must handle duplicates gracefully.

Outbox Architecture Flow

An event starts its life in the database. Your application only ever touches the database—the relay handles everything downstream, which keeps your write path uncluttered.

flowchart TD
    subgraph Application["Application Layer"]
        App[Application Code]
        DB[(Business<br/>Table)]
        Outbox[(Outbox<br/>Table)]
        App --> |BEGIN TRANSACTION| T1
        T1 --> |INSERT| DB
        T1 --> |INSERT| Outbox
        T1 --> |COMMIT| Commit1[Commit]
        Commit1 --> |Atomic write| DB
        Commit1 --> |Atomic write| Outbox
    end

    subgraph Relay["Relay / Publisher"]
        Poll[Poll outbox<br/>entries]
        Fetch[Fetch unpublished]
        Publish[Publish to broker]
        Mark[Mark as published]
        Poll --> Fetch
        Fetch --> Publish
        Publish --> Mark
    end

    subgraph Broker["Message Broker"]
        Topic[Event Topic]
    end

    subgraph Consumer["Consumer Services"]
        Handler[Event Handler]
        Process[Process event<br/>idempotently]
        Handler --> Process
    end

    Outbox --> |1. Poll| Poll
    Publish --> |2. Publish event| Topic
    Topic --> |3. Deliver| Handler

Flow Steps:

Step Component Action
1 Application Writes business data + outbox entry in single transaction
2 Relay Polls outbox table for unpublished entries
3 Relay Publishes event to message broker
4 Relay Marks outbox entry as published
5 Broker Delivers event to subscribed consumers
6 Consumer Processes event idempotently

Key Properties:

  • Atomicity: Business data and outbox entry commit together or not at all
  • Reliability: Relay retries ensure eventual delivery
  • At-least-once: Events may be published multiple times if relay crashes after publish but before mark
  • Idempotency required: Consumers must deduplicate using event IDs

For more on event-driven patterns, see Event-Driven Architecture.

Design and Architecture

The cost is delay. The broker does not get the message until the relay polls, processes, and publishes. For most use cases this is fine, milliseconds or seconds. But if your consumers need immediate notification, the outbox pattern is the wrong tool.

The outbox table design determines how fast the relay can run and how painful operations become. Pick the right indexes and partitioning strategy upfront, and polling stays snappy even as the table grows. The sections below cover what actually works in production.

Why Not Just Use Transactions?

Traditional database transactions cannot span both the database and the message broker. The broker is external. You cannot begin a transaction, write to the database, write to the broker, and commit as one unit.

The outbox works around this by keeping everything in the database. The transaction boundary covers the business data and the outbox entry. The relay publishes from the outbox to the broker outside the transaction. Atomicity for the write, eventual consistency for the notification.

You lose immediate notification. The broker does not get the message until the relay polls, processes, and publishes. For low-latency requirements, this delay matters. For most use cases, milliseconds or seconds is acceptable.

Outbox Table Schema Design

The outbox table schema matters for performance. A poorly designed outbox table can become a bottleneck under high write load.

Basic Schema

CREATE TABLE outbox (
    id              UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    event_type      VARCHAR(255) NOT NULL,
    aggregate_type  VARCHAR(255),          -- e.g., 'Order', 'Payment'
    aggregate_id    VARCHAR(255) NOT NULL,   -- e.g., order-123
    payload         JSONB NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    published_at    TIMESTAMPTZ,             -- NULL = unpublished
    published       BOOLEAN NOT NULL DEFAULT false,

    -- Idempotency key for deduplication
    event_id       VARCHAR(255) UNIQUE NOT NULL
);

-- Index for relay polling: find unpublished entries efficiently
CREATE INDEX idx_outbox_unpublished
    ON outbox (created_at)
    WHERE published = false;

-- Index for aggregate lookups (useful for debugging/replays)
CREATE INDEX idx_outbox_aggregate
    ON outbox (aggregate_type, aggregate_id, created_at);

Partitioning Strategies

For high-throughput systems, partitioning the outbox table prevents it from becoming a bottleneck.

By time range (recommended for most cases):

-- Partition by day for easy TTL cleanup and efficient querying
CREATE TABLE outbox (
    id              UUID DEFAULT gen_random_uuid(),
    event_type      VARCHAR(255) NOT NULL,
    aggregate_type  VARCHAR(255),
    aggregate_id    VARCHAR(255) NOT NULL,
    payload         JSONB NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL,
    published_at    TIMESTAMPTZ,
    published       BOOLEAN NOT NULL DEFAULT false,
    event_id        VARCHAR(255) NOT NULL,
    PRIMARY KEY (id, created_at)
) PARTITION BY RANGE (created_at);

-- Create daily partitions. A partitioned-table UNIQUE constraint must include
-- the partition key, so event_id uniqueness needs a separate registry if it
-- must be enforced globally across partitions.
CREATE TABLE outbox_2026_03_01 PARTITION OF outbox
    FOR VALUES FROM ('2026-03-01') TO ('2026-03-02');
-- ... etc

Partition pruning benefits: Queries for unpublished entries in a specific time range only scan the relevant partition. Old partitions can be detached and dropped for efficient cleanup.

Avoiding Hot Partitions

Time partitioning does not remove a write hotspot in the newest partition. Consider additional sharding only after measuring that bottleneck; partitioning by status alone does not separate rows by date.

Index Design for Relay Performance

The relay needs to find unpublished entries fast. A partial index on created_at where published = false supports the critical path:

  • Relay queries: WHERE published = false ORDER BY created_at LIMIT N
  • The index helps PostgreSQL find and order candidate rows, but fetching id and event payload columns still requires table access unless the index includes them.
  • Whether PostgreSQL can use an index-only scan also depends on the table’s visibility map.
-- Verify index is used with: EXPLAIN ( buffers, analyze ) SELECT ...
SELECT id FROM outbox
WHERE published = false
ORDER BY created_at
LIMIT 100;

Schema Design Summary

Concern Recommendation
Primary key UUID (no sequence contention)
Event ID Business meaningful ID (order-123-created), not random UUID
Payload JSONB (flexible, queryable)
Index for relay (created_at) WHERE published=false
Partitioning By time range (daily partitions)
Idempotency UNIQUE constraint on event_id

Different approaches work in different situations.

Polling Relay

The simplest approach polls the outbox table at regular intervals.

class PollingRelay:
    def __init__(self, db, broker, interval_ms=100):
        self.db = db
        self.broker = broker
        self.interval = interval_ms / 1000

    def run(self):
        while True:
            events = self.db.fetch_unpublished(limit=100, order="created_at")
            for event in events:
                self.publish_and_mark(event)
            sleep(self.interval)

    def publish_and_mark(self, event):
        try:
            self.broker.publish(event["event_type"], event["payload"])
            self.db.mark_published(event["id"])
        except Exception as error:
            self.log_publish_failure(event["id"], error)
            # Leave the row unpublished for a later retry.

Polling is simple and reliable. The tradeoff is latency. With a 100ms interval, events publish within 100ms at best. Reduce the interval and database load goes up.

Change Data Capture

A more sophisticated approach uses change data capture (CDC), which reads committed changes from the database log. Debezium can capture inserts in an outbox table and route them to a broker. PostgreSQL LISTEN/NOTIFY can wake a relay so it can query for pending rows, but it is a notification mechanism, not a durable change log; the relay still needs a recovery path that checks the table.

class CdcRelay:
    def __init__(self, db, broker):
        self.db = db
        self.broker = broker
        self.db.listen("outbox_inserts", self.handle_notification)

    def handle_notification(self):
        events = self.db.fetch_unpublished(limit=100)
        for event in events:
            self.publish_and_mark(event)

CDC cuts latency significantly. Events publish seconds or milliseconds after commit rather than on the next poll cycle. The cost is operational complexity. CDC tools add infrastructure and their own failure modes.

Polling vs CDC vs Transactional Outbox Comparison

Aspect Polling Relay CDC Relay Database Trigger + Relay
Latency Interval-based (100ms-1s typical) Near real-time (ms-s) Low (within transaction commit)
Database Load Higher (continuous polling) Low (event-driven) Lowest (no polling)
Infrastructure Simple process CDC tool (Debezium, AWS DMS) Trigger/procedure plus delivery process
Reliability High (simple, testable) High (but CDC tool has own failure modes) Highest (DB handles all)
Operational Complexity Low High (CDC tool management) Medium (DB procedures and relay)
Replay Capability Yes (mark-based) Yes (CDC log retention) Limited (WAL-based)
Cross-Database Support Yes Limited (CDC per-DB) No (DB-specific)
When to Use Low throughput, simple setup High throughput, low latency needs Atomic DB-side outbox writes

Deleting Instead of Marking

Some implementations delete outbox entries after successful publish rather than marking them.

def publish_and_mark(event):
    self.broker.publish(event["event_type"], event["payload"])
    self.db.delete_outbox_entry(event["id"])

This keeps the outbox table small. The cost is losing replay capability. Marking as published preserves the data and lets you replay from the beginning if needed.

Handling Failures and Idempotency

The outbox supports at-least-once delivery when the relay keeps retrying unpublished rows. If a publish succeeds but marking the row fails, the event may be published again. Consumers must handle duplicates.

Idempotent consumers solve this. The consumer tracks which events it has already processed, using an event ID in a processed_events table or cache.

class IdempotentConsumer:
    def __init__(self, broker, db):
        self.broker = broker
        self.db = db

    def handle_order_created(self, payload):
        event_id = payload["event_id"]
        order = payload["order"]
        with self.db.transaction():
            if not self.db.insert_processed_if_absent(event_id):
                return  # Skip duplicate
            self.process_order(order)

The idempotency table needs its own cleanup to avoid growing forever. Purge processed events older than a certain age, or use a sliding window.

Idempotent Consumer Effects

The outbox provides at-least-once delivery, not exactly-once message transmission. A consumer can make its local effect idempotent by recording the event ID and applying the business change in the same database transaction. If the event arrives again, the recorded ID lets the consumer skip that duplicate. Effects outside that transaction, such as sending an email or calling another service, need their own retry and deduplication strategy.

For a broader view of delivery semantics, see Saga Pattern which often pairs with outbox for distributed transaction handling.

Operational Considerations

The outbox table grows if the relay falls behind or publishes repeatedly fail. Monitor for stuck entries.

SELECT COUNT(*) FROM outbox
WHERE published = false
AND created_at < NOW() - INTERVAL '5 minutes';

Alert on this query catches relay failures or persistent publish errors.

Consider the relay’s failure modes too. If the relay crashes mid-batch, some events get published but not marked. The next relay instance reprocesses them. As long as consumers are idempotent, this is fine.

For high-throughput systems, batching helps. The relay fetches multiple events, publishes them in a single broker batch if supported, then marks them all published. Fewer broker round trips.

Relay High Availability

A single relay is a single point of failure. If it crashes, events pile up in the outbox until someone notices. For production systems, you need multiple relay instances.

Active-passive (leader election):

Use a lock in the database or Redis to elect an active relay. The active relay polls and publishes. Passive relays standby, ready to take over if the active dies.

import redis
import threading

class HaRelay:
    def __init__(self, db, broker, lock_key, ttl_seconds=30):
        self.db = db
        self.broker = broker
        self.lock_key = lock_key
        self.ttl = ttl_seconds
        self.redis = redis.Redis()

    def run(self):
        while True:
            acquired = self.redis.set(self.lock_key, socket.gethostname(),
                                       nx=True, ex=self.ttl)
            if acquired:
                self.do_work()  # Only active relay does work
            else:
                time.sleep(1)   # Standby, wait for lock release

    def do_work(self):
        while True:
            try:
                events = self.db.fetch_unpublished(limit=100)
                for event in events:
                    self.broker.publish(event["event_type"], event["payload"])
                    self.db.mark_published(event["id"])
            except redis.ConnectionError:
                break  # Lost lock, go back to standby

The nx=True means only one instance can hold the lock. The ex=self.ttl means the lock auto-releases if the active relay crashes and doesn’t renew. Passive instances detect the lock is gone and one of them acquires it.

Active-active (multiple relays with sharding):

Multiple relays can run simultaneously if they coordinate which outbox rows each one handles. Use SELECT ... FOR UPDATE SKIP LOCKED to let each relay claim a different batch:

class ShardedRelay:
    def process_outbox(self):
        # FOR UPDATE SKIP LOCKED: grab rows without blocking other relays
        events = db.fetch_unpublished(
            "SELECT * FROM outbox WHERE published = false "
            "ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED"
        )
        for event in events:
            try:
                self.broker.publish(event["event_type"], event["payload"])
                self.db.mark_published(event["id"])
            except Exception as error:
                self.log_publish_failure(event["id"], error)
                # Leave the row unpublished for a later retry.

Both relays run concurrently, each grabbing different rows. No coordination overhead, no leader election needed. Throughput scales with relay count.

Outbox Entry TTL and Cleanup

The outbox table grows over time. Published entries that are never cleaned up become dead weight. You need a cleanup strategy.

Soft delete with published flag (already in schema):

# Mark as published, don't delete
db.mark_published(event["id"])

Hard delete for old published entries (run as a scheduled job):

-- Clean up published entries older than 7 days
DELETE FROM outbox
WHERE published = true
AND published_at < NOW() - INTERVAL '7 days';

Automated cleanup with partition expiry (if using time-based partitioning):

-- Detach and drop old partitions
ALTER TABLE outbox DETACH PARTITION outbox_2026_02_01;
DROP TABLE outbox_2026_02_01;

Partition pruning makes cleanup essentially free — dropping a partition is O(1) rather than deleting millions of rows.

Retention policy guidance:

Use case TTL for published entries
Financial transactions 30-90 days (audit requirements)
Order events 7-14 days
High-volume event logs 1-3 days
Debug/replay scenarios Keep all for 24h, then batch archive

Never delete unpublished entries — they represent events that haven’t been delivered yet. If you see unpublished entries older than your poll interval, something is wrong with the relay.

Combining with Saga

The outbox pattern pairs naturally with saga. In a saga, each step updates a service and publishes an event to trigger the next step. Without outbox, the publish can fail and leave the saga stuck.

With outbox, the step’s update and outbox write are atomic. The relay publishes the event that triggers the next step. If the relay is down, events queue up and publish when it recovers. The saga makes progress as soon as the relay is available.

For details on saga implementation, see Saga Pattern.

The core problem outbox solves for saga is broker unavailability mid-transaction. In a distributed transaction spanning multiple services, if Step 2 publishes its “step complete” event but the broker is down, Step 3 never fires. The saga hangs. With outbox, Step 2 writes its completion and the outbox entry atomically. The relay retries until the broker accepts the event. Step 3 fires when the relay recovers.

The failure semantics differ from single-service patterns. In a buy-and-sell saga, each step is atomic within its service, but the cross-service coordination depends on events surviving broker outages. If the relay is down for 30 seconds during a 5-step saga, steps queue up and fire in sequence once the relay returns. The saga completes, just with added latency.

Compensations add complexity in outbox-based sagas. A compensation for Step 3 might execute after Step 4 has already run, depending on relay timing. Design compensations to be idempotent and aware of subsequent step state. Forward-only sagas where each step is irreversible work best with this approach. Sagas that need tight compensation ordering need additional coordination logic beyond outbox alone.

When to Use and When Not to Use the Outbox Pattern

Use the outbox when a committed database change must eventually produce an event, and the business write and event record can share one database transaction. This atomic boundary prevents a crash between the two writes from silently losing a notification. The relay publishes after commit, so consumers may see the event hundreds of milliseconds or seconds later.

It fits high-reliability workflows such as payments, inventory changes, and event-driven services where a missed event leaves another service stale. It is not justified for best-effort notifications, disposable telemetry, or a simple workflow that does not publish domain events. Skip it when the event must be visible synchronously with the database write; the relay cannot provide that guarantee.

The cost is a relay to deploy and monitor, retained outbox rows, and idempotency in consumers. If those costs exceed the harm from a missed notification, a direct best-effort publish may be enough.

For an overview of distributed transaction patterns, see Distributed Transactions.

CQRS and Event Sourcing

CQRS and event sourcing hit the same reliability wall that outbox was built to crack. In CQRS, getting events from the write side to the read side reliably is the whole challenge. Event sourcing has the same problem: your event store needs to tell consumers whenever new events land. Without outbox, both patterns end up relying on direct writes or polling that can quietly break and leave your systems out of sync.

Outbox fixes this. The event store and outbox entry commit together, so events don’t vanish between appending to the log and notifying consumers. Your consumers stay decoupled from your storage layer. Change the schema without worrying about breaking downstream processors.

The sections below dig into each pattern separately, then set them against direct writes and retry queues head-to-head. Five failure scenarios follow to show how outbox actually behaves when production decides to misbehave.

CQRS (Command Query Responsibility Segregation)

In CQRS, you separate read models from write models. Commands (writes) update the command side. Queries (reads) read from potentially different read models updated via events.

The outbox pattern is how you get events from the command side to the query side reliably:

Command Side                          Query Side
┌─────────────────┐                  ┌─────────────────┐
│ Order Service   │   outbox + relay  │ Read Model     │
│ - place_order  │ ───────────────► │ - OrderView    │
│ - update_status│   (reliable       │ - CustomerView │
│ - cancel       │    event stream)  │ - Analytics DB │
└─────────────────┘                  └─────────────────┘

Without outbox, you might update the command model and then directly update the read model. That direct call can fail and leave the systems out of sync. With outbox, the update and event publication are atomic — the read model eventually catches up, never misses an event.

Event Sourcing

Event sourcing takes CQRS further: instead of storing current state, you store the sequence of events that produced the current state. The event store is the source of truth.

The outbox pattern works well with event sourcing:

def place_order(order):
    with db.transaction():
        # Append to event log (the "event store")
        db.insert("events", {
            "event_type": "ORDER_PLACED",
            "aggregate_id": order.id,
            "payload": json.dumps(order),
            "sequence": get_next_sequence("Order", order.id)
        })

        # Also write to outbox for reliable notification
        db.insert("outbox", {
            "event_type": "ORDER_PLACED",
            "aggregate_id": order.id,
            "payload": json.dumps(order),
            "event_id": f"order-{order.id}-placed"
        })

The event store is the authoritative log. The outbox ensures events propagate to consumers (read models, audit logs, other services). If you only use event sourcing without outbox, you rely on consumers to read the event store directly — which couples them to your storage mechanism and makes it harder to evolve the event schema.

The outbox decouples: your event store is an internal implementation detail. Consumers receive events through the broker like any other event-driven system.

Trade-Off Table

Dimension Outbox Pattern Direct Write (Dual Write) Retry Queue
Consistency Strong (atomic within DB transaction) Weak (split across systems) Medium (retries eventually succeed)
Complexity Medium (relay + idempotency) Low (direct publish) Medium (queue management)
Latency Higher (poll interval delay) Low (immediate) Medium (queue processing)
Failure Handling Self-healing (relay retries) Requires manual compensation Built-in (queue retries)
Operational Burden Medium (monitor relay, outbox size) Low (no additional processes) Medium (queue infrastructure)
Replay Capability Yes (marked entries preserved) No Limited (queue retention)
Duplicate handling Consumer deduplication required Consumer deduplication required Consumer deduplication required
Infrastructure Cost Medium (DB + relay process) Low (just broker) High (queue cluster)
Scalability High (sharded relays) Medium (broker scales) High (queue scales)

Production Failure Scenarios

Scenario 1: Network Partition Between Relay and Broker

What happens: The broker accepts an event, but the acknowledgement is lost before it reaches the relay. The relay cannot know whether the publish succeeded, so it retries. The broker may then deliver the same event twice.

Impact: If consumers are not idempotent, they may process the same order twice.

Mitigation: Design consumers to be idempotent. Use the event ID to deduplicate.

Scenario 2: Database Crash During Outbox Write

What happens: The application starts a transaction, inserts the business data, inserts the outbox entry, but crashes before the transaction commits. The database rolls back automatically. No event published. No business data written. System is consistent.

Impact: None. This is actually correct behavior — the crash prevented an inconsistent state.

Mitigation: None needed. This is why atomic transactions matter.

Scenario 3: Relay Crash After Publish Before Mark

What happens: Relay publishes 100 events to the broker. Broker accepts all 100. Before the relay can mark any as published, it crashes. On restart, the relay re-publishes all 100 events.

Impact: Consumers receive 200 events total (100 original + 100 duplicates). If consumers track processed event_ids, duplicates are discarded harmlessly. If not, you have double-processing.

Mitigation: Idempotent consumers with deduplication. Consider using transactional sentry patterns where the mark is part of the same transaction as the publish acknowledgment.

Scenario 4: Outbox Table Grows Unbounded

What happens: The relay is down for maintenance. Business continues. Outbox entries pile up. After a week, millions of rows exist. The relay starts up and struggles to catch up while the table scan degrades database performance.

Mitigation: Monitor unpublished event count and age. Set up alerts for stuck relays. Consider partitioning by time and dropping old partitions instead of deleting rows. Archive cold data to a separate table before dropping partitions.

Scenario 5: Payload Schema Change Breaks Consumers

What happens: A new version of the application changes the event payload schema. Old consumers that rely on specific fields break when they receive the new format.

Mitigation: Use versioned event types (e.g., order_created_v1, order_created_v2). Maintain backward compatibility within the same major version. Implement schema registry and validation at the relay layer.

Scenario 6: Poison Row Blocks Relay Progress

What happens: A malformed or permanently rejected event fails on every retry. If the relay stops the batch at that row, newer events can remain unpublished behind it.

Mitigation: Track attempts, alert on the row’s age, and move repeatedly failing events to a reviewed dead-letter or quarantine path. Preserve the event ID and failure details so operators can repair and replay it; do not silently mark it published.

Quick Recap Checklist

  • Business data and outbox entry written in same database transaction
  • Relay process runs separately from application
  • Consumers are idempotent and deduplicate by event_id
  • Outbox table has appropriate indexes for relay polling
  • Published entries are cleaned up on a schedule
  • Unpublished entry count and age are monitored
  • Relay is deployed in HA configuration (active-passive or sharded)
  • Event payload does not contain sensitive data (PII, passwords, tokens)

Observability Checklist

The outbox pattern introduces a separate relay process and additional database writes. Without monitoring, problems in the relay go unnoticed until events pile up and consumers fall behind.

Metrics

  • Outbox lag: count and age of unpublished entries
  • Relay lag: time from row creation to successful publish, split by relay or shard
  • Oldest unpublished row age and retry count, so one poison row cannot hide behind a healthy average
  • Relay publish rate: events published per second
  • Publish failure rate: events that failed to publish
  • End-to-end latency: time from outbox write to consumer delivery
  • Consumer lag: how far behind consumers are from producer rate
  • Batch size histogram: how many events per publish cycle

Key Queries

-- How many unpublished events are piling up?
SELECT COUNT(*) FROM outbox WHERE published = false;

-- How old is the oldest unpublished event? (should be < 1 minute)
SELECT MAX(NOW() - created_at) FROM outbox WHERE published = false;

-- How many published events per minute?
SELECT DATE_TRUNC('minute', published_at), COUNT(*)
FROM outbox WHERE published = true
GROUP BY 1 ORDER BY 1 DESC LIMIT 10;

Logs

  • Log each batch of events published with count and duration
  • Log publish failures with error details and retry count
  • Log when relay acquires or loses leadership lock
  • Include event_id, aggregate_type, and aggregate_id in all logs for correlation

Alerts

  • Alert when unpublished event count exceeds threshold (relay is behind)
  • Alert when oldest unpublished event age exceeds threshold (relay is stuck)
  • Alert when one event exceeds its retry limit or repeatedly fails, and expose the event ID for investigation
  • Alert when publish failure rate exceeds baseline
  • Alert when relay loses leadership and standby doesn’t pick up (HA failure)

Tracing

from opentelemetry import trace

tracer = trace.get_tracer(__name__)

class OutboxRelay:
    def process_batch(self, events):
        with tracer.start_as_current_span("outbox.publish_batch") as span:
            span.set_attribute("batch.size", len(events))

            for event in events:
                with tracer.start_as_current_span("outbox.publish_event") as event_span:
                    event_span.set_attribute("event.id", event["event_id"])
                    event_span.set_attribute("event.type", event["event_type"])
                    event_span.set_attribute("aggregate.type", event["aggregate_type"])

                    try:
                        self.broker.publish(event["event_type"], event["payload"])
                        self.db.mark_published(event["id"])
                        event_span.set_status(trace.Status.OK)
                    except Exception as e:
                        event_span.record_exception(e)
                        event_span.set_status(trace.Status.ERROR)
                        raise

Security and Compliance Notes

The outbox pattern writes event payloads to a database table that a separate process reads and publishes. Misconfigured security can leak sensitive data through the event stream.

  • Encrypt outbox table data at rest (database-level encryption or application-level encryption of sensitive payload fields)
  • Give the application permission to insert outbox rows and the relay only the read, claim, and publish-state permissions it needs; avoid broad table or broker-admin access
  • Encrypt relay-to-broker communication (TLS on broker connection)
  • Authenticate relay-to-broker connections (broker authentication)
  • Authorize which relays can publish to which topics (broker ACLs)
  • Validate payload schema before publishing (reject payloads that are too large or malformed)
  • Audit log all publish operations with event metadata
  • Do not put sensitive data (passwords, tokens, PII) directly in event payloads; use references or obfuscated IDs instead
  • Set payload retention from replay and audit requirements, then delete or redact expired sensitive data from outbox storage and backups according to policy
  • Rate-limit the relay’s publish rate to prevent accidental or malicious overload of consumers

Common Pitfalls / Anti-Patterns

  • Deleting unpublished rows during cleanup can permanently lose events. Retention jobs should only remove rows that are safely published or explicitly quarantined.
  • A cleanup window shorter than the recovery or replay window can remove events operators still need. Set retention from the incident recovery and audit requirements, not just table size.
  • Parallel relays can publish events for one aggregate out of order unless the claim and partitioning strategy preserves per-aggregate order. Keep sequence numbers and make consumers reject or buffer gaps where order matters.
  • Retrying a poison row forever can stall a batch and inflate lag. Bound retries, alert on the row, and provide a deliberate repair and replay path.
  • A healthy average publish latency can hide one stale row. Alert on the oldest unpublished event as well as backlog count and average lag.

Interview Questions

1. What problem does the outbox pattern solve, and why is it needed?
The outbox pattern solves the dual-write problem, where an application must write to a database and publish a message to a broker as two separate operations. Without outbox, either operation can fail after the other succeeds, leaving the system inconsistent. The outbox pattern uses the database transaction as a single atomic boundary—both the business data and the outbox entry commit together, or neither does.
2. How does the relay/publisher component work in the outbox pattern?
The relay is a separate process that polls the outbox table for unpublished entries, publishes them to the message broker, and marks them as published. It runs independently from the application, so if it crashes, the application continues to write to the database normally. On restart, the relay picks up where it left off and publishes any missed events.
3. What delivery semantics does the outbox pattern provide, and why?
The outbox pattern provides at-least-once delivery. Events may be published more than once if the relay cannot record a successful publish. Consumers should deduplicate by event ID. If the deduplication record and local business change commit in one transaction, that local effect can be applied once even though the message may arrive repeatedly.
4. What is idempotency, and why is it critical for consumers in the outbox pattern?
Idempotency means processing the same event multiple times produces the same result as processing it once. Consumers track which event_ids they have already processed, typically in a processed_events table or cache. When a duplicate event arrives, the consumer skips it. Without idempotency, duplicate delivery would cause double-processing bugs.
5. What is the key index design for the outbox table, and why?
The partial index on `created_at` where `published = false` helps the relay find and order pending rows. It does not cover a query that also selects `id` or payload columns, so PostgreSQL must fetch those values from the table. An index-only scan is possible only when the needed columns are in the index and visibility information allows it.
6. How does CDC (Change Data Capture) compare to polling relay for outbox publishing?
Polling runs interval-based queries and is straightforward, with latency tied to the poll interval. PostgreSQL LISTEN/NOTIFY can wake a relay to query pending rows, while a tool such as Debezium reads database-log changes and tracks connector progress. Notifications alone are not durable CDC, and both approaches need recovery and duplicate handling. CDC adds operational complexity but can reduce polling delay.
7. What are the main approaches for relay high availability?
Two approaches: active-passive with leader election (using database locks or Redis) where one relay instance is active and others standby, and active-active with sharding (using SELECT FOR UPDATE SKIP LOCKED) where multiple relays claim different batches concurrently. Active-passive is simpler; active-active scales throughput with relay count but requires more coordination logic.
8. Why might you partition the outbox table, and what partitioning strategy do you recommend?
Partitioning can help with retention and queries when they align with the partition key, but it is not an automatic fix for high write load. Time-range partitions make old published data easier to remove, while unpublished rows must remain available for relay recovery. Choose the interval from write volume and retention needs, and measure before adding another partitioning dimension.
9. How does the outbox pattern combine with the Saga pattern?
Saga orchestrates distributed transactions as a sequence of steps, where each step publishes an event to trigger the next step. Without outbox, a publish failure can leave the saga stuck mid-flight. With outbox, each step's update and outbox write are atomic, so the event that triggers the next step is guaranteed to be published when the relay is available. The saga makes progress as soon as the relay recovers.
10. What happens when the relay falls behind or crashes for an extended period?
Unpublished entries accumulate in the outbox table. The application continues normally, writing both business data and outbox entries. When the relay restarts, it processes the backlog. Problems occur if the backlog grows very large (performance degradation on large table scans) or if consumers have moved on (stale events). Monitoring unpublished count and age catches this early. Partitioning helps by making cleanup efficient.
11. What are the main operational concerns when running the outbox pattern in production?
Key concerns: monitoring outbox lag (unpublished count and age), ensuring relay HA (multiple instances or leader election), managing outbox table growth through partitioning and cleanup policies, handling payload schema evolution across versions, and ensuring consumers are idempotent. Also important: alerting on stuck relays, logging publish batches and failures, and tracing event flow end-to-end.
12. When should you NOT use the outbox pattern?
The outbox pattern adds operational complexity (relay process, monitoring, idempotency logic) and introduces latency (poll interval or CDC delay). Do not use it when: missed events are acceptable (low-value notifications), immediate broker notification is required (low-latency use cases), your reliability requirements are low, or the added complexity is not justified by the business impact. The tradeoff is simplicity versus guaranteed delivery.
13. How does outbox pattern connect to CQRS and event sourcing?
In CQRS, the outbox pattern provides reliable event propagation from command side to query side (read models). Without outbox, updating a read model directly can fail and leave systems out of sync. In event sourcing, the outbox decouples the event store (internal implementation) from consumers, which receive events through the broker like any event-driven system. The outbox pattern enables both patterns to evolve independently from their consumers.
14. What is the difference between marking as published versus deleting outbox entries?
Marking as published preserves the outbox entry for replay capability. If you need to reprocess events from the beginning (new consumer, debugging, disaster recovery), the marked entries support that. Deleting keeps the table small but loses replay capability—you cannot go back and reprocess past events. The tradeoff is storage and maintenance overhead versus replay flexibility.
15. How do you handle payload schema changes in an outbox-based system?
Use versioned event types (e.g., order_created_v1, order_created_v2). Maintain backward compatibility within major versions. Implement schema validation at the relay layer before publishing. Consider using a schema registry to track versions. When possible, add new fields rather than removing old ones. Consumers should ignore unknown fields to handle forward compatibility.
16. What security considerations are specific to the outbox pattern?
Event payloads in the outbox table can be read by the relay and delivered to brokers, so sensitive data (PII, passwords, tokens) must not appear directly in payloads—use references or obfuscated IDs instead. Additional concerns: encrypt outbox table at rest, TLS on relay-to-broker connections, broker authentication, topic ACLs to authorize which relays can publish where, and payload schema validation to reject malformed or oversized payloads.
17. What monitoring queries should you run against the outbox table in production?
Essential queries: count of unpublished entries (alert if above threshold), age of oldest unpublished entry (alert if above 1-2 minutes), published events per minute (throughput metric), and published_at age distribution (histogram of processing times). Log each batch publish with count and duration, and log relay leadership changes in HA setups. Include event_id, aggregate_type, and aggregate_id in all logs for correlation.
18. What is FOR UPDATE SKIP LOCKED and why is it useful for outbox relays?
FOR UPDATE SKIP LOCKED is a PostgreSQL clause that lets concurrent transactions claim different rows without blocking or deadlocking. In an active-active relay setup, each relay instance uses this clause to grab a batch of unpublished outbox rows atomically. Each relay gets a different batch, processes and marks its own, and concurrency scales without central coordination.
19. How do database triggers or stored procedures fit with the outbox pattern?
A database trigger or stored procedure can write an outbox row as part of the business transaction; it does not safely publish directly to an external broker as part of that transaction. A polling relay or CDC connector still handles delivery. Database-side logic can keep the write atomic but may be less portable and harder to evolve than application code.
20. What are the retention policy considerations for outbox entries?
Retention depends on use case: financial transactions typically need 30-90 days for audit compliance; order events typically 7-14 days; high-volume logs 1-3 days; debug/replay scenarios keep all for 24 hours then archive. Never delete unpublished entries—they represent undelivered events that indicate relay problems. Use partition-based cleanup for efficiency (drop old partitions rather than DELETE rows).

Further Reading

Conclusion

The outbox pattern keeps a database change and its event record in the same transaction, then lets a relay publish committed events to the broker. That makes event loss recoverable, but delivery can be delayed or repeated, so consumers need idempotency and operators need to watch relay lag and outbox growth.

The pattern adds relay infrastructure, monitoring, and cleanup work. In return, a committed event stays available for retry if the broker is unavailable. The outbox itself does not guarantee delivery if the relay never recovers, so alerts on unpublished-event age and count are part of the design.

Whether that tradeoff makes sense depends on the cost of losing an event and the operational capacity of the team running the relay.

Category

Related Posts

TCC: Try-Confirm-Cancel Pattern for Distributed Transactions

Learn the Try-Confirm-Cancel pattern for distributed transactions. Explore how TCC differs from 2PC and saga, with implementation examples and real-world use cases.

#distributed-systems #transactions #saga

Distributed Transactions: ACID vs BASE Trade-offs

Explore distributed transaction patterns: ACID vs BASE trade-offs, two-phase commit, saga pattern, eventual consistency, and choosing the right model.

#distributed-systems #transactions #consistency

Choosing RPC, gRPC, or Event APIs for Service Integration

Compare REST-style RPC, gRPC, and event APIs by interaction pattern, latency, coupling, tooling, and operational needs before choosing an integration style.

#api-design #grpc #rpc