Technologies that implement this pattern: PostgreSQL · Apache Kafka · DynamoDB · Apache Flink · Elasticsearch · Redis
Go deeper:
- When the event log itself should be the source of truth rather than a side effect of a table write, see CQRS & Event Sourcing for projections, rebuilds and event versioning.
Why This Matters#
Every service that writes to a database and then tells someone about it has a dual-write bug, whether or not it has shipped yet. The code looks harmless: db.commit() then kafka.send(). It works in every test and in 99.99% of production requests. The other 0.01% are the process that crashed between the two lines, the broker that timed out after the commit, the deploy that killed the pod mid-request. At 2,000 writes/s, 0.01% is 17,000 silently divergent records a day: orders the warehouse never heard about, prices the search index never saw, refunds the ledger never booked.
Most candidates treat this as a messaging question: "publish an event to Kafka after the write." Staff engineers treat it as an atomicity-boundary question. There is exactly one place in the system that can commit atomically, and that is the database transaction. The design question is not which broker but how do I get the fact that this transaction committed out of the database without a second commit that can fail independently? There are only two honest answers: write the event into the same transaction (the outbox) or read the database's own commit log (change data capture). Everything else is a dual write with better logging.
The second reframe: the outbox is not a queue, it is a contract. CDC on raw tables couples every downstream consumer to your column names; an outbox table lets you publish a deliberate, versioned event shape while still getting the atomicity of a local commit. Choosing between them is choosing who pays when the schema changes.
If you can walk an interviewer from "where is the atomic boundary" to "how does the commit leave the database" to "what does the consumer do with duplicates and reordering" to "who gets paged when the relay stops," you are answering at Staff level.
The 60-Second Version#
- Dual writes fail at a measurable rate, not "rarely." Any write path with two independent commits diverges on every crash, timeout or deploy between them. At 0.01% of 2,000 writes/s that is ~17K inconsistent records per day, and none of them throw an error.
- The outbox costs one extra row per transaction. An
INSERTinto an outbox table inside the business transaction adds ~0.1–0.5ms and one WAL record. That is the entire price of atomicity on the write path. - The relay delivers at least once — always. Whether a poller or a CDC connector moves rows to the broker, a crash after publish and before checkpoint republishes. Consumers must dedupe on
event_id; budget for ~0.1–1% duplicates during failovers. - Lag is the failure mode, and it is measured in seconds. A healthy CDC relay runs 100ms–2s behind commit. Alert on
outbox.publish_lag_msp99 > 5s and on the age of the oldest unpublished row, not on row count. - An abandoned replication slot will fill the primary's disk. PostgreSQL retains WAL for an unconsumed logical slot indefinitely by default. At 50 MB/s of WAL, a 2 TB volume has ~11 hours before the database stops accepting writes. Cap it with
max_slot_wal_keep_sizeand page on slot lag. - Order is per key, never global. Key every event by aggregate ID so one order's events land on one partition in commit order. Global ordering across 48 partitions does not exist, and a design that needs it has a different problem.
The Problem#
A checkout service writes an order to PostgreSQL and needs four other systems to react: fulfilment must ship it, search must index it, the ledger must book revenue, and email must confirm it. The naive implementation commits the order and then publishes OrderPlaced. If the publish fails, the order exists and nothing downstream knows. If you reverse the order and publish first, a failed commit leaves the warehouse shipping an order that does not exist. Wrapping both in a distributed transaction (XA, two-phase commit) needs a broker that participates in one — Kafka does not — and turns every checkout into a coordinated commit across two systems with independent failure modes. Retrying the publish in-process loses the retry when the process dies. A nightly reconciliation job finds the gaps 18 hours late. The job is to make "this row committed" and "this event will be delivered" the same fact, then deliver it reliably and in per-key order to consumers that tolerate seeing it twice.
Case Studies That Use This Pattern#
- Payments — Ledger entries and payment-state events must never diverge; the outbox makes "charge recorded" and "charge announced" one commit
- Idempotency — The consumer half of the pattern: at-least-once delivery is only safe when every handler dedupes on an event ID
- Notifications — "Send the confirmation email" is a side effect of an order commit, not a second write the API may forget
- Search Engine — Keeping an index in step with the system of record through CDC instead of application-level dual writes
- News Feed — Post-created events drive fan-out; a lost event is a post nobody sees
- Hotel Booking — Reservation commits drive inventory, billing and partner sync; divergence is double-booking
- Message Broker — The broker the relay publishes into; retention and partitioning set the replay and ordering guarantees
- Distributed Cache — CDC-driven invalidation closes the race that application-side delete-after-write leaves open
Which Problem Are We Solving?#
"Publish events reliably" hides four goals with different designs. Name them and commit before drawing a box.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Domain events to other teams (OrderPlaced, PaymentCaptured) | Consumers are other services with their own release cycles; event shape is a public contract | Outbox table in the same transaction + CDC or polling relay; versioned schema in a registry | Relay lag; schema change breaks consumers; duplicate delivery | Every committed fact delivered at least once, in per-aggregate order, within an agreed lag (e.g. p99 < 5s) |
| Derived read models you own (search index, cache, materialized views) | Same team owns source and sink; full-row state is what's needed | Table-level CDC (no outbox), idempotent upserts keyed by primary key and source LSN | Snapshot/backfill races with live changes; index drifts silently | Sink converges to source; drift measured and bounded |
| Data platform replication (warehouse, lake, analytics) | Hundreds of tables, schema churn, high volume; minutes of lag acceptable | Log-based CDC into compacted topics, schema registry, consumers own transformation | Schema change breaks the pipeline; slot or binlog retention runs out | Complete history including deletes; hourly reconciliation counts |
| Cross-service workflow step (reserve inventory, then charge) | A business process spans services and needs a reaction to each step | Outbox per service plus a saga or workflow engine consuming the events | Missing compensation; a step's event lost or processed twice | Every step happens exactly once in effect, via idempotent handlers |
🎯 Staff Move: "I'll treat this as intent one: other teams consume OrderPlaced, so the event shape is a contract I don't want tied to my table layout. That means an outbox table written in the order transaction, relayed by CDC into Kafka keyed by order ID. Search is intent two — my team owns it — so it can read table-level CDC directly and skip the outbox."
The Core Tradeoff#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Dual write (commit, then publish) | Simplest code; lowest latency | Diverges on every crash or timeout between the two calls; no error is ever raised | Downstream teams discovering missing orders; whoever runs reconciliation |
| Distributed transaction (2PC/XA) | True atomicity across systems | Most brokers don't support it; coordinator failure blocks commits; latency of the slowest participant | Every write's p99; the on-call when the coordinator hangs |
| Outbox + polling publisher | Atomic with the business write; no extra infrastructure; easy to reason about | Poll interval adds latency (100ms–1s); polling query load on the primary; ordering needs care with concurrent pollers | The primary database's CPU and the team owning the poller |
| Outbox + log-based CDC relay | Atomic; ~100ms–2s lag; no polling load; outbox rows can be deleted in the same transaction | New infrastructure (Debezium, Kafka Connect); replication slot risk to the primary's disk; connector failover republishes | Platform team running connectors; DBA watching slot lag |
| Table-level CDC (no outbox) | Zero application changes; captures every change including deletes and changes made by scripts | Consumers coupled to physical schema; intermediate states and internal columns leak; events lack business intent | Every consumer team at each schema migration |
| Event sourcing (the log is the store) | No dual write by construction; full history | Every read needs a projection; schema evolution of old events is forever; steep team learning curve | The owning team, for the lifetime of the system |
| Listen-to-yourself (publish first, consume your own event to write the DB) | Single write, to the log | Read-your-writes breaks — the API cannot return the committed state; DB write can lag or fail after the client got a 200 | Users who refresh and don't see their change |
Staff Default Position#
Commit the event in the same local transaction as the state change, relay it from the commit log, and make every consumer idempotent on the event ID — never publish to a broker from request-handling code after a commit.
The default stack: an outbox table (event_id, aggregate_type, aggregate_id, event_type, payload, created_at) written in the business transaction; a log-based CDC connector (Debezium on the PostgreSQL WAL or MySQL binlog) routing outbox rows to a topic per aggregate type, keyed by aggregate_id so per-entity order holds; consumers dedupe on event_id with a processed-events table or an idempotent upsert written in the same transaction as their side effect. Below ~200 events/s, or where adding a connector is a quarter of platform work, a polling publisher with FOR UPDATE SKIP LOCKED is the right first version. Lag is alerted in seconds; replication slot retention is capped; the outbox schema is registered and versioned like an API. Derived read models your own team owns can use table-level CDC directly.
When to Deviate#
- The consumer only needs eventual repair, not prompt delivery — a nightly analytics export or a monthly invoice run. A reconciliation scan over
updated_at(with soft deletes) is simpler than a relay. Say what it misses: hard deletes and intermediate states. - The write and the event go to the same store — DynamoDB Streams or a single Kafka-backed event-sourced service. The store's own change log already is the outbox; adding a table is ceremony. Check the stream's retention (DynamoDB Streams keeps records for 24 hours) before you depend on it for recovery.
- Throughput beyond what one database's log can relay — a single CDC connector per database is single-threaded per slot; around 20–50K events/s per source you shard the source, not the relay. If the volume is that high and the database isn't the system of record anyway, write to the log first and make the database a projection.
- The event must be visible before the database write — rare, usually telemetry or audit where loss is worse than a phantom. Publish to the log first with idempotent downstream writes, and accept that the API cannot promise read-your-writes.
One Question, Three Levels#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Write to the DB, then publish to Kafka, with retries" | "Where's the atomic boundary? The event goes in the same transaction, and the relay reads the log." | "How many teams hand-roll this? One paved outbox library and one connector platform, or 40 subtly different relays?" |
| Delivery semantics | "Kafka gives us exactly-once" | "Relay is at-least-once; consumers dedupe on event ID in the same transaction as their effect" | Sets an org rule: every event carries an ID and a source position; idempotent consumption is a review requirement |
| Ordering | Assumes events arrive in commit order | Keys by aggregate ID; per-key order only; consumers handle version gaps | Decides which domains get per-key ordering contracts and which are explicitly unordered |
| Schema | Publishes the ORM entity as JSON | Outbox payload is a versioned contract separate from the table; registry enforces compatibility | Owns the event catalogue: ownership, deprecation windows, breaking-change process across teams |
| Failure | "If Kafka is down we retry" | "If the connector stops, the slot retains WAL; at 50 MB/s we have ~11 hours before the primary's disk fills. Cap it and page." | Treats CDC as tier-0 infrastructure with an error budget, game days and a documented replay path |
| Ownership | The service team owns everything | Producer owns the outbox schema; platform owns connectors; consumers own dedupe and DLQs | Draws the contract between platform and product teams; funds the connector platform as shared infra |
Why "First move" separates levels
The L5 answer — commit then publish with retries — is what most production code actually does, and it gets downleveled because the retry lives in the same process that just crashed. The Staff candidate identifies the database transaction as the only atomic boundary and moves the event inside it. The Principal candidate notices that every team will face the same question and that the org's real risk is twenty hand-written relays with twenty different bugs, not one service's dual write.
Why "Delivery semantics" separates levels
Kafka's transactional producer gives exactly-once within Kafka — read from a topic, write to a topic, commit offsets atomically. It cannot make a PostgreSQL commit and a Kafka write one atomic act, and it does nothing for a consumer whose side effect is an HTTP call or a row in another database. The Staff answer accepts at-least-once at the relay and pushes exactly-once-in-effect to the consumer, where the dedupe check and the side effect share one transaction. That is the only place it can actually be enforced.
Why "Failure" separates levels
Seniors think about the broker being down. Staff engineers think about the relay being down, because that failure is silent to the application and dangerous to the database: a logical replication slot pins WAL until it is consumed. The service keeps committing happily while the primary's disk fills toward a full outage of the system of record — caused by a component that was supposed to be downstream of it. Knowing that a CDC connector can take down its source is the signal interviewers listen for.
Where the Design Splits#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Outbox vs Table-Level CDC | A deliberate event contract (extra table, extra code) vs zero-touch capture that couples consumers to the schema |
| 2 | Polling Relay vs Log-Tailing Relay | No new infrastructure and simple ops vs sub-second lag and no query load, with a slot that can hurt the primary |
| 3 | Thin Events vs Fat Events | Small notifications that make consumers call back vs full state that makes the payload a schema contract |
| 4 | Dedupe at the Consumer vs Exactly-Once Machinery | Idempotent handlers everywhere vs transactional pipelines that only cover part of the path |
| 5 | Outbox Retention: Delete Now vs Keep for Replay | An empty table and a small WAL footprint vs a queryable history you can replay without the broker |
Fault Line 1: Outbox vs Table-Level CDC#
Table-level CDC reads every insert, update and delete on orders and publishes the row image. It needs no application change and catches writes made by migrations, admin scripts and other services sharing the database. But it publishes storage, not intent: a status column flip from 3 to 4 instead of OrderShipped, internal columns consumers start depending on, and three update events for one business action. An outbox publishes exactly the event the domain meant, with a shape you version independently of the table. Who pays: table CDC makes every consumer pay at every schema migration; the outbox makes the producing team pay one insert and an event-design review per new event. Staff default: outbox for anything another team consumes; table CDC for read models your own team owns (search, cache, analytics replicas). Deviate when: you need a complete audit of every row change regardless of which code path made it — table CDC is the only thing that catches the 2am manual UPDATE.
Fault Line 2: Polling Relay vs Log-Tailing Relay#
A polling publisher selects unpublished outbox rows every 100–500ms, publishes them, and marks or deletes them. A log-tailing relay (Debezium, a custom WAL reader, DynamoDB Streams) reads the commit log and never queries the table. Who pays: polling pays in latency (half the interval on average), in index churn and query load on the primary (~2–10 queries/s per poller), and in ordering subtlety when multiple pollers run; log tailing pays in a new piece of tier-0 infrastructure, connector operations and the replication slot that retains WAL if the relay stops. Staff default: log tailing once you run more than a couple of services on the pattern or need p99 lag under 1s; polling for a first version or a single low-volume service. Deviate when: the database is a managed offering that doesn't expose logical replication, or the org has no team to own connectors — a well-built poller beats an unowned Debezium cluster.
Fault Line 3: Thin Events vs Fat Events#
A thin event says "order 991 changed, version 7." Consumers call the order service for the state. A fat event carries the full order. Who pays: thin events make the producer pay in read load (every consumer calls back, often in a burst right after the event) and couple consumers to its availability; fat events make the producer pay in contract surface — every field in the payload is a field you cannot rename — and make the payload size a concern above ~100 KB per message. Staff default: fat-enough events — the fields consumers need to act without calling back, plus the aggregate version — published under a registered schema with backward-compatible evolution. Deviate when: the data is sensitive (PII, card data) and should not fan out to every topic subscriber; send a thin event and make consumers fetch through an authorized API.
Fault Line 4: Dedupe at the Consumer vs Exactly-Once Machinery#
Kafka transactions and Flink's two-phase commit sinks give exactly-once for data that stays inside them. They do not cover the CDC relay's republish after a connector restart, a consumer calling an email provider, or a write into another team's database. Who pays: relying on exactly-once machinery pays in false confidence and a gap at the edges; consumer-side dedupe pays a processed-events table (~50 bytes per event, retained for at least the replay window) and one extra indexed lookup per message. Staff default: every consumer dedupes on event_id (or applies an upsert guarded by source version) inside the same transaction as its effect; use transactional pipelines as an optimization inside stream processing, not as the correctness argument. Deviate when: the side effect is naturally idempotent — setting a value to the latest version, a PUT of full state — then the version check is the dedupe.
Fault Line 5: Outbox Retention — Delete Now vs Keep for Replay#
With log-based relays you can insert and delete the outbox row in the same transaction: the WAL carries the insert, the table stays empty, and there is nothing to vacuum. With a retained outbox (delete after N days by a janitor), you can replay a range of events from SQL without touching the broker. Who pays: delete-now pays in having no replay source except the broker's retention; retain-for-replay pays in table bloat — 2,000 events/s × 1 KB × 7 days is ~1.2 TB — and vacuum or partition-drop maintenance. Staff default: with CDC, keep a short retention (24–72h, daily partitions dropped whole, never row-by-row deletes) as a replay buffer that outlives a connector incident; the broker's 7-day retention is the longer replay source. Deviate when: a polling relay owns the table — then retention is mandatory because "published" is a column the poller updates.
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "Save to the database, then publish to Kafka" | "Hasn't been paged for the event that was never sent" | "That's a dual write. The event goes into an outbox row in the same transaction, and a relay reads it from the log." |
| "We'll use Kafka exactly-once so there are no duplicates" | "Thinks Kafka transactions span PostgreSQL" | "The relay is at-least-once. Every consumer dedupes on event ID inside the transaction that applies the effect." |
| "Just retry the publish if it fails" | "The retry lives in a process that can die" | "A retry in memory dies with the pod. Durability of the intent has to come from the database commit itself." |
| "Use CDC on the orders table and let consumers subscribe" | "Every column is now a public API" | "Table CDC for my own read models. Other teams get a versioned outbox event, so I can still refactor the table." |
| "Events arrive in order" | "Hasn't seen a partition rebalance" | "Order is per aggregate, by keying on order ID. Consumers check the aggregate version and park anything that arrives ahead of a gap." |
| "Monitor the outbox table size" | "Count, not time; misses the slot" | "Alert on publish lag p99 and oldest unpublished row age, and separately on replication slot retained bytes — that one can take down the primary." |
Quick Reference#
Staff Sentence Templates#
"The only atomic boundary here is the [database] transaction. So [event] is written to the outbox in that transaction, and [relay] moves it to [topic] keyed by [aggregate ID]. Delivery is at-least-once; [consumer team] dedupes on event ID."
"Relay lag is normally [X] ms. I'll page on p99 publish lag over [Y] s and on replication slot retention over [Z] GB, because at [W] MB/s of WAL that's [hours] before the primary's disk is at risk."
"Table-level CDC is fine for [read model] because my team owns both ends. [Other team] gets a versioned [event] from the outbox, so a column rename on my side is not an outage on theirs."
"Ordering is per [aggregate]. If a consumer sees version [n+2] before [n+1], it [parks the event / fetches current state], and it never assumes global order across partitions."
Implementation Deep Dive#
1. The Outbox Table and the Business Transaction — PostgreSQL#
The write path changes by one statement. The outbox insert shares the order's transaction, so either both exist or neither does.
CREATE TABLE outbox (
event_id uuid PRIMARY KEY, -- consumers dedupe on this
aggregate_type text NOT NULL, -- routes to topic: 'order' -> orders.events
aggregate_id text NOT NULL, -- message key: per-order ordering
aggregate_ver bigint NOT NULL, -- lets consumers detect gaps and stale events
event_type text NOT NULL, -- 'OrderPlaced', 'OrderCancelled'
schema_ver int NOT NULL, -- payload contract version
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
) PARTITION BY RANGE (created_at); -- daily partitions, dropped after 72h
BEGIN;
UPDATE orders SET status = 'PLACED', version = version + 1
WHERE id = $order_id AND version = $expected_ver; -- optimistic check on the aggregate
INSERT INTO outbox (event_id, aggregate_type, aggregate_id, aggregate_ver,
event_type, schema_ver, payload)
VALUES (gen_random_uuid(), 'order', $order_id, $expected_ver + 1,
'OrderPlaced', 3, $payload_json);
COMMIT; -- one fsync, both facts durable
Why the extra columns matter: aggregate_ver is what lets a consumer tell "event for version 8 arrived before version 7" from "duplicate of version 7." schema_ver lets a consumer reject or upcast a payload it does not understand instead of silently mis-parsing it. Daily partitions mean retention is a DROP TABLE of yesterday's partition — no row-by-row deletes, no vacuum debt on a table that receives 170M inserts a day at 2,000/s.
🎯 Staff Insight: Build the payload inside the transaction from the same in-memory state you just wrote, not by re-reading the row after commit. A post-commit read can see a later concurrent update and publish a payload that never corresponds to
aggregate_ver.
2. Log-Based Relay — Debezium Outbox Event Router#
The connector reads the WAL through a logical replication slot and routes outbox inserts to a topic per aggregate type, using the aggregate ID as the Kafka message key.
# Kafka Connect connector config (abridged)
connector.class = io.debezium.connector.postgresql.PostgresConnector
plugin.name = pgoutput
slot.name = outbox_orders
table.include.list = public.outbox
transforms = outbox
transforms.outbox.type = io.debezium.transforms.outbox.EventRouter
transforms.outbox.route.by.field = aggregate_type # 'order' -> topic outbox.event.order
transforms.outbox.table.field.event.key = aggregate_id # Kafka key = per-order partition
heartbeat.interval.ms = 10000 # advances the slot on quiet databases
# Producer side of Connect
producer.override.acks = all
producer.override.enable.idempotence = true
# PostgreSQL primary — protect the database from its own relay
max_slot_wal_keep_size = 200GB # PG13+: slot is invalidated instead of filling the disk
# alert: pg_replication_slots retained bytes > 20GB (warn), > 100GB (page)
What the router gives you: only INSERTs on the outbox are published; deletes are filtered, which is why inserting and deleting the row in one transaction still works — the insert is in the WAL even if the table is empty a microsecond later. The event ID travels as a message header for consumer dedupe.
The trade you make with max_slot_wal_keep_size: without it, a dead connector pins WAL until the disk fills and the primary stops. With it, a dead connector past the cap gets its slot invalidated, and you must re-snapshot or replay from the retained outbox partitions. That is the correct trade — losing the relay's position is recoverable; losing the system of record is an outage — but someone has to sign off on it and own the recovery runbook.
3. Polling Publisher — The First Version That Is Still Correct#
When there is no CDC platform, a poller is fine if it does three things right: it claims rows without blocking other pollers, it publishes before it marks, and it preserves per-aggregate order.
POLL_INTERVAL = 200ms
BATCH = 500
loop every POLL_INTERVAL:
BEGIN
rows = SELECT * FROM outbox
WHERE published_at IS NULL
ORDER BY created_at, event_id
LIMIT BATCH
FOR UPDATE SKIP LOCKED -- concurrent pollers take disjoint rows
for r in rows:
producer.send(topic_for(r.aggregate_type), key=r.aggregate_id,
value=r.payload, headers={event_id: r.event_id})
producer.flush() -- wait for acks=all on every message
UPDATE outbox SET published_at = now() WHERE event_id = ANY(ids(rows))
COMMIT -- crash before here => republish (at-least-once)
janitor hourly:
DELETE FROM outbox WHERE published_at < now() - interval '72 hours' -- or drop partitions
index: CREATE INDEX ON outbox (created_at) WHERE published_at IS NULL -- stays tiny
The ordering trap: with SKIP LOCKED and two pollers, poller A can lock order 991's version 7 while poller B grabs version 8 and publishes it first. If per-aggregate order matters, partition the work — each poller owns a hash range of aggregate_id (WHERE hashtext(aggregate_id) % 4 = $my_slot) — or run a single active poller with a lease and accept its throughput ceiling (~2–5K events/s with batched sends).
🎯 Staff Move: "I'll start with a polling relay because we don't have a Connect cluster and volume is 150 events per second. The poller owns hash ranges of order ID so per-order ordering holds. If we add three more services on this pattern, that's the trigger to fund a CDC platform instead of four pollers."
4. The Consumer Half — Inbox Dedupe in the Effect's Transaction#
At-least-once delivery is only safe if the dedupe check and the side effect commit together. Checking in Redis and then writing to PostgreSQL is another dual write.
-- consumer's own database
CREATE TABLE processed_events (
consumer text NOT NULL,
event_id uuid NOT NULL,
processed_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (consumer, event_id)
);
-- per message
BEGIN;
INSERT INTO processed_events (consumer, event_id) VALUES ('fulfilment', $event_id)
ON CONFLICT DO NOTHING;
-- if 0 rows inserted: duplicate -> COMMIT and ack, do nothing else
INSERT INTO shipments (order_id, ...) VALUES (...);
-- or, for state replication: guard on version instead of a dedupe table
-- UPDATE order_view SET ... , ver = $agg_ver WHERE order_id = $id AND ver < $agg_ver;
COMMIT;
-- then commit the Kafka offset; a crash between COMMIT and offset commit = harmless redelivery
Retention rule: keep processed_events at least as long as the longest replay you will ever do (broker retention, typically 7 days). At 2,000 events/s that is ~1.2B rows × ~60 bytes ≈ 70 GB — partition it by day like the outbox. For side effects outside your database (email, payment provider), pass event_id as the provider's idempotency key so the dedupe happens on their side too (Idempotency).
Technique Comparison
| Technique | Atomic With Write | Typical Lag | Load on Primary | Ordering | Ops Burden | Best For |
|---|---|---|---|---|---|---|
| Dual write | No | ~0 | None | Per call | None | Nothing that matters |
| Outbox + polling | Yes | 100ms–1s | Polling queries, index churn | Per key with hash-range pollers | Low | First version, < ~2K events/s |
| Outbox + log CDC | Yes | 100ms–2s | WAL decoding, slot retention | Per key via message key | Connector platform | Domain events at any scale |
| Table-level CDC | Yes (it is the commit) | 100ms–2s | Same as above | Per primary key | Connector platform + schema coupling | Own read models, warehouse feeds |
| Store change stream (DynamoDB Streams) | Yes | ~sub-second | None | Per item key | Low; 24h retention | Single-store services |
| Reconciliation scan | Eventually | Minutes–hours | Periodic scans | None | Low | Backstop and audits |
Architecture Diagram#
How to narrate it: the only arrow the checkout request waits on is the one into PostgreSQL. Everything to the right is asynchronous and at-least-once. The slot is drawn on purpose — it's the component that couples the relay's health back to the primary, so it has a cap and a page. Each consumer owns its dedupe and its DLQ; the platform team owns the connector and the lag alert; the order team owns the schema in the registry.
Failure Scenarios#
1. The Silent Dual Write — 0.02% of Orders Never Shipped#
An order service commits to MySQL and then publishes OrderPlaced with a 3-retry loop. It handles 1,500 orders/s at peak.
Week 1 Deploys roll pods with a 30s termination grace period. In-flight publishes
after commit are killed on SIGKILL for requests still running.
Week 1-6 ~0.02% of orders committed without an event: ~25K orders over six weeks.
Week 6 Customer support notices a pattern: "order confirmed, never shipped."
Week 6 Reconciliation script compares orders vs shipments: 25,312 orphans.
Week 7 Manual replay; 4,100 customers already cancelled or charged back.
Detection: a reconciliation metric that should have existed — orders_without_shipment_after_1h — plus publish.after_commit_failures. Neither was alerting.
Blast radius: every downstream consumer of order events; refund and support cost ~$400K.
Mitigation: replay from a reconciliation diff; idempotent fulfilment made the replay safe.
Prevention: outbox in the order transaction; hourly reconciliation as a permanent backstop with a page above 0 orphans older than 1h.
Owner: order team owns the outbox; fulfilment owns the reconciliation alert on its side of the contract.
🎯 Staff Insight: The dual write did not fail loudly once — it failed quietly on every deploy. The tell in an interview is "it only happens on deploys." Deploys happen 20 times a day.
2. The Abandoned Slot — Primary Disk Full in 9 Hours#
A team prototypes a second Debezium connector against production for a search experiment, then deletes the connector but not its replication slot.
Fri 17:00 Experiment connector deleted. Slot 'search_poc' remains, inactive.
Fri 17:00 WAL generation ~60 MB/s at peak. Slot pins every segment from now on.
Sat 02:00 Retained WAL 1.6 TB of 2 TB volume. Disk alert at 80% fires to a low-urgency channel.
Sat 02:20 Disk full. PostgreSQL PANICs on WAL write; primary refuses all writes.
Sat 02:20 Checkout, accounts, payments down: they share the database.
Sat 03:05 On-call identifies the inactive slot, drops it, frees WAL, restarts.
Sat 03:30 Writes recover. 70 minutes of full write outage.
Detection: pg_replication_slots with active = false for > 15 min (page); replication_slot.retained_bytes > 20 GB (warn), > 100 GB (page); disk free-space forecast < 4h (page).
Blast radius: the entire primary — every service sharing it — caused by a component that was supposed to be read-only.
Mitigation: drop the slot; restart; re-snapshot the experiment if still wanted.
Prevention: max_slot_wal_keep_size on every primary; slots created only through the CDC platform with an owner label; inactive-slot alert.
Owner: CDC platform owns slot lifecycle; the database team owns the cap and the disk alert.
3. The Rebalance Reorder — Cancelled Orders Shipped#
A team moves from one poller to four to handle a sale, using FOR UPDATE SKIP LOCKED without hash-range ownership.
t=0 Sale starts: 4,000 order events/s. Four pollers race on the same outbox.
t=+3min Customer places and cancels order 77120 within 400ms.
Poller A locks OrderPlaced (v1); poller B locks OrderCancelled (v2).
B's batch is smaller and publishes first.
t=+3min Fulfilment receives v2 (cancel: no shipment exists, no-op), then v1 (create shipment).
t=+2d 1,940 cancelled orders shipped during the 6-hour sale.
Detection: consumer-side events.version_gap and events.stale_version counters; shipments_for_cancelled_orders reconciliation.
Blast radius: every quick place-then-cancel during the sale; ~$180K in shipped goods.
Mitigation: halt shipments for orders whose current status is cancelled; recall where possible.
Prevention: pollers own hash ranges of aggregate_id; consumers apply events only when aggregate_ver = current + 1 and park the rest for up to 30s before fetching current state.
Owner: order team owns relay ordering; fulfilment owns version checking.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Dual write divergence | Reconciliation: orphans_older_than_1h > 0 | Every downstream consumer, silently | Replay from diff; move to outbox | Producing team |
| Connector stopped | outbox.publish_lag_ms p99 > 5s; oldest unpublished row age | All consumers late; slot growth | Restart connector; failover worker | CDC platform |
| Slot retention runaway | replication_slot.retained_bytes > 100 GB; inactive slot > 15 min | Primary write outage | Drop or advance slot; WAL cap | CDC platform + DBA |
| Duplicate burst after failover | consumer.duplicates_skipped spike | None if consumers dedupe; double effects if not | Inbox dedupe in effect transaction | Each consumer team |
| Per-key reordering | events.version_gap > 0 | Wrong final state for affected aggregates | Hash-range relay ownership; version-gated apply | Producer + consumer |
| Breaking schema change | Registry compatibility check fails; consumer deserialization errors | All consumers of the topic | Block at CI; new event version alongside old | Producing team |
| Poison event | DLQ depth > 0; consumer stuck on one offset | One consumer's partition | DLQ with owner and replay tool | Consumer team |
| Outbox table bloat (polling) | Table size growth; vacuum lag | Primary performance | Partition by day, drop partitions | Producing team |
Beyond Staff: The Principal View#
Why L7 Sees This Problem Differently#
A Staff engineer removes the dual write from one service. A Principal engineer counts how many services publish events and finds that each has solved it differently: three hand-written pollers with different ordering bugs, two teams on table-level CDC whose consumers broke at the last migration, one team on Kafka transactions that believes it has exactly-once end to end, and a dozen services still committing then publishing. The org's real exposure is not any single relay; it is that "how a committed fact leaves a database" has no owner and no standard. CDC also quietly becomes tier-0: once search, cache invalidation, the warehouse and fraud scoring all depend on it, a connector outage is a multi-team incident and a runaway slot is a system-of-record outage. At L7 the work is making the outbox a paved road — one library, one connector platform, one event catalogue — and pricing CDC as the critical infrastructure it has become.
🧭 Principal Move: "We have 30 services emitting events and four different relay implementations. I'd fund one CDC platform with an outbox library on top, make idempotent consumption a review gate, and run the connector fleet as tier-0 with its own error budget — because it already is tier-0, we just haven't admitted it."
The Org-Level Fault Line#
One event platform with a mandated outbox contract vs per-team event publishing.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each team publishes however it likes | Autonomy; no platform dependency | Dual writes persist in long tail; incompatible event envelopes; nobody owns slot safety | Consumer teams reverse-engineering every producer; the DBA on the night a slot fills a disk |
| Central team owns all events and schemas | Consistency; one catalogue | Becomes a bottleneck for every new event; domain knowledge sits in the wrong team | Product velocity; the central team's backlog |
| Platform owns mechanism, domains own events | Shared outbox library, connector fleet, registry and lag alerts; teams own event design and compatibility | Platform must support PostgreSQL, MySQL and the managed stores everyone actually uses | Platform headcount (3–5 engineers) |
The Principal default is the third row: the platform guarantees "a row in your outbox reaches the bus within 5 seconds, at least once, in per-key order, without endangering your primary." Domain teams own what the events mean and keep them backward compatible.
Cost Model#
Assumptions: fully loaded engineer ~$25K/month; Kafka storage ~$0.30 per logical GB-month (3× replication); Connect worker ~$300/month (4 vCPU); 1 KB average event; 7-day topic retention.
| Scale | Event Volume | Relay Machinery | Infra Cost | People / On-call | Rough Monthly Total |
|---|---|---|---|---|---|
| Startup | 200 events/s, 3 services | Polling relay in each service; shared Kafka | Kafka ~$1K; ~120 GB retained ≈ $40 | 0.2 FTE; service teams on-call | ~$1K + ~$5K people |
| Growth | 10K events/s, 40 services | Debezium on Connect (12 workers), registry, lag alerts, outbox library | Connect ~$3.6K; ~6 TB retained ≈ $1.8K; Kafka brokers ~$15K | 2 FTE platform; shared rotation | ~$20K + ~$50K people |
| Large | 300K events/s, 400 services, 60 databases | Connector fleet per database cluster, automated slot management, replay tooling, event catalogue | Connect ~$40K; ~180 TB retained ≈ $55K; brokers ~$150K | 5-person platform team; dedicated tier-0 rotation | ~$250K + ~$125K people |
The Principal observation: the infrastructure is cheap relative to what it replaces. A single quarter of reconciliation engineering across 40 teams — each writing its own "find the orders that didn't sync" job — costs more than the growth-stage platform for a year. The expensive line item at scale is people on a tier-0 rotation, and that is the cost of admitting that CDC is load-bearing.
The 3-Year Evolution Path#
The Year 3 trigger is the one teams miss: by then cache invalidation, search and fraud all ride the same connectors. A two-hour connector outage is no longer "events are late"; it is stale prices, stale permissions and blind fraud scoring at the same time.
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Polling vs CDC relay | Two-way | The outbox table is the same; swap the relay behind it in a sprint |
| Event envelope (event_id, aggregate_id, version, schema_ver headers) | One-way | Every consumer parses it; changing it is an org-wide migration |
| Published event payload fields | One-way-ish | Additive changes are free; removing or renaming a field needs every consumer's sign-off and a deprecation window |
| Exposing table-level CDC of a core table to other teams | One-way | Your physical schema becomes their contract; you cannot refactor the table without a coordinated migration |
| Partition count and message key of an event topic | One-way-ish | Changing the key or count breaks per-key ordering during the transition |
| Outbox retention window | Two-way | Config change; drop partitions sooner or later |
| WAL cap per slot | Two-way | Config change; sets the trade between relay position and primary safety |
The Standard I'd Write#
RFC: Publishing Domain Events (v1)
Scope: Every service that changes state in a database and notifies another service or team of that change.
MUST:
- Events MUST be written in the same local transaction as the state change, via the platform outbox library or the store's native change stream. Publishing to a broker after commit from request code is prohibited.
- Every event MUST carry
event_id,aggregate_id,aggregate_version,event_typeandschema_version, and MUST be keyed byaggregate_id.- Payload schemas MUST be registered and pass backward-compatibility checks in CI. Removing a field requires a 90-day deprecation notice to registered consumers.
- Consumers MUST dedupe on
event_id(or gate onaggregate_version) in the same transaction as their side effect, and MUST own a DLQ with a replay path.- Every replication slot MUST be created through the CDC platform with an owner, and every primary MUST set a WAL retention cap.
SHOULD: Prefer fat-enough events that let consumers act without calling back; keep 72h of outbox partitions as a replay buffer; reconcile tier-1 flows hourly.
Exceptions: Filed with the event platform team, time-boxed to 2 quarters, with an owner and a reconciliation backstop.
Success metrics: zero incidents caused by dual writes; publish lag p99 < 5s on 99.9% of minutes; zero primary outages from slot retention; 100% of tier-1 consumers passing a duplicate-delivery test.
What I'd Tell the VP#
Today, when one of our systems records something important — an order, a payment, a refund — it tells the other systems in a separate step that can silently fail. Last quarter that cost us about 25,000 orders that never shipped and several hundred thousand dollars in refunds and support. I'm proposing a shared way for every team to record the change and the announcement as one step, plus a small platform team to run the piece that delivers those announcements. It costs about two engineers for two quarters and a modest infrastructure bill. In exchange, "we recorded it but nobody heard" stops being a class of incident, and we stop paying 40 teams to write their own cleanup scripts.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Paved road over heroics | "Four relay implementations means four ordering bugs. One library, one connector fleet." |
| Admits the tier | "CDC feeds cache invalidation and fraud now. It's tier-0 whether we staff it or not." |
| Contract thinking | "The envelope is a one-way door. Payload fields evolve additively with a 90-day deprecation." |
| Protects the system of record | "Every slot has an owner and every primary has a WAL cap. A relay must never take down its source." |
| Prices the alternative | "Forty teams writing reconciliation jobs costs more than the platform." |
Staff answers that L7 interviewers find insufficient:
- "I'd add an outbox to this service." — Correct locally; silent on the 29 other services and who owns the relay they'll all share.
- "We'll alert on connector lag." — Doesn't say who is on call for it or whether it's staffed as the tier-0 dependency it has become.
- "Consumers should be idempotent." — Doesn't make it enforceable: no envelope standard, no review gate, no duplicate-delivery test.
How Real Companies Built It#
Debezium: The Outbox Event Router as a Documented Pattern#
Debezium, the open-source CDC project, documents the outbox pattern as a first-class feature. Its outbox event router expects an outbox table with an event id, an aggregatetype used for routing (by default to a topic named outbox.event.<aggregatetype>), an aggregateid used as the Kafka message key so per-aggregate order is preserved within a partition, and a payload. The router treats the table as insert-only and filters out deletes. The project's introductory post on the pattern goes further: it persists the outbox event and removes it in the same transaction, because CDC tails the append-only log rather than reading the table, so the insert is captured even though the table stays empty. The post also states the root problem plainly: Kafka cannot be enlisted in a distributed XA transaction with the service's database (Debezium outbox event router, Reliable microservices data exchange with the outbox pattern).
Staff insight: This is the canonical reference for "outbox plus log-based relay," and the insert-then-delete trick is a sharp detail to drop in an interview: it shows you understand that CDC reads the log, not the table, which is also why an unconsumed slot retains log on the primary.
Netflix: DBLog and Watermark-Based Full-State Capture#
Netflix's DBLog is a CDC framework built to keep heterogeneous datastores in sync by combining transaction-log events with rows selected directly from tables, so that a consumer can receive the full state of a table and not just changes since the connector started. Its key mechanism is a watermark: the framework interleaves chunked table selects with log events, which lets dumps be paused, resumed, and triggered for all tables, one table, or specific primary keys, without stalling log processing. The paper states that the approach uses no locks and has minimal impact on the source, and that DBLog runs in production for tens of microservices at Netflix (DBLog: A Watermark Based Change-Data-Capture Framework).
Staff insight: Every CDC design eventually needs a backfill — a new consumer, a corrupted index, a slot that was invalidated. The naive approach (stop the world and snapshot) blocks writes or races with live changes. Naming lock-free, chunked, watermark-interleaved backfill shows you've thought about day 400 of the pipeline, not just day one.
Shopify: Log-Based CDC Across a Sharded Monolith#
Shopify replaced a query-based batch extraction system with log-based CDC: Debezium reads MySQL binary logs, Kafka Connect manages the connectors, and events land in Kafka topics partitioned by primary key, with a Kafka Streams application consolidating per-shard topics into one logical topic per table and a schema registry managing Avro schemas. The write-up describes 100+ MySQL shards, roughly 150 Debezium connectors, 400 TB+ of CDC data in Kafka, p99 latency under 10 seconds from MySQL insert to Kafka availability, and ~65,000 records/s on average during Black Friday 2020 with spikes to 100,000. Their stated reasons for leaving query-based extraction: it missed hard deletes, could not see intermediate row states, and relied on every write correctly updating updated_at (Capturing every change from Shopify's sharded monolith).
Staff insight: This is the honest case for table-level CDC as a data-platform feed — and the honest case against updated_at polling as a substitute. When an interviewer suggests "just query for recently updated rows," the three misses Shopify listed are the answer.
Practice Drill#
Prompt: "A marketplace's listing service writes listings to PostgreSQL and then calls the search team's indexing API and publishes a
ListingChangedevent that pricing, recommendations and fraud consume. Sellers complain that edited prices take hours to show in search, and fraud found that some delisted items were still being recommended. The listing service does 3,000 writes/s at peak. Redesign how changes leave the listing service."
Staff Answer
Both symptoms are dual writes: the commit, the indexing call and the publish are three independent operations, and every crash, timeout or deploy between them leaves one of the three behind. I'd make the database transaction the only atomic boundary. Listing updates insert a ListingChanged row into an outbox table in the same transaction — event_id, listing_id as the aggregate ID, aggregate_version, schema_version, and a fat-enough payload (price, status, category, seller ID) that pricing and fraud can act without calling back. A Debezium connector reads the outbox from the WAL and routes it to a listings.events topic, 48 partitions, keyed by listing_id, so all changes to one listing stay in commit order. Search stops being called synchronously: it's the listing team's own read model in practice, but it's another team's system, so it consumes the same topic and upserts with WHERE version < $new_version, which makes duplicates and stale events no-ops. Delisting is just an event with status = DELISTED, so recommendations and fraud get it through the same path — no separate call to forget. Each consumer dedupes on event_id in the transaction that applies its effect and owns a DLQ. Numbers: 3,000 events/s × ~1 KB is ~3 MB/s, trivial for one connector; expected publish lag 200ms–1s, alert at p99 > 5s; outbox in daily partitions kept 72h as a replay buffer (~780 GB at peak rate — acceptable, or 24h if disk is tight); max_slot_wal_keep_size at 200 GB with a page at 100 GB retained. Backfill for search uses a chunked, watermark-style snapshot so we can rebuild the index without stopping writes. Hourly reconciliation compares listing versions in PostgreSQL to search's indexed versions and pages on any listing more than 10 minutes behind. Owners: listing team owns the outbox schema in the registry; platform owns the connector and slot; each consumer owns dedupe and DLQ.
Why this is L6:
- Names the atomic boundary and removes all post-commit side effects from request code, including the synchronous indexing call.
- Pushes exactly-once-in-effect to consumers with version-gated upserts and inbox dedupe, and keeps per-listing order by keying on listing ID.
- Protects the primary from its own relay (WAL cap, slot alert) and plans the backfill and reconciliation backstop with owners.
What L7 adds:
- Notices pricing, search and fraud will each want this from many services, and proposes the outbox library and connector fleet as a shared platform rather than a listing-team project.
- Treats the
ListingChangedenvelope as a one-way door and sets an org deprecation policy for payload fields. - Prices it: one connector worker and ~1 TB of storage versus the support and fraud cost of stale listings, and classifies the CDC path as tier-1 because fraud now depends on it.
Staff Interview Application#
How to Introduce This Pattern#
"The only atomic thing I have is the database transaction, so the event goes in it — an outbox row next to the state change. A relay reads it from the commit log and publishes keyed by the entity ID. Delivery is at-least-once, so every consumer dedupes on the event ID in the same transaction as its own write. The thing I'll watch is relay lag in seconds and the replication slot, because a stuck relay can fill the primary's disk."
Lead with the atomic boundary, then the relay, then consumer idempotency, then the operational risk and who owns it.
When NOT to Use This Pattern#
- Nobody else needs to know promptly: a nightly report or monthly export can scan for changes. Say what it misses (hard deletes, intermediate states) and accept it explicitly.
- The store already has a native change stream that fits: DynamoDB Streams, a Kafka-native event-sourced service. The stream is the outbox; don't add a table.
- One process, one database: if the "event" is consumed by the same service in the same database, write both rows in one transaction and skip the bus entirely.
- Loss-tolerant telemetry: click logs and metrics go straight to the log with batching; nobody needs per-event atomicity with a database row.
- Cross-service business transactions needing compensation: the outbox gets each step's event out reliably, but the coordination is a separate design — see Workflows, Sagas & Compensation and Coordination Strategies.
Follow-Up Questions to Anticipate#
| Interviewer Asks | What They Are Testing | How to Respond |
|---|---|---|
| "Why not just retry the publish?" | Understanding of where durability lives | "The retry lives in memory in a process that can die after commit. The intent to publish has to be durable in the same commit as the change." |
| "Doesn't Kafka have exactly-once?" | Precision about semantics | "Within Kafka. It doesn't span PostgreSQL. The relay is at-least-once; consumers dedupe on event ID inside their own transaction." |
| "Outbox or plain CDC on the table?" | Coupling and contracts | "Table CDC for my own read models. Other teams get outbox events so my schema stays mine." |
| "What happens if the connector is down for a day?" | Operational depth | "Events are late, not lost, as long as the slot holds. The risk is WAL growth on the primary, so there's a cap, a page at 100 GB, and a replay path from retained outbox partitions if the slot is invalidated." |
| "How do consumers handle out-of-order events?" | Ordering model | "Per-key order only, via the message key. Consumers gate on aggregate version and park or refetch on a gap." |
| "How do you onboard a new consumer that needs all history?" | Backfill design | "Chunked snapshot interleaved with the live log — a watermark approach — or replay from a compacted topic. Never a stop-the-world dump." |
Scorecard#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Publishes after commit with retries | Identifies the transaction as the atomic boundary; outbox or CDC | Treats event publication as an org contract and platform |
| Semantics | Assumes exactly-once from the broker | At-least-once relay; consumer dedupe in the effect transaction | Makes idempotent consumption enforceable via envelope and review gates |
| Ordering | Assumes global order | Per-key order via message key; version-gated apply | Decides ordering contracts per domain; one-way door awareness on keys |
| Operations | Watches broker health | Lag in seconds, slot retention cap, reconciliation backstop | Runs CDC as tier-0 with error budget, game days and replay tooling |
| Ownership | Implicit | Producer owns schema, platform owns relay, consumers own DLQs | Funds the platform; sets deprecation and exception policy |
Strong Hire Signals
| Signal | What It Sounds Like |
|---|---|
| Names the dual write | "Commit then publish is two commits. One will eventually fail alone." |
| Places exactly-once correctly | "Exactly-once in effect, at the consumer, with the dedupe row in the same transaction." |
| Knows the slot risk | "A dead connector pins WAL. At 50 MB/s we have about 11 hours." |
| Separates contract from storage | "Outbox payload is versioned; the table can change under it." |
Lean No-Hire Signals
| Signal | Why It Misses the Bar |
|---|---|
| "Wrap the DB write and Kafka send in a transaction" | Doesn't know the broker can't join one; will ship a dual write |
| No answer for duplicates | At-least-once is guaranteed by every relay; double charges follow |
| Lag alert on row count only | Can't tell 2 seconds of lag from 2 hours, and misses the slot entirely |
Common False Positives: Knowing Debezium config keys ≠ knowing why the outbox exists. Saying "event-driven" ≠ having an atomic boundary. Drawing a saga ≠ getting each step's event out of the database reliably.
Capacity Planning Quick Reference#
Sizing the Relay#
outbox_rows_per_day = writes_per_sec × events_per_write × 86,400
outbox_storage = rows_per_day × avg_row_bytes × retention_days # 2K/s × 1KB × 3d ≈ 520 GB
slot_time_to_disk_full = free_disk_bytes / wal_bytes_per_sec # 2 TB / 50 MB/s ≈ 11 h
publish_lag = commit_to_relay_read + relay_batch_wait + broker_ack # ~100 ms – 2 s healthy
polling_avg_added_lag = poll_interval / 2 # 200 ms poll → ~100 ms
inbox_storage = events_per_sec × 86,400 × replay_days × ~60 bytes # 2K/s × 7d ≈ 70 GB
duplicate_window = events published between last checkpoint and crash # bounded by flush interval
topic_partitions ≥ max_consumer_parallelism needed for catch-up # plan 3× steady state
Key Numbers Worth Memorizing#
| Number | Context |
|---|---|
| ~0.1–0.5 ms | Added write latency of one outbox insert in the same transaction |
| 100 ms – 2 s | Healthy log-based relay lag from commit to broker |
| p99 > 5 s | Reasonable page threshold for publish lag on domain events |
| 100–500 ms | Typical polling interval; average added lag is half of it |
| ~2–5K events/s | Ceiling of a single well-batched polling publisher |
| 24 hours | DynamoDB Streams record retention |
| 72 h | Sensible outbox partition retention as a replay buffer under CDC |
| 7 days | Common topic retention, sets the consumer dedupe window |
| 2 TB ÷ 50 MB/s ≈ 11 h | Time a dead slot takes to fill a primary's disk |
| 0.01% | Dual-write divergence that looks rare and is 17K records/day at 2K writes/s |
Common Pitfalls Checklist#
- No request-handling code publishes to a broker after a database commit
- Outbox insert shares the business transaction; payload built from the written state
- Every event carries
event_id,aggregate_id,aggregate_version,schema_version - Events keyed by aggregate ID; nobody relies on global order
- Consumers dedupe or version-gate inside the transaction that applies the effect
- Every replication slot has an owner, and the primary has a WAL retention cap
- Lag alerted in seconds; inactive-slot and retained-bytes alerts page
- Outbox retention by partition drop, not row deletes
- Payload schemas registered with compatibility checks in CI
- Hourly reconciliation backstop on tier-1 flows, with a page on any orphan older than 1h