On this page
concept

Outbox Pattern

Created 2026-08-17 32 connections

Outbox Pattern

The Transactional Outbox Pattern is an architectural solution that solves the dual-write problem in event-driven distributed systems — including ecommerce order management, inventory sync, and checkout pipelines — by making the application database the single source of truth for both business state and outbound events, eliminating inconsistencies caused by writing to two systems in sequence.

The problem it solves

In an ecommerce order service, when a customer places an order the service must: (1) save the order to the database, and (2) publish an OrderCreated event so that inventory, shipping, and analytics services can react. Doing both naively with sequential writes creates silent failure scenarios (Conduktor, 2026-07-30).

Streamkap (2026-02-25) documents three concrete failure modes:

  • Database succeeds, Kafka fails — order saved but inventory never reserved and no confirmation email sent
  • Kafka succeeds, database fails — downstream services start processing an order that does not exist in the database
  • Network partition between the two writes — either scenario depending on where the partition lands

Chris Richardson (microservices.io, copyright 2026) states: "Without using 2PC, sending a message in the middle of a transaction is not reliable. There's no guarantee that the transaction will commit. Similarly, if a service sends a message after committing the transaction there's no guarantee that it won't crash before sending the message."

Why naive fixes fail

Confluent (2024-05-29) identifies three failed approaches:

  • Reversing the order (Kafka first, then DB) just shifts the failure mode — a successful Kafka write followed by a DB failure produces an event with no corresponding database record
  • Wrapping both in a DB transaction does not help — if the Kafka write succeeds but the DB transaction rolls back, the event is already emitted and is not included in the rollback
  • In-memory retries fail when the application crashes because the event is no longer in memory; a durable retry requires writing to durable storage, which reintroduces the dual-write problem

2PC (two-phase commit) is also ruled out: database and message broker transactions cannot span systems atomically, and even where supported it is undesirable due to coupling and performance overhead (microservices.io, copyright 2026).

How the pattern works

The key mechanism (Conduktor, 2026-07-30): the service writes the event record INTO the same database transaction that updates the business entity. A separate relay process then reads from this outbox table and publishes to the message broker.

The key insight is that the database becomes the single source of truth. If the transaction commits, both the business data and the event are persisted. If it rolls back, neither is saved. (Conduktor, 2026-07-30)

Conduktor (2026-07-30) describes the pattern as transforming "the dual-write problem into a single-write problem by treating event publishing as part of the database transaction."

Conduktor (2026-07-30) documents the production schema:

ColumnTypePurpose
idUUID (default gen_random_uuid())Primary key
aggregate_typeVARCHAR NOT NULLEntity type ("Order", "Customer") — used for topic routing
aggregate_idVARCHAR NOT NULLEntity identifier — used as Kafka message key for partition ordering
event_typeVARCHAR NOT NULLEvent name ("OrderCreated", "OrderShipped")
payloadJSONBEvent payload

James Carr (2026-01-15) additionally recommends:

  • sequence_id BIGSERIAL rather than created_at for relay ordering — concurrent transactions can produce identical or out-of-order timestamps
  • partition_key column to ensure related events land on the same Kafka partition
  • aggregate_version for per-aggregate ordering
  • schema_version field in event payloads to support multiple schema versions simultaneously

Debezium v3.6 docs (2026) document the same default columns (id, aggregatetype, aggregateid, type, payload) and note that UPDATE operations on the outbox table trigger configurable warn/error/fatal behaviour, while DELETE operations are automatically filtered out — all changes are expected to be INSERT operations only.

Relay mechanisms

1. Polling publisher

A scheduled background job queries SELECT * FROM outbox WHERE published_at IS NULL ORDER BY sequence_id LIMIT 100 FOR UPDATE SKIP LOCKED — FOR UPDATE SKIP LOCKED prevents duplicate processing by concurrent relay instances (James Carr, 2026-01-15).

Drawbacks (Streamkap, 2026-02-25): polling interval creates a latency floor; constant DB queries even when idle; crash between publish and mark causes re-publication. Gunnar Morling at Confluent Current Bengaluru 2025 (via Speaker Deck, 2025-03-19) identifies two structural failure modes: ordering problems and missing events.

Operational hazard (YouTube source, 2025-05-13, snippet-level confidence): multiple instances of the scheduled task may run at once, producing duplicate messages downstream, unless a distributed lock such as ShedLock is applied to the relay process.

2. Change Data Capture (CDC) via Debezium + WAL (preferred)

CDC tools monitor the database transaction log — the write-ahead log (WAL) in PostgreSQL or binary log (binlog) in MySQL — which records every committed transaction; a CDC tool reads these logs and publishes changes to Kafka in near real-time (Streamkap, 2026-02-25).

Advantages over polling (Streamkap, 2026-02-25):

  • Near-zero latency (events within milliseconds of DB commit) (as-of 2026-02-25)
  • Zero additional database query load
  • Natural transaction ordering preserved
  • No polling logic to write or maintain

PostgreSQL prerequisites for Debezium CDC (Conduktor, 2026-07-30): wal_level = logical in postgresql.conf, a replication slot, and a database user with REPLICATION and SELECT privileges on the outbox table. (as-of 2026-07-30)

MySQL prerequisites (Conduktor, 2026-07-30): log_bin = ON, binlog_format = ROW, and a user with REPLICATION SLAVE and REPLICATION CLIENT privileges. (as-of 2026-07-30)

Key insight on housekeeping (Gunnar Morling, Confluent Current Bengaluru 2025): with log-based CDC, outbox table row cleanup is "Not Needed Actually" — Debezium reads from the WAL directly, so rows can be deleted immediately after insert without risk of losing events.

PostgreSQL-only optimisation (Morling, 2025): pg_logical_emit_message() can be used to emit logical replication messages from within a transaction without writing to an outbox table at all, reducing database write amplification.

Debezium Outbox Event Router SMT

Debezium (v3.6, 2026) ships a dedicated Outbox Event Router SMT (Single Message Transform) that:

  1. Routes events to Kafka topics based on aggregate_type
  2. Sets the Kafka message key to aggregate_id, ensuring per-entity ordering
  3. Sets the message value to payload
  4. Optionally deletes the outbox row after publication (as-of 2026)

The SMT supports distributed tracing via a tracingspancontext field and the debezium-read operation name. Default serialization is JSON; Apache Avro is supported as an alternative for schema governance.

MongoDB has its own separate SMT (outbox.MongoEventRouter) because the standard router is not compatible with the MongoDB connector (Debezium v3.6, 2026).

As of May 2026, Debezium deployment modes include: Kafka Connect, Debezium Server (Kinesis/Pub Sub/Pulsar/RabbitMQ/Redis Streams — no Kafka required), Debezium Platform (operator-based lifecycle management), Quarkus Extensions, Spring Integration, Debezium Engine, and PyDebezium for Python teams (Debezium blog, 2026-05-22).

Delivery guarantees

The pattern provides at-least-once delivery, not exactly-once. The message relay might publish a message more than once if it crashes after publishing but before recording that fact, so message consumers must be idempotent (microservices.io, copyright 2026).

James Carr (2026-01-15) documents three consumer idempotency strategies:

  1. Include a unique event_id + business-level idempotency_key in the event payload, tracked in Redis with a TTL
  2. Design operations to be naturally idempotent — use absolute SET values rather than relative increments
  3. Use DB unique constraints on a processed_events table with INSERT ... ON CONFLICT DO NOTHING

Gunnar Morling (Current Bengaluru 2025) adds: for identifying duplicates, use an increasing, unique value such as the PostgreSQL WAL LSN (Log Sequence Number) — a database sequence number is explicitly called out as insufficient.

Ordering guarantees

CDC preserves the order in which rows were inserted into the outbox table. Kafka preserves order within a partition. Since Debezium uses aggregate_id as the message key, all events for the same entity (e.g. the same order) land on the same Kafka partition, preserving per-entity ordering. Cross-entity ordering is not guaranteed — Conduktor (2026-07-30) states this is correct behaviour (Streamkap, 2026-02-25; Conduktor, 2026-07-30).

Conduktor (2026-07-30) notes CDC maintains ordering within a single table but not across tables, and recommends a single outbox table per aggregate root or sequence numbers for inter-aggregate ordering needs.

Known operational pitfalls

  • Table bloat: the outbox table will accumulate rows without a cleanup job (except under log-based CDC — see above). James Carr (2026-01-15) recommends a partial index on (sequence_id) WHERE published_at IS NULL for polling efficiency, plus a scheduled deletion job.
  • Schema evolution: outbox event schemas should be treated as API contracts. Conduktor (2026-07-30) recommends including a schema_version field and adding fields rather than removing or renaming without a migration plan.
  • Large payload anti-pattern: Streamkap (2026-02-25) recommends storing large binary objects in object storage and putting only a reference URL in the event, citing Kafka's default message size limit and replication lag.
  • Multiple relay instances: without FOR UPDATE SKIP LOCKED or a distributed lock, multiple polling relays can process the same row simultaneously, producing duplicate events (YouTube source, 2025-05-13; James Carr, 2026-01-15).
  • Operational complexity of self-managed Debezium: requires configuration of logical replication slots, connector restart policies, replication lag monitoring, offset storage, and connector upgrades (Streamkap, 2026-02-25). (as-of 2026-02-25)

Relationship to complementary patterns

Conduktor (2026-07-30) frames outbox and Saga Pattern as complementary: the outbox guarantees that a single service reliably publishes events when local state changes; the Saga coordinates a sequence of local transactions across multiple services, with each service using the outbox internally.

Confluent (2024-05-29) identifies three valid solutions to the dual-write problem: the transactional outbox pattern, event sourcing, and the listen-to-yourself pattern.

Conduktor (2026-07-30) contrasts Outbox+CDC vs 2PC:

  • Outbox+CDC: asynchronous, eventual consistency, high availability, loose coupling
  • 2PC: synchronous, strong consistency, reduced availability, tight coupling

Gunnar Morling (Current Bengaluru 2025) evaluates three alternatives, each with explicit cons:

  • Listen-to-Yourself (write Kafka first, read back): no synchronous read-your-own-writes
  • 2PC via KIP-939: reduced availability (both systems must be up simultaneously)
  • Stream processing: complexities around transactional consistency

Ecommerce platform context

commercetools caveat (docs.commercetools.com/api/projects/subscriptions, undated): commercetools subscriptions provide at-most-once delivery — "undelivered notifications are not retained and will not be delivered." This means practitioners building on commercetools cannot rely on the platform's subscription layer for outbox-style at-least-once delivery; a client-side outbox or Messages API polling fallback is required, though commercetools provides no documented guidance on this. (as-of release notes current to 2025-06-04)

Shopify historical case (Shopify Engineering, 2021-03-12):

Shopify replaced query-based CDC (Longboat) with Debezium log-based CDC on Kafka Connect after finding that updated_at-based extraction missed hard deletes, intermediate row states, and perpetually-updating rows. Over Black Friday/Cyber Monday 2020, Shopify processed ~65,000 records/second on average, with spikes to 100,000 records/second, running 150 Debezium connectors across 12 Kubernetes pods; p99 latency from MySQL insertion to Kafka availability was less than 10 seconds. (as-of 2021-03-12)

Contradictions

Historical lineage

James Carr (2026-01-15) traces the pattern's lineage to the early 2000s: Oracle AQ stored procedures writing XML to a database table with triggers firing to downstream queues; to RabbitMQ sidecars (local broker instances shovelling to a central cluster); to emergency local SQLite journals written at 2am during RabbitMQ outage incidents — concluding: "this pattern keeps showing up because the dual-write problem is fundamental."

Key terms

TermMeaning
Dual-write problemThe inconsistency that arises when a service must atomically update two separate systems (e.g. database + message broker) and one write fails
Outbox tableA dedicated database table that stores events as part of the same transaction that updates business entities
Relay / message relayA background process that reads from the outbox table and publishes events to the message broker
Polling publisherA relay that uses periodic scheduled queries to find unprocessed outbox rows
Transaction log tailingA relay that uses CDC tools (Debezium) to read the database WAL and detect new outbox rows
Outbox Event Router SMTDebezium's built-in Single Message Transform that routes outbox rows to per-aggregate Kafka topics
aggregate_idThe entity identifier used as the Kafka message key, ensuring per-entity event ordering within a partition
At-least-once deliveryEvents are guaranteed to be delivered at least once, but potentially more than once — consumers must be idempotent
pg_logical_emit_message()PostgreSQL-specific function for emitting logical replication messages without writing to an outbox table
LSN (Log Sequence Number)PostgreSQL's monotonically increasing WAL position; recommended by Morling as a duplicate-detection key for consumers

Benchmarks (as-of 2021-03-12 — Shopify, stale-risk)

  • Shopify CDC throughput: ~65,000 records/second average, 100,000 rec/sec peak (BFCM 2020)
  • Shopify p99 latency: <10 seconds (MySQL insertion to Kafka availability)
  • Shopify scale: 150 Debezium connectors, 12 Kubernetes pods, 400TB+ CDC data

No 2026 independent benchmark data available for outbox relay latency or throughput from named ecommerce systems.

Research agent · 2026-08-17