Transactional outbox and inbox / idempotent consumer
Since Chapter 7, this series has named a specific, unaddressed gap: order-service commits its transaction, then publishes to Kafka in a separate step — and a crash between those two steps means a committed order with no event ever published for it. This chapter closes that gap. It’s a focused treatment inside this series’ Northwind story; for a full six-chapter deep dive into outbox design, relay strategies, and ordering, see this site’s dedicated Transactional Outbox series.
1. Problem the Pattern Solves
Northwind’s platform team finally reproduces the gap under load testing: they kill an order-service pod, deliberately, in the narrow window after its database transaction commits but before the Kafka send() call returns. The order exists, correctly, in order_schema — but inventory-service and payment-service never receive OrderPlaced, because the pod died before the message left the process. The order sits in PENDING forever, with no reservation attempted and no payment charged, invisible until a customer complains that an order they placed never progressed.
The reverse problem exists too, more subtly: if order-service published to Kafka before committing its transaction (an easy mistake to make, and the exact bug Chapter 1’s AFTER_COMMIT discipline exists to prevent), a transaction rollback after a successful publish would mean inventory-service reserves stock for an order that, from order-service’s point of view, never happened at all.
Forces in tension:
- Atomicity across two different systems vs. no distributed transaction. The database commit and the Kafka publish are two separate systems with no shared transaction manager spanning both (a 2PC-style XA transaction to Kafka is impractical and rarely used in production, as Chapter 1 already noted about distributed transactions generally) — yet the business needs them to behave as if they were atomic.
- Reliability vs. latency. Guaranteeing the event is eventually published, even after a crash, requires a relay mechanism that adds a small amount of latency (polling or log-tailing) between the commit and the message actually reaching Kafka — a trade-off against the “publish immediately” illusion Chapter 7’s code suggested.
- Consumer-side duplication vs. broker guarantees. Even with a reliable outbox guaranteeing at-least-once delivery, Kafka’s own semantics mean a message can still be delivered more than once to a consumer (a rebalance, a consumer crash after processing but before committing its offset) — the outbox alone doesn’t solve duplicate processing, only reliable sending.
2. Core Idea
Transactional outbox: instead of publishing to Kafka directly from application code, write the event to an outbox table in the same local database transaction as the business data change — since both writes are now part of one ACID transaction, they’re atomic with respect to each other by construction. A separate relay process then reads the outbox table and actually publishes to Kafka, retrying until it succeeds, and only then marking (or deleting) the outbox row.
Inbox / idempotent consumer: the mirror-image discipline on the receiving side — since Kafka’s at-least-once delivery means a consumer may see the same message more than once, the consumer records which message IDs it has already processed (an inbox table, or an idempotency check against its own business data) and skips duplicates, making “at least once delivery” behave like “effectively once processing” from the business’s point of view.
Commonly confused with:
- Change data capture (CDC), e.g. Debezium. CDC reads the database’s write-ahead log directly to detect changes, and can be used as the relay mechanism for an outbox (tailing the outbox table’s inserts) without a polling query — a legitimate, often more efficient implementation choice covered in the dedicated outbox series, not a different pattern.
- Event sourcing (Chapter 9). Event sourcing makes the event log itself the source of truth for a service’s state. An outbox table is a transient delivery mechanism — rows are relayed and then cleaned up; it’s an implementation detail of reliable publishing, not a durable historical record.
order-service’sorder_eventstable (Chapter 9) and itsoutboxtable (this chapter) can coexist: one is permanent history, the other is a short-lived delivery queue. - Retrying a failed publish call in application code. A simple
try/retryaroundkafkaTemplate.send()doesn’t survive a process crash between the database commit and the retry loop completing — the outbox pattern’s guarantee comes specifically from the event being durably persisted in the same transaction as the business change, not from retry logic alone.
3. When to Use It
Strong indicators:
- Any service publishing events after a local database write, where losing an event due to a crash between commit and publish is a real, consequential problem — which describes essentially every event publisher in this series since Chapter 7, Northwind’s
order-servicevery much included. - A demonstrated or plausible crash window between commit and publish — this chapter’s load-testing scenario is exactly the kind of evidence that turns “theoretically possible” into “worth fixing now.”
Concrete use cases:
- E-commerce, as here: order placement, payment capture, and inventory changes all need their resulting events reliably published, not just “usually” published.
- Financial services: a ledger entry and the notification event about it must never diverge — a lost “transaction posted” event after a successfully committed ledger entry is a direct reconciliation failure.
- Healthcare: a lab result recorded in a clinical system that fails to reliably notify a physician’s dashboard is a patient-safety-relevant gap, not just an engineering inconvenience.
- Any saga participant (Chapter 11): every step in a saga depends on its resulting event actually reaching the next participant — an unreliable publish silently breaks saga correctness in exactly the way Chapter 11 assumed was solved.
Prerequisites:
- A relational database already used transactionally for the business write (true for every service in this series since Chapter 1) — the outbox table lives in that same database specifically so the two writes share one transaction.
- A relay process (polling, as shown below, or CDC-based) that is itself monitored — an outbox with a dead relay is just a table that grows forever while publishing nothing.
- Idempotent consumers on the receiving end (this chapter’s inbox half) — an outbox guarantees at-least-once delivery, which is only useful if consumers handle redelivery correctly.
4. When Not to Use It
- The event’s loss is genuinely inconsequential. A low-stakes, best-effort notification (an internal Slack ping on order volume, say) may not justify the added table, relay process, and monitoring — weigh the actual cost of an occasionally-lost message against the pattern’s overhead.
- A service doesn’t yet publish any events after a local write. Introducing an outbox table for a service that has no current event-publishing need is pure anticipation — add it when the publishing need exists, not before.
- Very low event volume with a simple, monitored manual reconciliation process already in place. If a service publishes rarely enough that occasional manual reconciliation (finding and republishing missed events) is genuinely cheaper than building and operating a relay, the outbox’s operational cost may not be justified — an honest, volume-dependent judgment call, not a universal rule.
- Risk/overengineering signal: building a fully general, reusable “outbox framework” across every service before any specific service has demonstrated the crash-window problem in practice — as with every pattern in this series, let a real, demonstrated gap (as Section 1’s load test provided) drive adoption, not the pattern’s general appeal.
5. Implementation Example
The outbox table, in order-service’s own database, written in the same transaction as the order itself:
CREATE TABLE outbox ( id UUID PRIMARY KEY, aggregate_id UUID NOT NULL, event_type TEXT NOT NULL, payload JSONB NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), published_at TIMESTAMPTZ);CREATE INDEX idx_outbox_unpublished ON outbox (created_at) WHERE published_at IS NULL;package `in`.o612.eng.northwind.order.internal.outbox
import org.springframework.jdbc.core.JdbcTemplateimport org.springframework.stereotype.Componentimport java.util.UUID
@Componentinternal class OutboxWriter(private val jdbc: JdbcTemplate, private val serializer: EventSerializer) { fun write(aggregateId: UUID, eventType: String, payload: Any) { jdbc.update( "INSERT INTO outbox (id, aggregate_id, event_type, payload) VALUES (?, ?, ?, ?::jsonb)", UUID.randomUUID(), aggregateId, eventType, serializer.toJson(payload), ) }}package `in`.o612.eng.northwind.order.internal
import `in`.o612.eng.northwind.order.api.OrderPlacedimport `in`.o612.eng.northwind.order.internal.outbox.OutboxWriterimport org.springframework.stereotype.Serviceimport org.springframework.transaction.annotation.Transactional
@Serviceinternal class OrderService( private val orderRepository: OrderRepository, private val outboxWriter: OutboxWriter,) { @Transactional fun placeOrder(request: PlaceOrderCommand): Order { val order = orderRepository.save(Order.pending(request)) // Same transaction, same commit-or-rollback unit as the order // insert above — this is the entire guarantee the pattern provides. outboxWriter.write(order.id, "OrderPlaced", OrderPlaced(order.id, order.items)) return order }}Notice what disappeared compared to Chapter 7’s version: no ApplicationEventPublisher, no AFTER_COMMIT listener bridging to Kafka. The outbox row’s insertion is the commit-time guarantee now — there’s no longer a separate step after commit that could fail independently.
The relay, polling for unpublished rows and publishing them — a separate, simple scheduled component, deliberately not part of the request-handling path:
package `in`.o612.eng.northwind.order.internal.outbox
import org.springframework.jdbc.core.JdbcTemplateimport org.springframework.kafka.core.KafkaTemplateimport org.springframework.scheduling.annotation.Scheduledimport org.springframework.stereotype.Componentimport java.util.UUID
@Componentclass OutboxRelay( private val jdbc: JdbcTemplate, private val kafkaTemplate: KafkaTemplate<String, String>,) { @Scheduled(fixedDelay = 500) fun relayPendingEvents() { val pending = jdbc.queryForList( "SELECT id, aggregate_id, event_type, payload FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT 100" ) pending.forEach { row -> val topic = topicFor(row["event_type"] as String) kafkaTemplate.send(topic, row["aggregate_id"].toString(), row["payload"] as String).get() jdbc.update("UPDATE outbox SET published_at = now() WHERE id = ?", row["id"] as UUID) } }
private fun topicFor(eventType: String) = when (eventType) { "OrderPlaced" -> "orders.events"; else -> error("Unknown event type $eventType") }}.get() on the Kafka send makes the relay wait for broker acknowledgment before marking the row published — if the broker is unreachable, the row stays unpublished and is retried on the next poll, which is exactly the reliability guarantee the pattern promises. (The dedicated outbox series covers a log-tailing/CDC-based relay as a lower-latency alternative to this polling approach, and the ordering guarantees that need extra care once multiple relay instances run concurrently.)
The inbox / idempotent consumer, on inventory-service’s side, guarding against the redelivery Kafka’s at-least-once semantics guarantee will eventually happen:
CREATE TABLE inbox ( message_id UUID PRIMARY KEY, processed_at TIMESTAMPTZ NOT NULL DEFAULT now());package `in`.o612.eng.northwind.inventory.internal
import `in`.o612.eng.northwind.order.api.OrderPlacedimport org.springframework.dao.DuplicateKeyExceptionimport org.springframework.jdbc.core.JdbcTemplateimport org.springframework.kafka.annotation.KafkaListenerimport org.springframework.stereotype.Componentimport org.springframework.transaction.annotation.Transactionalimport java.util.UUID
@Componentinternal class IdempotentOrderPlacedListener( private val jdbc: JdbcTemplate, private val inventoryService: InventoryService,) { @KafkaListener(topics = ["orders.events"], groupId = "inventory-service") @Transactional fun onOrderPlaced(event: OrderPlaced, @org.springframework.kafka.support.KafkaHeaders headers: Map<String, Any>) { val messageId = UUID.fromString(headers["message-id"] as String) try { // Insert-or-fail: if this message_id is already in the inbox, // the insert throws and we skip processing — same transaction // as the business effect below, so both succeed or both roll back. jdbc.update("INSERT INTO inbox (message_id) VALUES (?)", messageId) } catch (e: DuplicateKeyException) { return // already processed — this is a redelivery, not an error } inventoryService.reserveStock(event.orderId, event.items.toReservationRequests()) }}The inbox insert and the actual business effect (reserveStock) happen in the same local transaction — this mirrors the outbox pattern’s core trick exactly, just on the receiving end: making two things atomic with respect to each other by putting them in one transaction, rather than trying to coordinate two separate systems.
6. Step-by-Step Flow
- Client action.
POST /orders, unchanged. - API request.
order-servicewrites the order and the outbox row in one transaction, then returns — the same latency profile as Chapter 7’s version, since the outbox insert is a fast local write, not a network call. - Service behavior. The relay, running independently on its own schedule, picks up the unpublished row.
- Database interaction. The relay’s
UPDATE ... SET published_atonly happens after Kafka acknowledges the send — a crash of the relay itself before that update simply means the row is picked up again on the next poll, safely (the relay’s own publish call being effectively idempotent from Kafka’s perspective, modulo Section 7’s producer-idempotence note). - Inter-service communication.
inventory-servicereceives the event, possibly more than once over the flow’s lifetime. - Error or failure handling. A redelivered message is caught by the inbox check and skipped — no double stock decrement, closing exactly the idempotency requirement this series has flagged as a design obligation since Chapter 1.
- Observability signals. Track outbox table depth (unpublished row count) and age of the oldest unpublished row — a growing, aging backlog means the relay is stuck or the broker is unreachable, a direct, early signal before anyone notices a missing downstream effect.
- Final response/outcome. Every committed order is guaranteed, eventually, to have its event published — and every consumer processes that event’s business effect exactly once, regardless of how many times Kafka happens to redeliver it.
7. Production Concerns
- Timeouts, retries, idempotency. Configure the Kafka producer itself with
enable.idempotence=trueto avoid the relay’s own retries producing duplicate messages at the broker level, on top of the inbox pattern handling duplicates at the consumer level — belt and suspenders, addressing duplication at both the producer and consumer ends. - Data consistency and transaction boundaries. The outbox and inbox patterns are precisely about making a local write and a messaging effect atomic with respect to each other, without a distributed transaction — the core technique (piggyback the messaging concern onto a transaction you already have) generalizes to any “local write plus reliable external effect” problem, not just Kafka publishing.
- Outbox table growth and cleanup. Published rows should be deleted or archived on a schedule (a separate, low-priority scheduled job) — an ever-growing outbox table, even one whose rows are all
published_at IS NOT NULL, eventually becomes a vacuum and query-performance problem in Postgres. - Relay availability and scaling. A single relay instance is a throughput bottleneck and a soft single point of failure (a stuck relay stalls all publishing); running more than one relay instance safely requires a locking or partitioning scheme (
SELECT ... FOR UPDATE SKIP LOCKEDis a common, effective technique) to avoid two relays racing to publish the same row — covered in depth in the dedicated outbox series. - API versioning. Unaffected directly, but note that the outbox’s
payloadcolumn now carries the same event schema whose evolution discipline Chapters 7–8 already established — no new versioning surface, just a new place the existing contract is serialized. - Authentication and service-to-service trust. No new surface — the relay uses the same Kafka producer credentials the service would have used publishing directly.
- Logging, metrics, tracing, correlation IDs. Propagate the correlation ID into the outbox row’s payload (or a dedicated column) so it survives the relay hop and appears on the eventually-published Kafka message exactly as it would have if published synchronously.
- Kubernetes deployment. The relay can run as a scheduled task inside the same deployable as the service (as shown, via
@Scheduled) for simplicity at Northwind’s current scale, or as its own deployable for independent scaling once outbox volume justifies it — evaluate this the same way Chapter 2 evaluated any extraction, on measured need. - Testing strategy. Test the outbox write and the business write’s atomicity directly (force a rollback after the outbox insert and confirm neither persists); test the relay against an embedded/Testcontainers Kafka; test the inbox’s duplicate-skip behavior by delivering the same message twice in a test and asserting the business effect happened exactly once.
- Migration strategy. Introduce the outbox for the publisher side and the inbox for the consumer side together, per event flow, starting with the flow that has the most consequential crash-window risk (Northwind started with
OrderPlaced, given its role at the head of the saga from Chapter 11) — not as a platform-wide rewrite in one release.
8. Common Mistakes
- Publishing directly from application code “as a fallback” alongside the outbox. Keeping a direct
kafkaTemplate.send()call as a belt-and-suspenders measure alongside the outbox reintroduces exactly the crash-window gap the outbox exists to close, plus a risk of double-publishing. Fix: the outbox table is the only path to publishing once adopted — remove the direct call entirely. - No cleanup of published outbox rows. Leaving every historical row in the table forever, even marked
published_at, degrades query and vacuum performance over time. Fix: a scheduled cleanup job deleting or archiving rows older than a retention window. - Treating the inbox as optional “since the outbox already guarantees delivery.” The outbox guarantees the message is sent, at least once — it says nothing about how many times a consumer processes it. Skipping the inbox reintroduces duplicate-processing bugs (double stock decrements) that the outbox alone cannot prevent. Fix: implement both halves — outbox on the publisher, inbox on every consumer that isn’t already naturally idempotent.
- Running multiple relay instances without a locking scheme. Naively running two relay replicas for availability, without
SKIP LOCKEDor equivalent partitioning, causes both to pick up and publish the same unpublished row, producing duplicate messages (survivable with a correct inbox, but wasteful and worth avoiding at the source). Fix: use row-level locking (FOR UPDATE SKIP LOCKED) or partition the outbox across relay instances explicitly. - Ignoring outbox backlog as an operational signal. Not monitoring unpublished row count or age means a stuck relay (a crashed pod, a misconfigured Kafka connection) goes unnoticed until a business-side symptom (an order stuck in
PENDING) surfaces the problem indirectly and late. Fix: alert on backlog depth and age directly, as a leading indicator, not a lagging one. - Assuming Kafka’s
enable.idempotence=truemakes the consumer-side inbox unnecessary. Producer idempotence prevents the producer’s own retries from creating duplicates at the broker — it does nothing about a consumer crashing after processing a message but before committing its offset, which still causes redelivery. Fix: understand these as two separate guarantees addressing two separate failure points, both needed.
9. Decision Guide
| Problem signal | Use this pattern? | Why | Alternative |
|---|---|---|---|
| Service publishes an event after a local database write, and losing that event matters | Yes | Closes the commit-then-publish crash window with a mechanism that survives a crash | — |
| Event’s occasional loss is genuinely low-stakes | Maybe not | Relay infrastructure and monitoring cost may exceed the value of the guarantee | Direct publish with monitoring/alerting on anomalies |
| A consumer of any at-least-once broker | Yes (inbox half) | Redelivery will happen eventually; processing must be idempotent regardless of publisher reliability | — |
| Service doesn’t yet publish any events | No, not yet | Nothing to protect until the publishing need exists | Add when the first event-publishing requirement arrives |
| High outbox/relay volume, single relay instance becoming a bottleneck | Yes, plus scale the relay | The core pattern still applies; add locking/partitioning for multiple relay instances | — |
10. Hands-On Exercise
Extend it: add outbox-backed publishing to payment-service for its PaymentCaptured/PaymentFailed events (currently, per Chapter 11’s saga implementation, published directly) — apply the same commit-time guarantee to the saga’s second step.
Simulate a failure: in a test or local environment, insert a row directly into order-service’s outbox table without going through placeOrder at all, and confirm the relay picks it up and publishes it on its next poll — proving the relay’s correctness is independent of how the row got there. Then simulate a relay crash between the Kafka send and the UPDATE ... published_at and confirm the row is republished (and that the inbox on the consumer side correctly deduplicates the resulting redelivery).
Decision question, with justification required: Northwind’s outbox relay currently polls every 500ms. The team is considering switching to a CDC-based relay (Debezium tailing the outbox table’s write-ahead log) for lower latency and reduced database polling load. Given the current 500ms latency is well within customer-visible tolerance for order confirmation, what specific, measured signal would justify the added operational complexity of running Debezium? Name the trade-offs from Section 1.
11. Key Takeaways
- The transactional outbox closes the gap between “committed a local transaction” and “successfully published the resulting event” by writing the event to a table in the same transaction as the business change, then relaying it separately, with retry, until Kafka acknowledges it.
- The inbox / idempotent consumer pattern is the necessary mirror image — an outbox guarantees at-least-once sending, which only becomes safe once every consumer handles at-least-once receiving correctly.
- Both patterns use the same core trick: making two things atomic with respect to each other by putting them in one local database transaction, instead of trying to coordinate two separate systems directly.
- Remove any direct publish call once the outbox is adopted — keeping both creates duplication risk and defeats the guarantee.
- Monitor outbox backlog depth and age as a leading operational signal — a stuck relay is otherwise invisible until a business-side symptom surfaces it much later.
- Multiple relay instances need explicit locking or partitioning to avoid duplicate publishing; a correct inbox tolerates this, but it’s worth avoiding at the source.
- This pattern directly protects the sagas built in Chapter 11 — a saga’s correctness assumed every step’s event reliably reaches the next participant, an assumption this chapter is what actually makes true.