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.
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
idand 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
Further Reading
- Event-Driven Architecture - Patterns and principles for building event-driven systems
- Saga Pattern - Distributed transaction orchestration with outbox
- Distributed Transactions - Overview of transactional patterns across systems
- Debezium Documentation - CDC tools for reliable outbox publishing
- PostgreSQL NOTIFY documentation - Notification behavior and queue limits
- PostgreSQL table partitioning - Partitioned-table constraint limitations
- Debezium Outbox Event Router - CDC routing for outbox tables
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 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.
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.