Why this matters
- Dual write (commit ledger then publish to Kafka) loses messages if the broker call fails after DB commit — ShardPay lost settlement events in early architecture before outbox adoption.
- Outbox is the standard pattern for reliable event publishing from transactional services.
- Inbox dedup on consumers completes the story for at-least-once end-to-end effectively-once processing.
- Interviewers ask how you atomically commit transfer row and event intent without 2PC across Postgres and Kafka.
The dual-write problem
Naive flow: BEGIN → insert transfer → COMMIT → kafka.send(). JVM dies after commit, before send — downstream never learns of transfer. Alternative: send first, then DB — broker has event DB lacks. Both orders fail one leg.
ShardPay needs atomicity between ledger state and event intent, not between two independent systems. The outbox table lives in the same Postgres database as transfers.
Walkthrough: crash after commit
Ledger transaction inserts transfer row and outbox row {eventId, payload, status=PENDING} in one commit. JVM crashes before relay runs. On restart relay polls WHERE status=PENDING, publishes to Kafka, marks SENT. No lost event — worst case delayed publish seconds until relay catches up.
Key points
- Outbox table — append-only pending events in the same DB as business data, committed in one transaction. ShardPay's
outbox_eventsholds Avro payload bytes, topic, partition key, andeventIdfor traceability. - Relay / CDC — polling publisher or Debezium CDC tail reads committed outbox rows and publishes to broker. ShardPay uses Debezium on
outbox_eventsfor sub-second relay latency; polling fallback in disaster mode. - Inbox table — consumer-side dedup store keyed by
eventId; prevents duplicate apply on at-least-once redelivery. ShardPay pairs every outboxeventIdwith inbox unique constraint on projector services. - Dual-write problem — attempting separate commits to DB and broker without coordination; guaranteed partial failure under crash. Outbox eliminates dual write by making broker publish asynchronous and retryable after DB truth exists.
- Ordering with outbox — relay publishes in
created_atorder per shard when using single relay partition; Debezium preserves binlog order. ShardPay includessequence_numper aggregate for consumers that need strict per-transfer order beyond broker defaults.
Inbox pattern on the consumer
Outbox fixes producer-side atomicity; inbox fixes consumer-side idempotency. Settlement projector: BEGIN → check inbox for eventId → apply credit → insert inbox row → COMMIT → then Kafka offset commit (or same transaction in transactional consumer setups).
Duplicate delivery hits inbox first — second apply skipped. Combined with outbox relay retries, the pipeline achieves effectively-once business outcomes.
Walkthrough: relay publishes twice
Debezium delivers same outbox row change twice during connector restart. Kafka has duplicate records with same eventId. Settlement consumer processes first, inbox records eventId. Second delivery no-ops at inbox — merchant balance unchanged second time. Ops sees duplicate Kafka records in debug topic but single ledger effect.
// ShardPay ledger — transactional outbox write (Java 17)
@Transactional
public TransferResult commitTransfer(TransferRequest req) {
Transfer transfer = transferRepo.insert(req);
OutboxEvent event = OutboxEvent.builder()
.eventId(UUID.randomUUID().toString())
.topic("shardpay.transfer.completed.v1")
.partitionKey(transfer.merchantAccountId())
.payload(serializer.toBytes(TransferCompleted.from(transfer)))
.status(OutboxStatus.PENDING)
.build();
outboxRepo.insert(event);
return TransferResult.from(transfer);
// Debezium publishes after commit; no kafka.send() in request thread
}
// Consumer inbox — paired with outbox eventId
@Transactional
public void applySettlement(SettlementPosted event, String eventId) {
if (inboxRepo.exists(eventId)) return;
merchantLedger.credit(event);
inboxRepo.insert(eventId, Instant.now());
}
Outbox operational concerns
Monitor relay lag (now - oldest PENDING), outbox table growth, and poison messages stuck in FAILED. ShardPay alerts when lag exceeds 30 seconds — merchants see delayed notifications but money is correct in DB.
Retention: mark SENT rows archived after 7 days. Idempotency TTL on inbox independent of outbox archive policy.
Kafka trackSee outbox pattern with Kafka
Quick recall
Everything you need if you only revisit this box.
- Outbox atomically commits business data and event intent in one DB transaction.
- Relay or CDC publishes asynchronously; crash-safe because DB is source of truth.
- Inbox dedup on consumers makes at-least-once delivery effectively-once for ledger projections.
Test yourself
Answer these before moving on — recall is what makes it stick.