Design a Transactional Outbox and Event-Driven Ledger
This page is one interview loop in three rounds. All three rounds design the same system. Each round opens with the interviewer raising the scope, and the design from the round before has to evolve to meet it.
| Round 1: Mid-level | Round 2: Senior | Round 3: Architect | |
|---|---|---|---|
| Story | An order service must tell fulfillment and email about every paid order | A ledger service (balances and transfers) announces every movement to billing, analytics and audit | 40 services in 3 regions share one event backbone; events are contracts |
| Level (Amazon) | SDE II (L5) | Senior SDE (L6) | Principal (L7) |
| Volume | 1M events/day; ~50/s in the busiest minute | 50M events/day: ~579/s average, 6,000/s peak; 30K balance reads/s | 2B events/day: ~23K/s average, 70K/s peak |
| Footprint | 1 region; one PostgreSQL database; one Kafka cluster | 1 region, 3 AZs; Aurora PostgreSQL + Amazon MSK | 3 regions; 40 service databases; one MSK cluster per region |
| Targets | No lost events, no ghost events; delivered within 5 s; 99.9% | Commit P99 < 15 ms; dispatch P99 < 100 ms; no table bloat; 99.99% | Contracts that never break consumers; replay a year of events safely |
| Reading time | ~35 min | ~40 min | ~45 min |
You can start at any round. Rounds 2 and 3 open with a "Where we left off" summary that catches you up.
Loop Opener: What Is the Transactional Outbox?
You Already Know One: a Letter Tray Inside the Vault
Picture a bank clerk. When she moves money, she writes the entry in the ledger book, and in the same motion she drops a notification letter into a tray that sits inside the vault. Later a courier empties the tray and delivers the letters. If the vault door never closes (she tears up the page because the customer changed their mind), the letter is torn up with it. The ledger and the letters can't disagree, because they were written in one place, in one motion.
| Thing | In the bank picture | In our system |
|---|---|---|
| The ledger book | The record of what happened | The service's own tables (orders, balances, ledger entries) |
| The letter tray | Letters waiting to go out, kept in the vault | An outbox table in the same database |
| "One motion" | Writing the entry and the letter together | One database transaction that writes both |
| The courier | Takes letters out of the tray and delivers them | A relay (or a log reader) that publishes outbox rows to a message broker |
| The recipient | Reads the letter and acts on it | A consumer service, which must cope with a letter arriving twice |
What Makes It Hard
- Two systems can't commit together. The database and the message broker (Kafka, SQS) each have their own idea of "saved". There is no practical way to make one commit wait for the other.
- Networks time out. A publish that times out may or may not have happened. A commit whose reply is lost may or may not have happened.
- Consumers see things twice. Every reliable delivery system retries, and every retry can repeat something that already worked.
- Order matters per thing. Fulfillment must not see "shipped" before "paid" for the same order, even though thousands of orders are in flight at once.
We protect three invariants throughout the loop:
- No lost events: every committed change produces its event, eventually.
- No ghost events: no event is ever published for a change that rolled back.
- Effects happen once: each consumer acts on each event exactly once in effect, even though the event may be delivered more than once, and in order for each entity.
The Question the Whole Loop Answers
How do we make a data change and the event about it always agree, and make every consumer act on it exactly once in effect?
The answer gets sharper every round:
- Round 1: an outbox table written in the same transaction, a relay that publishes and deletes, idempotent consumers, and per-order ordering.
- Round 2: a ledger at 6,000 events a second: fast dispatch, correct locking for money, poison events, broker outages, and a table that never bloats.
- Round 3: a company-wide backbone: event contracts, change data capture at scale, replayable history, regions, and an honest answer to "is it exactly-once?"
The outbox is the quiet pattern under several other loops. The payment processing loop and the digital wallet loop use it to publish ledger changes, and the notification system loop relies on its producers using it. Here we teach it on its own, with its trade-offs in full.
Round 1 · Mid-level · "Order Paid → Fulfillment Must Know"
~35 min · SDE II (L5) · 1 region · 1M events/day, ~50/s in the busiest minute · delivered within 5 s · 99.9%
R1.1 Establish Design Scope
The interviewer says: "When an order is paid, fulfillment has to start packing and the customer gets a confirmation email. Today the order service writes the order and then calls Kafka, and sometimes fulfillment never hears about a paid order. Fix it." We ask before we draw.
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| Who consumes the event? | Fulfillment and the email service, owned by other teams. | Two independent consumers, each reading at its own pace: a broker with consumer groups, not a direct call. |
| How fast must they see it? | Within seconds. | A relay that polls once a second is fast enough this round. Milliseconds come in Round 2. |
| Can an event be delivered twice? | Yes, as long as nothing happens twice because of it. | We promise at-least-once delivery and make consumers idempotent (step 1.3). |
| Must events for one order arrive in order? | Yes. "Shipped" before "paid" makes fulfillment do the wrong thing. | Per-order ordering (step 1.5). Ordering across different orders doesn't matter. |
| Do we already run a broker? | Yes: Kafka on Amazon MSK, run by the platform team. | We publish to a topic; we don't build a queue. |
| How many orders? | About 250,000 paid orders a day. Each order produces about 4 events: placed, paid, shipped, cancelled or delivered. | About 1M events a day (R1.7). |
Out of scope for this round:
- Millisecond dispatch.
- Heavy write contention on the same rows.
- Changing event shapes without breaking consumers.
R1.2 Functional Requirements, Derived Step by Step
| Phrase from the problem | Requirement |
|---|---|
| "When an order is paid" | Changing the order's state and recording the event about it happen atomically: both or neither |
| "Fulfillment has to start packing" | Every event reaches every subscribed consumer, eventually |
| "Sometimes fulfillment never hears" | No event is lost, even if the process crashes or the broker is down |
| "Nothing happens twice" | Consumers apply each event's effect once, even when it's delivered twice |
| "Shipped before paid is wrong" | Events for one order are delivered and applied in order |
Not yet: millisecond dispatch, contention on hot rows, schema evolution.
R1.3 Non-Functional Requirements: the Questions
- No lost events. A committed order change must produce its event, whatever crashes.
- No ghost events. An event for a change that rolled back must never be published. A ghost "paid" event ships goods nobody paid for.
- Bounded delay. P99 under 5 seconds from commit to the consumer applying the event, in normal operation.
- The database must not bloat. The outbox must stay small, not grow by a row forever for every event.
- Availability. 99.9% for the order write path: about 43.8 minutes a month. A broker outage must not take the write path down.
R1.4 The API
The client-facing API doesn't change much: the event is a side effect of a normal write.
Pay for an order
httpPOST /v1/orders/o_48213/pay HTTP/1.1 Host: api.shop.example Authorization: Bearer <access token> Idempotency-Key: 2b7e4c1a-9f3d-4a8e-b6c2-5d1e0f7a3b94 Content-Type: application/json { "payment_method_id": "pm_771", "amount_minor": 4599, "currency": "USD" }
httpHTTP/1.1 201 Created Content-Type: application/json { "payment_id": "pay_90311", "order_id": "o_48213", "order_status": "PAID", "order_version": 2 }
201 Created because a payment resource was created. A retry with the same Idempotency-Key returns the same body (the payment processing loop covers how). order_version goes up by one on every change to the order.
The event envelope. Every event, whatever its type, travels in the same outer shape. Consumers read the envelope first and only then the payload.
json{ "event_id": "0192a6f4-7c1e-7b3a-9d42-5e8f1a2b3c4d", "event_type": "order.paid", "schema_version": 1, "aggregate_type": "order", "aggregate_id": "o_48213", "aggregate_version": 2, "occurred_at": "2026-09-27T12:01:07.412Z", "producer": "order-service", "correlation_id": "req_5c1d9e", "payload": { "order_id": "o_48213", "customer_id": "c_1207", "amount_minor": 4599, "currency": "USD", "items": [ { "sku": "SKU-3310", "qty": 1 } ] } }
| Field | Why it's there |
|---|---|
event_id | A unique ID per event (a UUID). Consumers use it to spot duplicates |
event_type, schema_version | What happened, and which version of its shape, so consumers can parse it safely |
aggregate_type, aggregate_id | The aggregate is the entity the event is about (here, one order). Ordering and partitioning are per aggregate |
aggregate_version | The order's version after this change: 1, 2, 3 ... with no gaps. Lets consumers detect duplicates and out-of-order delivery (step 1.5) |
occurred_at | When the change committed, for humans and reports. Never used for ordering (clocks differ between machines) |
correlation_id | Ties the event back to the request that caused it, for tracing |
Recap
- One write endpoint; the event is a side effect.
- Every event has an ID, an aggregate ID and a gap-free aggregate version.
- About 1M events a day, delivered within seconds, at least once, in order per order.
Let's build it, starting with what the team has today.
R1.5 Design Evolution: From "Commit, Then Publish" to an Outbox
Every step follows the same pattern: a problem, your turn to think, the answer, and what it costs us. The cost is usually the next problem.
Step 1.0: The Baseline
The order service commits the order, then publishes an event to Kafka.
Synthesizing vector architecture diagram...
Two separate writes to two separate systems. Nothing ties write 2 to write 1.
What's good about it: it's two lines of code, and it usually works.
What it costs us: this is a dual write, two systems written one after the other with no shared transaction. When the second write fails, the two systems disagree forever.
Step 1.1: "The Order Committed but the Publish Timed Out: the Event Is Lost"
The problem: the order commits as PAID. The Kafka publish times out, and then the pod is killed by a deploy before it retries. Fulfillment never hears about the order. The customer paid and nothing ships.
What would you do?
The failure the outbox removes, as a timeline:
Synthesizing vector architecture diagram...
Dual write: the database is right, the stream is wrong, and nothing will ever notice.
With the outbox, the write is one transaction:
sqlBEGIN; UPDATE orders SET status = 'PAID', version = version + 1 WHERE order_id = 'o_48213' AND status = 'PLACED' RETURNING version; -- returns 2 INSERT INTO outbox (event_id, bucket, aggregate_type, aggregate_id, aggregate_version, event_type, schema_version, payload) VALUES ('0192a6f4-...', 5, 'order', 'o_48213', 2, 'order.paid', 1, '{...}'); COMMIT;
The UPDATE comes first on purpose: it locks the order row, so a second change to the same order waits here until this transaction commits. That matters for ordering in step 1.5.
Primitive: Change Data Capture and the Outbox Pattern · Drill: The dual-write that broke search consistency (the commit-succeeded-publish-failed question is answered here; the CDC-versus-polling question in step 3.2 ('Why CDC instead of a relay polling every 500 ms?'), with the comparison table in step 2.1)
Step 1.2: "Who Moves Events From the Outbox to the Broker?"
The problem: rows pile up in outbox. Something must get them to Kafka.
What would you do?
Step 1.3: "The Relay Crashed After Publishing but Before Deleting"
The problem: the relay publishes 50 events, Kafka acknowledges them, and the relay is killed before its DELETE commits. The rows are still in outbox. When the relay restarts, it publishes them again.
What would you do?
Consumer pseudocode (fulfillment):
textfor each record polled from orders.events (auto-commit off): BEGIN INSERT INTO inbox (event_id, received_at) VALUES (record.event_id, now()) ON CONFLICT (event_id) DO NOTHING if no row was inserted: COMMIT, then commit the offset, continue -- duplicate apply the effect (create or update the fulfillment job) COMMIT commit the Kafka offset for this record
Primitive: Database Isolation Levels, ACID & Concurrency Anomalies
Step 1.4: "Two Relay Instances Published the Same Rows"
The problem: we run two relay tasks so one can fail. Both run SELECT ... ORDER BY seq LIMIT 50 at the same moment, both get the same 50 rows, and both publish them.
What would you do?
SKIP LOCKED deliberately gives an inconsistent view of the table. That's exactly right for handing out work, and exactly wrong for anything that must see every row, like a report.
Step 1.5: "An Order's 'Shipped' Arrived Before Its 'Paid'"
The problem: with two relays claiming batches independently, version 3 (order.shipped) of an order reached Kafka before version 2 (order.paid). Fulfillment saw "shipped" for an order it never saw paid.
What would you do?
The relay loop:
textevery 1 s, and again immediately while work was found: BEGIN b = SELECT bucket FROM outbox_buckets WHERE EXISTS (SELECT 1 FROM outbox o WHERE o.bucket = outbox_buckets.bucket) ORDER BY random() LIMIT 1 FOR UPDATE SKIP LOCKED if no bucket: COMMIT and sleep rows = SELECT * FROM outbox WHERE bucket = b ORDER BY seq LIMIT 200 send each row to Kafka topic orders.events, key = aggregate_id wait for every acknowledgement (5 s timeout: ROLLBACK and back off) DELETE FROM outbox WHERE seq = ANY(the published seqs) COMMIT
The relay holds a database transaction open while it publishes. That breaks our usual rule of "no network calls inside a transaction", and it's the one exception we allow: the transaction locks only one bucket row and outbox rows (never an order), it runs on the relay's own two connections, and the publish has a hard timeout. A slow broker delays events; it can't block orders.
Synthesizing vector architecture diagram...
Both versions of order o_48213 live in bucket 5, so exactly one relay publishes them, in order. Relay B works on a different bucket in parallel.
Step 1.6: "The Outbox Table Grew to Millions of Rows"
The problem: a teammate changed the relay to mark rows sent = true "so we keep a history". Six months later the outbox has 180 million rows and the relay's query is slow.
What would you do?
The growth math:
| Design | Rows kept | Size |
|---|---|---|
| Mark processed, keep forever | 1M a day | 1,000,000 × 750 B = 750 MB/day ≈ 274 GB/year |
| Delete on success, 1 s polling, at 50/s | ~50 waiting + a batch in flight | ≈ 50–100 rows × 750 B ≈ 40–75 KB |
Per-table autovacuum settings make vacuum run on this table after a fixed number of dead rows, not a percentage of a tiny table:
sqlALTER TABLE outbox SET (autovacuum_vacuum_scale_factor = 0, autovacuum_vacuum_threshold = 1000);
Round 1 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.0 | (baseline) | Commit, then publish | Dual write |
| 1.1 | Lost event after commit | Outbox row in the same transaction | A table to drain |
| 1.2 | Moving rows to Kafka | Relay: read, publish, delete | Up to 1 s of delay; a new process |
| 1.3 | Crash between publish and delete | At-least-once + inbox in the consumer's transaction | An inbox per consumer |
| 1.4 | Two relays publish the same rows | FOR UPDATE SKIP LOCKED claims | Order across relays |
| 1.5 | Out-of-order events | Buckets claimed whole; key by aggregate; version check at the consumer | Parallelism capped by bucket count |
| 1.6 | Outbox growth | Delete on success; per-table autovacuum | A dead tuple per event to vacuum |
R1.6 Architecture v1
Synthesizing vector architecture diagram...
Follow an event left to right: it's born inside the order transaction, waits in the outbox until a relay claims its bucket, and reaches each consumer group separately. Each consumer deduplicates in its own database.
The pieces:
- Order service: stateless tasks on ECS Fargate. It only ever writes to its own database.
- Aurora PostgreSQL: one writer and one reader in a second AZ for failover. The outbox lives here, next to the orders.
- Relay: two small Fargate tasks with their own two-connection pool.
- Amazon MSK: the platform team's Kafka cluster. Our topic has 12 partitions and replication factor 3; the producer uses
acks=all, so an acknowledged record is on every in-sync replica. - Consumers: each is a Kafka consumer group; each group gets every event, and within a group each partition is read by one member at a time.
Schemas
sql-- Order service database CREATE TABLE orders ( order_id TEXT PRIMARY KEY, -- 'o_48213' order_num BIGINT NOT NULL UNIQUE, -- numeric part, used for the bucket customer_id TEXT NOT NULL, status TEXT NOT NULL CHECK (status IN ('PLACED','PAID','SHIPPED','DELIVERED','CANCELLED')), version BIGINT NOT NULL DEFAULT 1, -- +1 on every change; becomes aggregate_version amount_minor BIGINT NOT NULL, updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE TABLE outbox ( seq BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, event_id UUID NOT NULL, bucket SMALLINT NOT NULL CHECK (bucket BETWEEN 0 AND 15), -- order_num mod 16 aggregate_type TEXT NOT NULL, aggregate_id TEXT NOT NULL, aggregate_version BIGINT NOT NULL, event_type TEXT NOT NULL, schema_version INT NOT NULL, payload JSONB NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), UNIQUE (aggregate_type, aggregate_id, aggregate_version) -- a version is announced once ); CREATE INDEX outbox_by_bucket ON outbox (bucket, seq); CREATE TABLE outbox_buckets (bucket SMALLINT PRIMARY KEY); -- 16 rows: 0..15 -- Fulfillment's database (the email service has the same inbox) CREATE TABLE inbox ( event_id UUID PRIMARY KEY, received_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX inbox_by_time ON inbox (received_at); -- for the daily prune CREATE TABLE fulfillment_orders ( order_id TEXT PRIMARY KEY, last_order_version BIGINT NOT NULL, -- the per-aggregate high-water mark; never pruned state TEXT NOT NULL );
Two design notes:
- The inbox is not partitioned. It would be tempting to partition it by day and drop old days, but PostgreSQL requires a unique constraint on a partitioned table to include the partition key, so
event_idalone could no longer be unique across days. We prune it with a batched dailyDELETEinstead (R1.7). - The consumer's version high-water mark lives on its own row for the order, which is never pruned while the order matters. So even after an inbox row is deleted, a very late duplicate is still rejected by the version check.
Trace 1: the happy path
Synthesizing vector architecture diagram...
The client gets its answer as soon as the order commits; the event follows within about a second.
Trace 2: the relay crashes after publishing. The relay publishes v2, Kafka acknowledges, and the relay is killed before its DELETE commits. Its transaction rolls back; the row and bucket 5 are free again. A second later the other relay claims bucket 5 and publishes v2 again. Fulfillment receives v2 twice.
Trace 3: the duplicate is ignored. The second copy of v2 reaches fulfillment. Its INSERT INTO inbox hits the existing event_id and inserts nothing, so fulfillment commits without doing anything and moves its offset forward. Even without the inbox, the version check would skip it: v2 is not greater than the stored last_order_version of 2.
R1.7 Numbers
Traffic
| Quantity | Math | Value |
|---|---|---|
| Events per day | 250,000 orders × 4 events | 1,000,000 |
| Average rate | 1,000,000 ÷ 86,400 s | ≈ 11.6/s |
| Busiest minute | assume 4× the daily average (evening peaks, promotions): 11.6 × 4 ≈ 46 | plan for 50/s |
The outbox
| Quantity | Math | Value |
|---|---|---|
| Row size | payload ~500 B + other columns ~130 B + row header and index entries ~120 B (estimate) | ≈ 750 B |
| Rows waiting at peak, 1 s polling | 50/s × 1 s, plus a batch in flight | ≈ 50–100 rows ≈ 40–75 KB |
| If kept forever | 1M × 750 B = 750 MB/day; × 365 | ≈ 274 GB/year |
| One-hour broker outage at peak | 50/s × 3,600 s = 180,000 rows × 750 B | ≈ 135 MB, then drained |
Relay load. When idle, each relay runs one claim query a second: 2 queries/s. At peak, 50 events a second spread over 16 buckets: the chance a given bucket gets none of 50 events in a second is , so about 15 buckets have work each second. Each costs about four statements (claim, read, delete, commit): about 60 statements/s, plus the idle claims. That's well under 100 statements a second: noise for the database.
Drain after an outage. A relay batch of 200 rows takes roughly 30 ms (read, publish with acks=all, delete, commit; a rough figure to confirm by test), so one relay drains about 6,600 rows/s and two about 13,000/s. The 180,000 rows from a one-hour peak outage drain in about 180,000 ÷ 13,000 ≈ 14 s.
The inbox. We keep inbox rows for 8 days: one more than the topic's 7-day retention, because a consumer that rewinds its offsets can only re-read what Kafka still holds. At about 120 bytes per row including its index entries (an estimate): 1M rows/day × 8 days × 120 B ≈ 0.96 GB per consumer. The daily prune deletes about 1M rows in batches of 10,000 (100 short statements), so it never holds locks for long.
Monthly cost (us-east-1 on-demand list prices, 730 hours a month; the consumers' own costs belong to their teams)
| Item | Math | Monthly |
|---|---|---|
Aurora writer + reader, db.r7g.large, Aurora Standard | 2 × $0.276/h × 730 h | ≈ $403 |
| Aurora I/O | assume ~40 I/Os per order: 250K × 40 × 30.4 days ≈ 304M × $0.20 per million | ≈ $61 |
| Aurora storage | orders ~0.5 GB/day → ~180 GB after a year × $0.10/GB-month | ≈ $18 |
MSK, 3 × kafka.m7g.large | 3 × $0.204/h × 730 h | ≈ $447 |
| MSK storage | 3 brokers × 100 GB × $0.10/GB-month | $30 |
| Fargate: order service, 3 tasks × 0.5 vCPU, 1 GB, Arm | 3 × (0.5 × $0.03238 + 1 × $0.00356)/h × 730 h | ≈ $43 |
| Fargate: relay, 2 tasks × 0.25 vCPU, 0.5 GB, Arm | 2 × (0.25 × $0.03238 + 0.5 × $0.00356)/h × 730 h | ≈ $14 |
| Application Load Balancer | $0.0225/h × 730 h + a few capacity units | ≈ $22 |
| CloudWatch metrics, alarms, logs | ≈ $20 | |
| Total | ≈ $1,060/month |
The MSK cluster is the biggest line, and in practice it's shared with other teams, so our share is smaller. The outbox itself costs almost nothing: a few kilobytes of table and two tiny relay tasks.
MSK storage check: 1M events × ~1 KB × 7 days × 3 replicas ≈ 21 GB across the cluster, well inside 300 GB.
R1.8 Trade-Offs
Polling vs a push signal
| Polling every second (our choice) | A push signal (wake the relay when a row is written) | |
|---|---|---|
| Delay | Up to one interval (1 s), about 0.5 s on average | Milliseconds |
| Load when idle | One cheap query per relay per second | None |
| What can go wrong | Nothing new | Signals get lost, so a polling sweep is still needed as a backstop |
| Fits this round? | Yes: "within seconds" | Not needed yet; Round 2 needs it |
Delete vs mark processed
| Delete on success (our choice) | Mark sent = true | |
|---|---|---|
| Table size | Only pending rows (kilobytes) | Grows by every event forever, unless something else deletes |
| Dead tuples | One per event, on a tiny table | One per event on a growing table, plus the live history |
| History of events | In the broker and archive | In the table (but in the wrong place: the hottest database) |
Outbox vs two-phase commit with the broker. Two-phase commit (2PC) asks every participant to "prepare" (promise it can commit), then tells all of them to commit. Kafka only recently added opt-in 2PC participation (KIP-939, transaction.two.phase.commit.enable, off by default and not something we can assume on our cluster); even with it, a prepared transaction holds its row locks until the coordinator decides, so a coordinator outage freezes those orders. The outbox gets the same "both or neither" guarantee from one local transaction, and pushes the cross-system step into a retryable relay whose only failure mode is delay or duplicates. See Two-Phase Commit and Saga Orchestration.
Kafka vs a queue. We use the Kafka cluster we already have. A pay-per-request queue service could carry 1M events a day for far less money, but Kafka's consumer groups (every consumer reads the full stream at its own pace), keyed ordering and replay by rewinding offsets are what later rounds build on. See Message Queues vs Event Streams.
R1.9 Failure Modes
| Failure | What you'd see | How the design responds |
|---|---|---|
| MSK unreachable for an hour | Relay publish errors; outbox rows and the oldest row's age grow | Orders keep committing: the write path never talks to Kafka. The relay backs off (1 s, 2 s, 4 s, up to 30 s). At peak the outbox grows to about 180,000 rows (135 MB) and drains in about 14 s once the relay's next attempt succeeds (at most 30 s after MSK returns, because of the backoff). Nothing is lost. |
| Both relay tasks down | Oldest outbox row's age climbs | Same as a broker outage: lag grows, nothing is lost. ECS restarts the tasks. An alarm pages when the oldest row is older than 30 s for 5 minutes. |
| Aurora writer failover | Order writes and relay transactions fail for under a minute (typically) | Uncommitted transactions roll back, including relay claims. Committed outbox rows survive on Aurora's shared storage and are published after the failover. Rows published but not yet deleted are published again: duplicates the consumers absorb. |
| Fulfillment crashes mid-processing | A redelivery after restart | Its transaction rolled back, and the offset was never committed, so Kafka redelivers the event and fulfillment processes it as new. If it crashed after its commit but before the offset commit, the inbox rejects the redelivery. |
| The email service crashes around sending | Occasionally, one duplicate email | Sending an email is an external side effect that can't join a database transaction. Email writes an email_outbox row in the same transaction as its inbox row, and a sender sends it and marks it sent. A crash between "sent" and "marked" resends once. For email we accept that; for money we would not (Round 3, step 3.5). |
| A version gap at a consumer | Fulfillment receives v4 but its last version is 2 | It doesn't apply v4; it parks v4 (and anything later for that order) behind the gap and keeps the partition moving, so other orders aren't held up. When v3 arrives, it applies v3, then the parked events in order. A gap older than a few seconds alarms: in Round 1 it should never happen, so it means a bug upstream. |
R1.10 Pillar Check
| Pillar | What Round 1 covers |
|---|---|
| Reliability | The order write path doesn't depend on the broker; events wait in the outbox during outages; relays retry with backoff; consumers are idempotent, so every retry is safe REL 4 · REL 5 · REL 11 |
| Performance Efficiency | Polling chosen because "within seconds" allows it; the outbox stays a few kilobytes, so the relay's queries stay fast PERF 1 · PERF 3 |
| Security | Consumers never touch the order database; the relay's database role can only read and delete outbox rows; TLS to MSK and encryption at rest on Aurora and MSK SEC 3 · SEC 8 · SEC 9 |
| Cost Optimization | About $1,060 a month, most of it a Kafka cluster we share; the outbox pattern itself costs a few dollars COST 5 |
| Operational Excellence | Light this round: alarms on the oldest outbox row's age and on consumer version gaps OPS 8 |
| Sustainability | Skipped this round: two tiny relay tasks and a table of a few kilobytes. |
R1.11 Round 1 Rubric and Follow-Ups
What a strong mid-level (L5) answer shows
- Names the dual-write problem and explains why neither "publish first" nor "retry harder" fixes it.
- Writes the event into an outbox table in the same transaction as the state change.
- Designs a relay and says honestly that delivery is at-least-once.
- Makes consumers idempotent with an inbox updated in the same transaction as the effect, and commits offsets only after that transaction.
- Explains
SKIP LOCKED, and sees that parallel relays break per-entity order. - Keeps the outbox small: delete on success, not "mark processed forever".
Follow-up questions
-
"Why not use
occurred_atas the Kafka ordering key?" Answer: Kafka doesn't order by timestamps; it orders by position within a partition, and the key only decides the partition. The key must be the aggregate ID so that one order's events share a partition. Timestamps from different machines can also disagree by milliseconds or more. -
"We want to go from 16 to 64 buckets. Can we just change the modulus?" Answer: not while rows are pending. An order's pending rows could end up split between an old bucket and a new one, and two relays could publish them out of order. Stop the relays, let the outbox drain to empty (or wait until it is), change the bucket function and the
outbox_bucketstable together, and restart. The same care applies to Kafka partitions: adding partitions changes which partition a key maps to. -
"Could the relay use a high-water mark,
WHERE seq > last_published, instead of deleting rows?" Answer: it looks simpler, but it's a trap.seqvalues are handed out when rows are inserted, not when their transactions commit. A transaction that tookseq100 can commit after one that tookseq101; a relay that already moved past 101 never sees 100. Deleting claimed rows (or reading the commit-ordered log, Round 3) has no such gap.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Publish after commit, retry on failure" | The process can die between the two; the event is lost silently. |
| "Publish before commit" | A rolled-back change leaves a ghost event. |
| "Exactly-once delivery" | The relay can crash between publish and delete; delivery is at-least-once, and consumers make it once in effect. |
| "Commit the offset, then write to the database" | A crash in between skips the event forever. |
| "Order events by timestamp" | Clocks differ; order comes from the aggregate's version and the partition. |
| "Mark rows processed and keep them" | 274 GB a year of dead weight in the hottest database, plus a dead tuple per update. |
Round 2 · Senior · "50M Events a Day on a Ledger Under Contention"
~40 min · Senior SDE (L6) · 1 region, 3 AZs · 50M events/day: ~579/s average, 6,000/s peak · 30K balance reads/s · commit P99 < 15 ms · dispatch P99 < 100 ms · 99.99%
R2.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 2. If you're starting here, it's everything you need from Round 1.
Round 1 in 60 seconds. "An order service had to tell fulfillment and email about every paid order, about 1M events a day and 50 a second at peak. Committing and then publishing is a dual write, so we insert each event into an outbox table in the same transaction as the order change: both commit or neither does, so there are no lost events and no ghost events. A relay polls every second, claims a bucket of rows with
FOR UPDATE SKIP LOCKED, publishes them to Kafka keyed by order ID, and deletes them in the same transaction. A crash between publish and delete republishes, so delivery is at-least-once, and each consumer records the event ID in an inbox and applies its effect in one transaction, committing the Kafka offset only after that. Order per order holds because one bucket is published by one relay at a time, Kafka keeps a key in one partition, and consumers apply only the next aggregate version. Delete-on-success keeps the outbox at a few kilobytes instead of 274 GB a year. About $1,060 a month. The open costs: up to a second of delay, and nothing yet about contention or poison events."
Architecture v1, compact
Synthesizing vector architecture diagram...
Round 1 in one picture: the event is born in the same transaction as the change, and every hop after that may repeat but never loses or reorders an order's events.
Round 1 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.1 | Lost event after commit | Outbox row in the same transaction | A table to drain |
| 1.2 | Moving rows to Kafka | Relay: read, publish, delete | Up to 1 s of delay |
| 1.3 | Crash between publish and delete | At-least-once; inbox in the consumer's transaction | An inbox per consumer |
| 1.4 | Two relays, same rows | FOR UPDATE SKIP LOCKED | Order across relays |
| 1.5 | Out-of-order events | Buckets, key by aggregate, version check | Parallelism capped by buckets |
| 1.6 | Outbox growth | Delete on success; autovacuum per table | A dead tuple per event |
Open costs: a second of polling delay; one database with light writes; no plan for bad events, broker outages or churn at scale.
R2.1 The Scope Raise
Interviewer: "The same pattern now sits under our ledger service: account balances and transfers between accounts. It produces 50 million events a day, and at peak we see about ten times the average rate. Billing, analytics and audit consume every movement, and billing wants to see them within milliseconds, not after the next poll. Many transfers hit the same accounts at the same moment. Last month one malformed event stopped a consumer for an hour. The brokers get rebooted for patching. And we can't afford the table bloat we saw on another team's outbox."
We ask back, and say what each answer changes.
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| What is one "event" here? | One posted ledger entry. A transfer moves money between two accounts, so it posts two entries and publishes two events. | 50M events a day is 25M transfers a day; at peak 6,000 events/s is 3,000 transfers/s (R2.6). Events are keyed by account. |
| What does "within milliseconds" mean? | Billing's P99 from commit to "published" under 100 ms. | A 1 s poll is out; we need a wake-up on commit (step 2.1). |
| How hot do accounts get? | Most are quiet. A few hundred business accounts take tens of transfers a second each; the biggest over 100 a second at peak. | Per-account serialization must be correct and cheap (step 2.2). Accounts far above that need the wallet loop's striping. |
| Can a consumer skip a bad event? | Audit must never silently skip one. Billing can delay one account, not all of them. | A dead-letter queue that blocks only the affected account, with an alarm and a replay path (step 2.3). |
| How long can the broker be unavailable? | Patching reboots one broker at a time; the worst full outage we've had was 10 minutes. | The outbox must absorb 10 minutes at peak without hurting transfers (step 2.4). |
| How many balance reads? | About 5 reads per event at peak: 30,000 a second. | Reads go to Aurora readers, not the writer (R2.6). |
| Product asked for five nines on transfers. | They'd like it. | 99.999% allows 26 seconds of downtime a month. One Aurora writer failover typically takes 30–60 seconds, and Aurora's own SLA for Multi-AZ clusters is 99.99%. We commit to 99.99% (about 4.4 minutes a month) and say why. |
Scope change
| Round 1 | Round 2 | |
|---|---|---|
| Domain | Orders | A double-entry ledger |
| Events | 1M/day, 50/s peak | 50M/day: ~579/s average, 6,000/s peak |
| Writes | Light, rarely on the same row | 3,000 transfers/s at peak; hot accounts |
| Reads | Not our concern | 30,000 balance reads/s |
| Dispatch delay | Within 5 s | P99 < 100 ms from commit to published |
| Commit latency | Not stated | P99 < 15 ms |
| Bad events | Not considered | One bad event must not stop the stream |
| Broker outages | "Lag grows" | 10 minutes at peak, absorbed |
| Availability | 99.9% (43.8 min/month) | 99.99% (4.4 min/month) |
R2.2 What Breaks in the Round 1 Design
| Round 1 choice | What breaks at the new scope |
|---|---|
| Relay polls every second | Average delay 0.5 s, worst about 1 s. Billing wants P99 under 100 ms. Polling every 10 ms would work but runs 100 queries a second per relay, forever. |
| Order row lock, taken implicitly | A transfer touches two rows. Taking them in the wrong order deadlocks, and "check the balance, then write" overdraws. |
| Consumer retries a failing event forever | One malformed event blocks its whole Kafka partition, which carries thousands of other accounts. |
| "The outbox buffers outages" | True, but at 6,000 events a second a 10-minute outage is 3.6M rows. We have to size it and drain it without starving transfers. |
| Delete on success with default autovacuum | One dead tuple per event: 50M a day, 6,000 a second at peak. If anything holds vacuum back, the relay's index scans slow down. |
| Events are the only record consumers have | For money, a consumer's copy that drifts from the ledger must be caught, not assumed away. |
The order we fix it in: dispatch latency (2.1), money correctness under contention (2.2), bad events (2.3), outages (2.4), churn (2.5), then audit (2.6).
R2.3 New Requirements and API Additions
Make a transfer
httpPOST /v1/transfers HTTP/1.1 Host: ledger.internal.example Authorization: Bearer <service token for the checkout service> Idempotency-Key: 5d0c7e2a-1b4f-4e9a-8c3d-7f6a2b1e9d05 Content-Type: application/json { "from_account_id": 1001, "to_account_id": 2002, "amount_minor": 2500, "currency": "USD", "reference": "inv_7731" }
httpHTTP/1.1 201 Created Content-Type: application/json { "transfer_id": 88210045, "status": "POSTED", "entries": [ { "account_id": 1001, "amount_minor": -2500, "balance_after_minor": 47500, "entry_seq": 57 }, { "account_id": 2002, "amount_minor": 2500, "balance_after_minor": 12500, "entry_seq": 12 } ] }
| Status | When |
|---|---|
201 Created | Posted, or a retry of one that was (same body) |
422 Unprocessable Entity | INSUFFICIENT_FUNDS, ACCOUNT_FROZEN, or IDEMPOTENCY_KEY_REUSED. Rejections are stored with the key, so a retry gets the same rejection |
503 Service Unavailable | Mid-failover; retry with the same key |
The event, one per posted entry. The aggregate is now the account, and aggregate_version is the account's entry_seq: 1, 2, 3 ... with no gaps, bumped only when an entry is posted.
json{ "event_id": "0192a7b1-4e2c-7d11-8a9b-3c5d6e7f8a90", "event_type": "ledger.entry_posted", "schema_version": 1, "aggregate_type": "account", "aggregate_id": "1001", "aggregate_version": 57, "occurred_at": "2026-09-27T18:04:11.093Z", "producer": "ledger-service", "correlation_id": "req_a81f", "payload": { "transfer_id": 88210045, "amount_minor": -2500, "currency": "USD", "balance_after_minor": 47500, "counter_account_id": 2002 } }
Dead-letter admin API (for operators; covered in step 2.3)
httpGET /v1/admin/dead-letters?consumer=billing&status=OPEN HTTP/1.1
httpHTTP/1.1 200 OK Content-Type: application/json { "items": [ { "dead_letter_id": "dl_311", "consumer": "billing", "account_id": "1001", "aggregate_version": 57, "reason": "SCHEMA_VALIDATION: missing currency", "parked_behind_it": 3, "first_seen": "2026-09-27T18:04:12Z" } ] }
httpPOST /v1/admin/dead-letters/dl_311/replay HTTP/1.1 Content-Type: application/json { "mode": "FIXED_CONSUMER" }
replay sends the dead-lettered event, then the events parked behind it for that account, back through the consumer in version order.
Per-account serialization rules (the contract every code path follows):
- Every change to an account's balance happens inside a transaction that holds that account's row lock until commit.
- When a transaction touches several accounts, it locks them in ascending
account_idorder. - No rule that spans rows may be checked without first locking a row that stands for it (step 2.2).
- Nothing slow runs while a lock is held: no calls to other services, no broker publishes.
R2.4 Design Evolution: Fast, Safe and Bounded
Step 2.1: "Consumers Want Events in Milliseconds, Not After the Next Poll"
The problem: billing wants P99 under 100 ms from commit to "published". The relay polls every second. What would you do? There are at least four ways to get events out of the database faster. Compare them, then pick one for 6,000 events a second.
The relay loop with the signal:
textafter each transfer COMMIT: signal.wake() -- non-blocking; merges with a pending wake-up each relay loop (2 per task, own pool of 2 connections): wait for a wake-up or 1 s, whichever comes first wait 5 ms -- linger: batch more commits together repeat: BEGIN b = claim one bucket that has rows (FOR UPDATE SKIP LOCKED) if none: COMMIT; stop repeating rows = SELECT ... FROM outbox WHERE bucket = b ORDER BY seq LIMIT 500 publish each row to ledger.entries, key = account_id; wait for all acks (2 s timeout: ROLLBACK, back off 100 ms doubling to 5 s) DELETE FROM outbox WHERE seq = ANY(the published seqs) COMMIT
Primitive: Change Data Capture and the Outbox Pattern
Step 2.2: "Two Transfers on One Account Overdrew It"
The problem: account 1001 holds 30.00. Two transfers of 25.00 from it arrive together. Both read "30.00, enough", and both post. The balance is −20.00. What would you do? There are five common ways to serialize work on an account. Which one, and what exactly does it guarantee?
When a rule spans rows: write skew. Suppose the ledger adds "at most 10,000.00 out of an account per day", checked by summing today's outgoing entries. Two transfers each read "9,990.00 sent today", each pass, and each insert a different new entry row. Under snapshot isolation (PostgreSQL's Repeatable Read), nothing conflicts, because snapshot isolation only catches two writers of the same row: this anomaly is write skew. Our design is safe anyway: both transfers first update the account's own row, so the second waits for the first, and its limit check runs after the first committed (it reads the entries with a fresh statement under Read Committed). Locking one row that stands for the whole rule is called materializing the conflict.
Why not run everything at Serializable? PostgreSQL's Serializable level (serializable snapshot isolation) does catch write skew, but by aborting transactions it can't prove safe. On busy accounts the abort-and-retry rate climbs exactly when load is highest, and it tracks extra read locks in memory. We choose Read Committed plus explicit locks on the rows that carry each rule, and we write down which rule each lock protects.
Advisory lock details, because they come up:
- The key is a 64-bit integer (
pg_advisory_xact_lock(bigint)), or two 32-bit integers. If account IDs are alreadyBIGINT, use the ID itself: no hashing, no collisions. For UUIDs, use a 64-bit hash (PostgreSQL'shashtextextended(text, seed)returns abigint), never a 32-bit one likehashtext. - Why 32 bits hurts: by the birthday bound, two of keys collide with probability about , which reaches 50% at about accounts. With 20M accounts we'd expect colliding pairs. With 64 bits and 100M accounts: , essentially none.
- A collision is a latency bug, not a correctness bug: two unrelated accounts queue behind each other. Sort by the lock key (not the account ID) when taking several, and a collision can't deadlock; a transaction that already holds a key gets it again at once.
- Memory: advisory locks live in PostgreSQL's shared lock table, which holds about
max_locks_per_transaction × (max_connections + max_prepared_transactions)entries. With the default of 64 and, say, 2,000 connections, that's 128,000 entries; at roughly 270 bytes each (an estimate) about 35 MB. They write no WAL and dirty no table pages.
Primitive: Database Isolation Levels, ACID & Concurrency Anomalies · Drill: ACID isolation write skew anomaly (both of its questions are answered just above)
Step 2.3: "One Bad Event Blocks Everything Behind It"
The problem: a deploy of the ledger service published 40 events with currency missing. Billing's consumer throws on the first one, retries it forever, and its whole Kafka partition stops. That partition carries about 1/64 of all accounts.
What would you do?
Synthesizing vector architecture diagram...
One account waits; the partition keeps moving.
Step 2.4: "The Broker Is Down for 10 Minutes"
The problem: the whole MSK cluster is unreachable for 10 minutes during a bad network change, at peak. What would you do?
Step 2.5: "Vacuum Can't Keep Up With Delete Churn"
The problem: every event is an insert and a delete: 6,000 dead tuples a second at peak, 50M a day. One afternoon, the relay's claim-and-read queries go from 2 ms to 400 ms, and dispatch lag climbs. The outbox has only 3,000 live rows. What would you do?
The schema has to allow the drop job. In PostgreSQL a primary key on a partitioned table must include the partition key, so the outbox key is (seq, created_at), not seq alone:
sqlCREATE TABLE outbox ( seq BIGINT GENERATED ALWAYS AS IDENTITY, -- identity on a partitioned table needs PostgreSQL 17+; on older versions use a bigserial/sequence default on the parent created_at TIMESTAMPTZ NOT NULL DEFAULT now(), bucket SMALLINT NOT NULL CHECK (bucket BETWEEN 0 AND 63), -- account_id mod 64 event_id UUID NOT NULL, aggregate_id BIGINT NOT NULL, -- account_id aggregate_version BIGINT NOT NULL, -- the account's entry_seq event_type TEXT NOT NULL, schema_version INT NOT NULL, payload JSONB NOT NULL, PRIMARY KEY (seq, created_at) ) PARTITION BY RANGE (created_at); CREATE INDEX outbox_by_bucket ON outbox (bucket, seq); -- created on every partition -- the job, for each new hourly partition: CREATE TABLE outbox_2026092718 PARTITION OF outbox FOR VALUES FROM ('2026-09-27 18:00+00') TO ('2026-09-27 19:00+00'); ALTER TABLE outbox_2026092718 SET (autovacuum_vacuum_scale_factor = 0, autovacuum_vacuum_threshold = 10000, autovacuum_vacuum_cost_delay = 0);
Autovacuum settings are set on each partition, because vacuum works on partitions, not on the parent. The Round 1 uniqueness rule on (aggregate, version) is dropped here for the same reason (it can't be enforced across partitions); the ledger's own UNIQUE (account_id, entry_seq) on ledger_entries does that job.
Step 2.6: "The Ledger Must Be Auditable End to End"
The problem: auditors ask: "Billing says account 1001 had 475.00 on Tuesday. How do you know billing's copy matches the ledger?" What would you do?
Round 2 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Polling is too slow | In-process coalescing signal + bucket claims + 1 s sweep | A sweep for missed pokes; claim-and-delete per batch |
| 2.2 | Overdrafts under concurrency | Conditional UPDATE row locks in ascending account order; locks that materialize cross-row rules | ~125–200 transfers/s per account ceiling |
| 2.3 | One bad event blocks a partition | Schema check at insert; consumer DLQ that blocks only one account; replay API | A stale account until resolved; DLQ on-call |
| 2.4 | 10-minute broker outage | Outbox absorbs 3.6M rows; capped drain in ~4.3 min | A lag spike |
| 2.5 | Delete churn and vacuum | Short transactions, batched deletes, per-partition autovacuum, hourly partitions dropped when empty | A partition job and alarms |
| 2.6 | Proving consumers match | Ledger as truth; per-account entry_seq reconciliation with a grace window | A daily job and an owner |
R2.5 Architecture v2
Synthesizing vector architecture diagram...
Writes and relay work go to the writer; reads and exports go to readers. The relay lives inside the ledger service so a commit can wake it in microseconds.
The pieces that changed from Round 1:
- Relay inside the service: two loops per task, 12 in total, each with its own pool of two connections (24 relay connections), separate from the request pool.
- 64 buckets and 64 partitions: enough parallelism for 6,000 events/s, and headroom for the drain.
- Aurora cluster: a
db.r7g.4xlargewriter in AZ a, adb.r7g.4xlargereader in AZ b as the failover target, and twodb.r7g.2xlargereaders in AZ c (sizing in R2.6). - MSK: 3 ×
kafka.m7g.xlarge, one per AZ, replication factor 3,min.insync.replicas=2, 7-day retention.
Ledger schema (the outbox is in step 2.5)
sqlCREATE TABLE accounts ( account_id BIGINT PRIMARY KEY, kind TEXT NOT NULL CHECK (kind IN ('USER', 'SYSTEM')), currency CHAR(3) NOT NULL, status TEXT NOT NULL DEFAULT 'ACTIVE' CHECK (status IN ('ACTIVE', 'FROZEN', 'CLOSED')), balance_minor BIGINT NOT NULL DEFAULT 0, entry_seq BIGINT NOT NULL DEFAULT 0, -- +1 only when an entry is posted status_version BIGINT NOT NULL DEFAULT 0, -- +1 on freezes and other status changes CONSTRAINT user_not_negative CHECK (kind = 'SYSTEM' OR balance_minor >= 0) ); CREATE TABLE transfers ( transfer_id BIGINT PRIMARY KEY, from_account BIGINT NOT NULL REFERENCES accounts, to_account BIGINT NOT NULL REFERENCES accounts, amount_minor BIGINT NOT NULL CHECK (amount_minor > 0), created_at TIMESTAMPTZ NOT NULL DEFAULT now(), CHECK (from_account <> to_account) ); CREATE TABLE ledger_entries ( transfer_id BIGINT NOT NULL REFERENCES transfers, account_id BIGINT NOT NULL REFERENCES accounts, entry_seq BIGINT NOT NULL, amount_minor BIGINT NOT NULL CHECK (amount_minor <> 0), -- signed; a transfer's two rows sum to zero balance_after BIGINT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (account_id, entry_seq) -- no gaps, no repeats per account ); CREATE TABLE idempotency_keys ( caller_id TEXT NOT NULL, idem_key TEXT NOT NULL, request_hash BYTEA NOT NULL, response JSONB, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (caller_id, idem_key) );
Idempotency keys are kept 30 days and pruned in batches: like the inbox, this table can't be partitioned by day without losing key uniqueness across days. Rejections and successes are both stored with their key. And because the event rows commit in the same transaction as the transfer, a retry that finds its key already stored doesn't need to publish anything again: the outbox already guarantees the original event goes out.
Trace 1: a transfer at peak
Synthesizing vector architecture diagram...
The client's answer never waits for the broker; the event follows within tens of milliseconds.
If either UPDATE returns zero rows (not enough money, or a frozen account), the handler rolls back to the savepoint, stores the 422 response on the idempotency row and commits. No entries and no events are written for a rejection, so the ledger balances on that branch too.
Trace 2: a poison event. Billing receives 1001 v57 with currency missing (a producer bug that slipped past validation). It dead-letters v57 and blocks account 1001 in one transaction, commits the offset and continues with other accounts. v58 for account 1001 arrives and is parked. The audit page fires; after the fix, the replay API sends v57 and v58 back in order. (Diagram in step 2.3.)
Trace 3: a broker outage.
Synthesizing vector architecture diagram...
The write path never stops; the only cost is lag, gone about 14.4 minutes after the outage began (10 min + up to 5 s backoff + 4.3 min drain).
R2.6 Numbers and Cost
Traffic
| Quantity | Math | Value |
|---|---|---|
| Events per day | given | 50,000,000 |
| Average event rate | 50,000,000 ÷ 86,400 s | ≈ 578.7/s |
| Peak event rate | 578.7 × 10 = 5,787 | plan for 6,000/s |
| Transfers | 2 events per transfer | 25M/day; 3,000/s at peak |
| Balance reads at peak | 6,000 × 5 | 30,000/s |
The outbox
| Case | Math | Size |
|---|---|---|
| Keep every row (mark processed) | 50M × 750 B = 37.5 GB/day; × 365 | 37.5 GB/day ≈ 13.7 TB/year |
| Delete on success, 2 s polling (for comparison) | 6,000/s × 2 s = 12,000 rows × 750 B | ≈ 9 MB |
| Delete on success, our signal (P50 ~15 ms, P99 < 100 ms) | 6,000/s × 0.1 s = 600 rows × 750 B | ≈ 450 KB at P99 |
| A missed poke (1 s sweep), worst case | 6,000 rows × 750 B | ≈ 4.5 MB |
| 10-minute broker outage at peak | 3.6M rows × 750 B | ≈ 2.7 GB, drained in ~4.3 min |
Relay load (rough). 12 loops, each waking about every 20 ms at peak and draining roughly one bucket batch of ~10 rows per wake-up, gives about 600 relay transactions a second at peak, each with four statements. We check that number in the load test; the linger is the knob that trades a few milliseconds of delay for fewer, bigger batches.
Writer sizing. We plan on about 500 short transactions per second per vCPU (a planning assumption, the same as the wallet loop's 1,000/s on a 2-vCPU instance; confirmed by a load test that runs our real transfer). At peak: 3,000 transfers + ~600 relay transactions ≈ 3,600/s ≈ 7.2 vCPUs. A db.r7g.4xlarge (16 vCPUs) runs at about 45% at peak, which leaves room for the drain after an outage.
Reader sizing. We plan on about 2,500 indexed single-row reads per second per vCPU (an assumption to load-test). 30,000 ÷ 2,500 = 12 vCPUs. Readers: one 4xlarge (16, AZ b) + two 2xlarge (8 each, both in AZ c) = 32 vCPUs, about 38% busy. We place them so that any single AZ loss leaves 16 reader vCPUs: lose AZ a and the AZ b reader is promoted to writer, leaving the two in AZ c; lose AZ b and the two in AZ c remain; lose AZ c and the AZ b reader remains. 16 vCPUs at peak is about 75% busy, still serving. (Two 2xlarge readers in different AZs would leave only 8 vCPUs after losing AZ a.)
Commit latency budget (P50, dependent steps add; planning figures)
| Step | Time |
|---|---|
| Insert idempotency key + savepoint | ~1 ms |
Debit UPDATE, then credit UPDATE | ~2 ms |
| Insert transfer, 2 entries, 2 outbox rows (3 statements) | ~1.5 ms |
| Store the response | ~0.5 ms |
COMMIT (Aurora writes the log to storage in 3 AZs and waits for a quorum) | ~3 ms |
| Total | ≈ 8 ms, leaving ~7 ms under the 15 ms P99 for lock waits and tail |
WAL volume (rough). A transfer writes about 3 KB of WAL records (two account updates, a transfer row, two entries, two outbox rows with ~500-byte payloads, a key row, index entries, the later outbox deletes and a commit record), before full-page images after checkpoints. Call it ~2 KB per event: 50M × 2 KB ≈ 100 GB/day, about 12 MB/s at peak. Nothing in this round retains WAL (no replication slots), but Round 3's CDC will read it, so we keep the number.
MSK sizing. At ~1 KB per record: 6 MB/s in at peak (0.58 MB/s average), 18 MB/s written across brokers with replication factor 3, and 18 MB/s read by three consumer groups. The drain after an outage pushes 20 MB/s in. Three kafka.m7g.xlarge brokers carry this with wide margin (we chose xlarge over large for drain headroom). Storage: 50M × 1 KB × 7 days × 3 replicas = 1,050 GB; we provision 500 GB per broker, 1,500 GB in total.
Monthly cost (us-east-1 on-demand list prices, 730 hours)
| Item | Math | Monthly |
|---|---|---|
| Aurora instances, I/O-Optimized | (2 × $2.212 + 2 × $1.106)/h × 1.3 × 730 h | ≈ $6,298 |
| Aurora storage, I/O-Optimized | ~7 TB after a year (25M transfers × 750 B × 365 ≈ 6.8 TB) × $0.225/GB-month | ≈ $1,575 |
MSK, 3 × kafka.m7g.xlarge | 3 × $0.408/h × 730 h | ≈ $894 |
| MSK storage | 1,500 GB × $0.10 | $150 |
| Fargate, ledger service, 6 tasks × 1 vCPU, 2 GB, Arm | 6 × (1 × $0.03238 + 2 × $0.00356)/h × 730 h | ≈ $173 |
| Fargate, DLQ tooling and reconciliation runner | ≈ $30 | |
| Application Load Balancer (rough: capacity units depend on bytes) | ≈ $150 | |
| CloudWatch | ≈ $100 | |
| S3 exports + Athena reconciliation | ≈ $70 | |
| Total | ≈ $9,440/month |
That's about $6.2 per million events (9,440 ÷ 1,520M events a month).
Why I/O-Optimized. On Aurora Standard, the same instances cost $4,844, storage $700, and I/O would be about $6,080 (assume ~30 billed I/Os per transfer and ~5 per event for the relay: (25M × 30 + 50M × 5) × 30.4 days ≈ 30.4B I/Os × $0.20 per million), about $11,624 in all, with I/O over half the bill. AWS suggests I/O-Optimized once I/O is above about a quarter of Aurora spend; here it saves about $3,750 a month.
Ledger history older than 12 months is exported to S3 and removed from Aurora, as in the wallet loop; the figures above assume that.
R2.7 Trade-Offs
The dispatch approaches (figures are rough and depend on hardware and payloads)
| Polling (1 s) | In-process signal + sweep (our choice) | LISTEN/NOTIFY + SKIP LOCKED | Log-based CDC | |
|---|---|---|---|---|
| Latency | ~0.5 s average, ≤ ~1 s | ~10–20 ms typical; ≤ 1 s if a poke is lost | Milliseconds while connected; the sweep interval if disconnected | Tens to hundreds of ms |
| Throughput limit | Claim-and-delete rate on the writer (thousands/s) | Same | Commits that NOTIFY are serialized by a global lock (one report: ~2.9K writes/s) | Tens of thousands/s per database (one task per PostgreSQL connector) |
| Table bloat | One dead tuple per event (delete) | Same | Same | None if the outbox is insert-only and partitions are dropped |
| WAL risk | None | None | None; notifications use their own queue (8 GB by default; if full, NOTIFY fails at commit) | A stalled slot retains WAL |
| Coordination | SKIP LOCKED bucket claims | Same | Same, plus a listener connection | One connector per database; Kafka partitions |
| Payload limit | None | None | Payload < 8,000 bytes: send IDs only | Kafka message size (1 MB default) |
| Complexity | Low | Low | Medium: session connections, no transaction pooling | High: connector, slots, snapshots |
| Right when | Seconds are fine | Milliseconds at thousands/s, one service | You can't touch the writing code, and commits are few | Very high rates or many services (Round 3) |
The serialization approaches
Row locks, conditional UPDATE (our choice) | Advisory locks | Optimistic (version) | Distributed lock | Single writer | |
|---|---|---|---|---|---|
| Overhead | The update we do anyway | In memory; no WAL | None until conflict | A network round trip per lock | Tiny per operation, in memory |
| Deadlocks | Possible; lock in ascending order | Possible; lock in ascending key order | None (retries instead) | Possible across keys | None |
| Hot account | Queues in order (~125–200/s per row) | Queues in order | Retry storms | Queues, plus network | Best: no locks at all |
| Forgotten by a code path? | Impossible: no update without it | Possible: it's voluntary | Possible if a path skips the version | Possible | Impossible: one writer |
| Failure safety | Released on commit, rollback or crash | Released on commit, rollback or crash (xact form) | Nothing held | A paused holder can outlive its lease | You build failover |
R2.8 Failure Modes
| Failure | What you'd see | How the design responds |
|---|---|---|
| Dual-write drift (a team bypasses the outbox and publishes after commit) | Reconciliation breaks: entries in the ledger with no matching event, or events for transfers that rolled back | The proof is the timeline in step 1.1: commit succeeds, publish times out, the process dies. The fix is the outbox, enforced by review and by the ledger role being the only producer allowed on ledger.entries. |
| A poison event | A dead letter; one account blocked in one consumer | Step 2.3: page, fix, replay in order. Other accounts are unaffected. |
Advisory-lock hash collisions (a service that locks by hashtext) | Lock waits on accounts that are quiet | Two unrelated accounts share a 32-bit key. Switch to the 64-bit form or the account's own BIGINT ID. |
| Replayed duplicates (a consumer rewinds its offsets by an hour) | A burst of events the consumer already applied | The per-account entry_seq check rejects every one; nothing is applied twice. |
| Relay lag during an Aurora failover | Transfers fail for typically under a minute; the relay's transactions abort | Callers retry with their idempotency keys. Committed outbox rows survive on Aurora's shared storage. The relays reconnect; pokes lost in the failover are covered by the 1 s sweep. Rows published but not deleted before the crash go out again; consumers skip them. |
| Losing an AZ | One third of tasks, one MSK broker, maybe the writer | Remaining tasks keep serving; MSK keeps accepting acks=all writes with two of three replicas; if the writer was in that AZ, Aurora promotes the 4xlarge reader in AZ b. Whichever AZ is lost, 16 reader vCPUs remain (R2.6), about 75% busy at peak. |
R2.9 Production Gotchas
| Gotcha | Symptom | Cause | Fix |
|---|---|---|---|
| Network calls inside a database transaction | Connection pool exhausted, lock waits climb when a downstream is slow | A fraud check or push notification between BEGIN and COMMIT on the request path | Nothing external inside request transactions. The relay's publish is the one bounded exception: own pool, only outbox rows locked, hard timeout. |
| A mutable outbox without vacuum tuning | Relay queries slow down; table and index grow | status = 'SENT' updates, or default autovacuum settings, or a long transaction pinning cleanup | Delete on success, per-partition autovacuum, hourly partition drops, transaction timeouts |
| Committing broker offsets before the database transaction | A consumer silently misses events after a crash | Auto-commit on, or commit offset then write | enable.auto.commit=false; commit the offset only after the database commit |
| Ordering broken by parallel relays | Consumers see version gaps and alarms | Relays claim rows, not buckets; or a new code path inserts the outbox row before locking the account | Claim whole buckets; insert outbox rows after the account update |
| A sequence high-water mark as the relay's cursor | Rare events never published | Sequence numbers are assigned at insert, not at commit, so a lower number can commit later | Claim and delete, or read the commit-ordered log |
| No future partitions | Transfers suddenly fail with "no partition found for row" | The partition-creation job stopped | Create 48 hours ahead; alarm below 24 |
R2.10 Pillar Check
| Pillar | What Round 2 adds |
|---|---|
| Reliability | The write path survives a 10-minute broker outage by design; bad events are contained to one account per consumer; every hop retries safely; a writer failover loses no committed event REL 4 · REL 5 · REL 10 · REL 11 |
| Performance Efficiency | Dispatch approach chosen by comparison, not habit; conditional updates as locks we'd pay for anyway; writer sized from the peak with drain headroom PERF 1 · PERF 3 |
| Security | Only the ledger role may produce to ledger.entries (MSK IAM access control); the relay role can read and delete outbox rows only; consumers can't write the ledger; TLS in transit and KMS encryption at rest SEC 3 · SEC 8 · SEC 9 |
| Cost Optimization | About $9,440 a month, $6.2 per million events; I/O-Optimized chosen from the I/O share; old ledger history moved to S3 COST 5 · COST 6 |
| Operational Excellence | Alarms with first actions (below); DLQ ownership and replay tooling; a partition job with its own alarm OPS 8 · OPS 10 |
| Sustainability | Light this round: Graviton instances; the outbox kept at kilobytes instead of 13.7 TB a year; history in S3, not on database storage SUS 3 · SUS 4 · SUS 5 |
Alarms and first actions
| Signal | Alarm | First action |
|---|---|---|
| Oldest outbox row's age | > 5 s for 2 min (warn); > 30 s (page) | Check MSK health and relay errors |
| Outbox row count | > 50,000 | Same; expect it during a broker outage |
| Oldest open transaction on the writer | > 10 min | Find and end it; it's blocking vacuum |
| Dead tuples in the current outbox partition | > 1M | Check the oldest transaction and autovacuum activity |
| Future outbox partitions | < 24 | Run the partition job; transfers will fail when they run out |
| Dead letters | audit: any (page); billing, analytics: > 10/min | Find the producer or consumer bug; fix; replay |
| Consumer lag per group | > 60 s for 5 min | Scale the consumer; check its database |
| Reconciliation breaks | any | Same-day investigation; repair the consumer from the ledger |
R2.11 Round 2 Rubric and Follow-Ups
What a senior (L6) answer adds over L5
- Compares polling, a commit-time signal,
LISTEN/NOTIFYand CDC with their real limits (payload size, lost notifications, poolers, the commit lock, slot WAL), and picks one with numbers. - Serializes money with locks it can explain: which rows, what order, what guarantee, and where write skew hides.
- Knows advisory locks are voluntary and 64-bit, and that collisions cost latency, not correctness.
- Contains poison events to one aggregate and keeps order for it, instead of stopping or skipping.
- Sizes an outage and its drain, and keeps transfers up while the broker is down.
- Knows why delete-heavy tables bloat, what pins vacuum, and designs partitions the schema can actually drop.
- Treats events as announcements and proves consumer copies against the ledger.
Follow-up questions
-
"Why not publish from the request thread right after commit, and keep the outbox only as a fallback?" Answer: it's a legitimate optimization, as long as the outbox row is still written in the transaction and deleted only by whoever publishes it. But now two publishers can race for the same row: the request thread publishes v57 while a relay also claims it, and ordering per account depends on who wins. Our poke gets within a few milliseconds of that with one publisher per bucket, so we don't take on the race.
-
"A consumer wants all events for a transfer together, not per account. What do you do?" Answer: add a
ledger.transfer_postedevent keyed bytransfer_id, written as a third outbox row in the same transaction. Keys decide ordering, so we pick the key per consumer need. Don't re-key the existing topic: every account consumer depends on its per-account order. -
"Could we skip the inbox entirely now?" Answer: for per-account consumers, yes: the stored
entry_seqrejects duplicates and detects gaps, and it's one row per account instead of one row per event (20M accounts × ~50 B ≈ 1 GB, versus 50M events/day × 8 days × ~120 B ≈ 48 GB per consumer). Consumers without per-aggregate state still need an inbox.
Interview gotchas from this round
| Gotcha | Why it's wrong |
|---|---|
| "Poll faster" | Latency down, idle load up, forever. Wake the relay on commit. |
"LISTEN/NOTIFY is a free push" | Lost while disconnected, 8,000-byte payloads, no transaction pooling, and a commit-time global lock. |
| "Advisory locks can't deadlock" | They can; take them in sorted key order. And a 32-bit hash collides at tens of thousands of keys. |
| "Skip the poison event" | Audit silently loses money movements. Dead-letter it, block that account, alarm, replay. |
| "Delete on success means zero dead tuples" | Every delete leaves a dead tuple; small tables and partition drops keep that cheap. |
| "Five nines" | 26 seconds a month; one Aurora failover spends it. |
Round 3 · Architect · "An Event Backbone for the Whole Company"
~45 min · Principal (L7) · 3 regions · 40 services · 2B events/day, 70K/s peak · contracts that never break consumers · replay a year of events · 99.99% per region
R3.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 3. If you're starting here, it's everything you need from Rounds 1 and 2.
Round 2 in 60 seconds. "Our ledger publishes 50M events a day, 6,000 a second at peak, one event per posted entry, keyed by account, with the account's gap-free
entry_seqas the version. Each transfer is one transaction: an idempotency key, conditional updates that lock both accounts in ascending order, two entries, two outbox rows. The relay lives inside the service: a coalescing poke after each commit wakes it, it claims whole buckets withSKIP LOCKED, publishes in order, deletes, and a 1-second sweep catches missed pokes, so dispatch is about 15 ms. We rejectedLISTEN/NOTIFYfor its commit-time global lock and lost notifications, and deferred CDC. Consumers checkentry_seq, dead-letter a bad event and block only that account, and a daily reconciliation compares each consumer's position against the ledger. A 10-minute broker outage becomes 3.6M outbox rows, drained in about 4 minutes. Hourly outbox partitions, dropped when empty, and short transactions keep vacuum ahead of 6,000 deletes a second. About $9,440 a month, 99.99%. Open costs: every team builds its own relay, nothing stops a producer from breaking a consumer, history lasts only as long as Kafka keeps it, and it's one region."
Architecture v2, compact
Synthesizing vector architecture diagram...
Round 2 in one picture: one well-run outbox for one service.
Steps so far
| Step | Problem | Component |
|---|---|---|
| 1.1–1.2 | Dual write | Outbox row in the same transaction; relay |
| 1.3–1.5 | Duplicates, two relays, order | Idempotent consumers; SKIP LOCKED bucket claims; key by aggregate; version check |
| 1.6 | Growth | Delete on success |
| 2.1 | Polling delay | In-process signal + sweep |
| 2.2 | Overdrafts | Conditional-update row locks in ascending order |
| 2.3 | Poison events | DLQ that blocks one account |
| 2.4–2.5 | Outages, vacuum | Outbox absorbs; hourly partitions; short transactions |
| 2.6 | Audit | Per-account reconciliation against the ledger |
Open costs: one service's relay, no contracts, 7 days of history, one region.
R3.1 The Scope Raise
Interviewer: "The ledger team's pattern worked, so now everyone wants it. We have 40 services, each with its own database, publishing events that other teams depend on. Last quarter a producer renamed a field and three consumers broke in production. Together we're heading for about 2 billion events a day, and the teams that copied your relay are fighting vacuum and relay lag on their busiest databases."
Interviewer: "The audit team wants to rebuild any read model from the event history, going back a year. We run in three regions, and some events are needed outside the region that produced them. And every team asks the same question: what is exactly-once and what isn't?"
We ask back, and say what each answer changes.
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| Are all 40 on PostgreSQL? | Yes: 30 on Aurora PostgreSQL, 10 on RDS for PostgreSQL. | One CDC approach fits all (step 3.2). And the WAL risk differs: RDS instances have a fixed disk that a stalled slot can fill; Aurora keeps WAL on its cluster volume. |
| Who owns an event's shape? | The producing team. Consumers are other teams. | Contracts need a registry and checks the producer can't skip (step 3.1). |
| How far back must we replay? | One year for read models. The ledger keeps its own 7-year history. | Kafka keeps 7 days; a year lives in an S3 archive (step 3.3). |
| Which regions, and how much crosses? | us-east-1, eu-west-1, ap-southeast-1, about 50/30/20 of traffic. Every customer and account has one home region. About 10% of events are needed in the other two regions. | Regional clusters, with only the global topics replicated (step 3.4). |
| What's the busiest single database? | The ledger in us-east-1: 100M events a day, about 12,000 a second at peak. | CDC per database must handle 12K/s (step 3.2). |
| What do teams mean by exactly-once? | Some think Kafka gives it end to end. | An honest, per-hop contract (step 3.5). |
| Is there a team to own this? | A platform team of six engineers. | A shared platform instead of 40 relays (step 3.6). |
Scope change
| Round 2 | Round 3 | |
|---|---|---|
| Producers | 1 service | 40 services, 120 databases (each service in each region) |
| Events | 50M/day, 6,000/s peak | 2B/day: ~23,150/s average, 70,000/s peak |
| Event shapes | One team decides | Contracts across teams; changes must not break consumers |
| History | 7 days in Kafka | 1 year, replayable |
| Regions | 1 | 3; 10% of events cross regions |
| Delivery promise | At-least-once, idempotent consumers | A written per-hop contract |
R3.2 What Breaks in the Round 2 Design
| Round 2 choice | What breaks at the new scope |
|---|---|
| Event shapes live in each producer's code | Nothing stops a rename or type change from reaching consumers. |
| A relay inside each service | 40 copies of claim-and-delete logic, each with its own bugs. On the busiest database, 12,000 deletes a second of churn plus relay transactions on the writer. |
| Dispatch by polling the outbox | At this scale we want the log-based approach, which brings replication slots, and a stalled slot on an RDS instance can fill its disk. |
| Kafka as the only history | 7 days. A read model with a bug found after 8 days can't be rebuilt from events. |
| One region | Events produced in Europe are invisible to consumers in Virginia. |
| "At-least-once" said per team | Some teams promise exactly-once to their own consumers and are wrong about it. |
The order we fix it in: contracts first (they break things today), then CDC (it removes 40 relays), history and replay, regions, the delivery contract, and the platform that owns all of it.
R3.3 New Requirements and API Additions
The envelope, version 2. Two fields make regions and failovers safe, and one ties the payload to its registered schema:
json{ "event_id": "0192b3c8-11d4-7e5a-b0c1-2d3e4f5a6b7c", "event_type": "ledger.entry_posted", "schema_version": 3, "schema_ref": "events/ledger.entry_posted@3", "aggregate_type": "account", "aggregate_id": "1001", "aggregate_version": 58, "writer_epoch": 1, "origin_region": "us-east-1", "occurred_at": "2026-09-27T19:12:40.551Z", "producer": "ledger-service", "correlation_id": "req_c02e", "traceparent": "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01", "payload": { "transfer_id": 88210391, "amount_minor": 1500, "currency": "USD", "balance_after_minor": 49000, "counter_account_id": 7702 } }
| New field | Why |
|---|---|
schema_ref | The registered schema this payload was validated against |
origin_region | Where the aggregate's writer lives; consumers in other regions keep order per origin |
writer_epoch | A counter bumped each time a service's database is promoted in another region. Consumers compare (writer_epoch, aggregate_version), which catches versions reused after a regional failover (step 3.4) |
traceparent | W3C trace context, so a trace can follow an event across services |
Schema registry rules (step 3.1): one schema per event type, compatibility mode FULL_ALL, checked in CI and at runtime.
Replay API (step 3.3), for the platform's replay service:
httpPOST /v1/replays HTTP/1.1 Content-Type: application/json { "source": { "topic": "ledger.entries", "region": "us-east-1", "from": "2026-08-01T00:00:00Z", "to": "2026-09-27T00:00:00Z" }, "target_topic": "replay.billing-rebuild-20260927", "start_from_snapshot": "billing-projection@2026-08-01", "max_events_per_second": 200000, "honor_quarantine_list": true }
httpHTTP/1.1 202 Accepted Content-Type: application/json { "replay_id": "rp_0927_01", "status": "RUNNING", "estimated_events": 5700000000 }
A replay never writes to the live topic. It writes to a dedicated replay topic that only the rebuilding consumer reads. The estimate here: August 1 to September 27 is 57 days × 100M ledger events a day = 5.7B events.
R3.4 Design Evolution: From One Outbox to a Backbone
Step 3.1: "A Producer Renamed a Field and Three Consumers Broke"
The problem: the orders team renamed amount_minor to total_minor in order.paid. Their tests passed. Three consumers in other teams started failing in production, and one of them silently wrote zeros.
What would you do?
Step 3.2: "Polling Relays Can't Keep Up at 2B a Day"
The problem: the busiest database (the ledger in us-east-1) now peaks at 12,000 events a second. Its relay deletes 12,000 rows a second and runs over a thousand claim transactions a second on the writer. Forty teams each maintain their own relay. What would you do?
The slot: the one thing that can hurt the database. A slot makes PostgreSQL keep every WAL segment from the slot's confirmed position onward. If the connector stalls (crashed, misconfigured, MSK unreachable), WAL piles up:
| Database | Where WAL lives | What a stalled slot does |
|---|---|---|
| RDS for PostgreSQL (10 services) | The instance's own storage | Fills the disk, and a full disk stops the database: a producer outage caused by a consumer of its log |
| Aurora PostgreSQL (30 services) | The cluster's shared storage volume | Grows storage and its bill; alarms matter as much, the cliff is further away |
The ledger's numbers (from Round 2's ~2 KB of WAL per event): 100M events/day × 2 KB ≈ 200 GB/day, about 8.3 GB an hour on average and 24 MB/s ≈ 86 GB an hour at the 12,000/s peak.
Defenses, in order:
- Alarm on time, then on bytes. Debezium writes a heartbeat row every 10 seconds (
heartbeat.action.queryinto a small heartbeat table that is part of the publication). The platform measures end-to-end lag as "now minus the commit time of the newest heartbeat row that has come through the connector": page at more than 60 s for 5 minutes. As a backstop, alarm on the CloudWatch metricOldestReplicationSlotLag: warn at 5 GB, page at 25 GB (about 17 minutes at the ledger's peak, 3 hours at its average). - Heartbeats also keep the slot moving. If a database's outbox is quiet but its other tables are busy, the connector has nothing to confirm, the slot never advances, and WAL piles up anyway. The heartbeat row gives it something to confirm.
- Cap retention where the engine allows it. PostgreSQL 13 and later has
max_slot_wal_keep_size(default −1, unlimited). Set below the free disk (for example 200 GB on an RDS volume with at least 400 GB free), PostgreSQL invalidates the slot instead of filling the disk. We set it on the RDS services, and check per engine version whether the Aurora parameter group exposes it. - Recover without trusting the stream. An invalidated or lost slot means the stream may have missed events. We create a new slot (it retains WAL from that point), republish the outbox rows (kept 3 days) written since the last event the topic has from that database, minus a margin, wait for the acks, and only then start the connector on the new slot. Republishing first keeps each partition in version order: the old gap lands before anything newer. Consumers skip the repeats. Because rows are kept 3 days and we page within minutes, the stream never loses an event it can't recover.
The partitioned outbox for CDC, in the column names the outbox router expects:
sqlCREATE TABLE outbox ( id UUID NOT NULL, -- event_id created_at TIMESTAMPTZ NOT NULL DEFAULT now(), aggregatetype TEXT NOT NULL, -- routes to a topic, e.g. 'ledger' -> ledger.entries aggregateid TEXT NOT NULL, -- the Kafka key type TEXT NOT NULL, -- event_type payload JSONB NOT NULL, -- the full envelope, incl. version and epoch PRIMARY KEY (id, created_at) -- must include the partition key ) PARTITION BY RANGE (created_at); CREATE PUBLICATION outbox_pub FOR TABLE outbox, cdc_heartbeat WITH (publish = 'insert', publish_via_partition_root = true);
publish = 'insert' publishes only inserts (the router ignores anything else anyway), and publish_via_partition_root makes rows from every daily partition appear as rows of outbox, so new partitions need no connector change. Dropping a partition produces no row events.
Two AWS facts that shape this: logical decoding runs only on the writer. Aurora PostgreSQL doesn't support logical decoding from readers, so the connector can't be moved off the writer. And the connector is single-task, so we plan on about 25,000 small events a second per connector (a planning figure we load-test). The ledger's 12,000/s peak fits; a database that outgrows it gets split (the wallet loop shards its ledger for the same reason), each part with its own slot and connector.
Why CDC instead of a relay polling every 500 ms? A poller adds a query every half-second to the busiest database, forever; it needs claims and deletes that churn the table and its index; it adds up to the polling interval in latency; and ordering gets subtle with several pollers. Reading the WAL adds almost no query load, publishes as soon as the commit is in the log, and preserves commit order. The price is the connector platform and the slot. At one service and 6,000/s (Round 2) the relay was the better deal; at 40 services and 12,000/s on one database, CDC is.
Primitive: Change Data Capture and the Outbox Pattern
Step 3.3: "Rebuild the Billing Read Model From Scratch"
The problem: billing's per-customer monthly totals were computed with a rounding bug for 6 weeks. Audit wants the read model rebuilt correctly from the history of events, and wants to know, for any account, "what was the balance on June 1?" What would you do?
Why point-in-time queries get slow in event-sourced systems. To know an aggregate's state at a moment, you replay its events up to that moment. An account with 2 million events means reading and applying 2 million events for one question. The fix is snapshots: store the aggregate's state every N events (we use 10,000), start from the snapshot just before the moment, and replay at most 10,000 events. For the ledger it's even simpler, because each entry already stores the balance after it.
Why choose event sourcing and CQRS over plain tables for a core ledger? For: a complete audit trail by construction, exact answers to "what was true at time T", and read models that can be rebuilt or added later from history. Against: reads are eventually consistent with writes (projections lag), events live forever so their schemas must stay readable forever, and the design is harder to build and to explain than updating a row. For a ledger, the audit trail and temporal queries are the product, so it's worth it. For a product catalog, usually not.
Primitive: Event Sourcing and CQRS · Drill: The banking ledger that took 3 days to replay (both of its questions are answered in the two paragraphs above)
Step 3.4: "Events Must Flow Between Regions"
The problem: a customer's home is eu-west-1, but the fraud team's model runs in us-east-1 and needs that customer's events. And if eu-west-1 fails, the services there fail over to another region. What would you do?
"Lost together" isn't the whole story. It's tempting to say "the outbox row and the state are lost together in a regional failover, so nothing disagrees". That holds for events that never left the old region. It doesn't hold for anything that left before replication: a 201 reply to a client, or an event the old region's connector already published and MSK Replicator already copied. Those are the forks in the table above. The lost-tail reconciliation lists them from both sides once the old region is reachable again: Aurora attempts to snapshot the old primary's storage at the point of failure (rds:unplanned-global-failover-...); when it succeeds, reconciliation reads the lost tail from it. Each fork is either re-applied through the service's API with its original idempotency key or compensated, owned by the service's team. The wallet loop applies the same idea to money.
Primitive: Cloud Disaster Recovery and Multi-Region Active-Active
Step 3.5: "Can We Get Exactly-Once?"
The problem: a team lead says: "Kafka has exactly-once now, so we don't need the inbox." Another team promises its partners "exactly-once delivery". What would you do?
Primitive: Database Isolation Levels, ACID & Concurrency Anomalies
Step 3.6: "Should Every Team Build This?"
The problem: 40 teams, each wiring its own connector, schema checks, DLQ, replay scripts and dashboards. Some do it well. What would you do?
Round 3 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 3.1 | Producers break consumers | Glue Schema Registry, FULL_ALL, CI and runtime checks; new types for breaking changes | Governance |
| 3.2 | Relays at 2B/day | Debezium CDC on the outbox via MSK Connect; insert-only daily partitions; slot alarms, heartbeats, caps, republish | A connector per database; the slot |
| 3.3 | Rebuild read models | S3 archive, replay to shadow projections, snapshots, quarantine list | Storage; replay time |
| 3.4 | Events across regions | Regional MSK, MSK Replicator for global topics, one writer per aggregate, writer_epoch | Lag, transfer charges, fork reconciliation |
| 3.5 | "Exactly-once?" | A per-hop contract; idempotency everywhere outside one transaction | Harder to sell |
| 3.6 | 40 teams | An event platform | A platform team |
R3.5 Global Architecture
One region
Synthesizing vector architecture diagram...
Everything a region needs lives in the region: each database's log feeds one connector, topics feed consumers and the archive, and replays run on their own topics so live consumers never see old events.
Three regions
Synthesizing vector architecture diagram...
Six one-way MSK Replicators, one per ordered pair of regions, copy only the global topics (about 10% of events). Local topics never leave their region.
Trace 1: a schema change caught in CI
Synthesizing vector architecture diagram...
The contract is enforced where the change is made, and again where the event is produced.
Trace 2: rebuilding billing's projection
Synthesizing vector architecture diagram...
5.7B events at 200,000/s is about 8 hours (R3.6). The old projection serves, flagged, until the switch.
Trace 3: a region fails over
Synthesizing vector architecture diagram...
The slot is created before writes open, so every new write reaches the stream through it. The republish reads the promoted database's own outbox and finishes before writes open and before the connector starts, so the pre-promotion tail lands in each partition ahead of every new-epoch event. The secondary clusters already run with rds.logical_replication=1 in their parameter group (it is a static parameter that needs a reboot, so it can't be switched on mid-failover). Only the lost tail (writes that never replicated) needs reconciliation.
Time budget for this failover (a target we rehearse, not a guarantee): detect and decide ≤ 10 min + fence ~1 min + promote the global cluster ~1–2 min (typical) + epoch and slot ~1 min + republish the tail and wait for acks ~1 min ≈ 15 min until writes resume; then start the connector (up to about 10 minutes on MSK Connect, an assumption) ≈ 25 min until new events flow. Events written in between wait safely in the outbox, and the slot keeps their WAL.
R3.6 Numbers and Cost
Traffic
| Quantity | Math | Value |
|---|---|---|
| Events per day | given | 2,000,000,000 |
| Average rate | 2 × 10⁹ ÷ 86,400 s | ≈ 23,150/s |
| Peak | assume 3× average: many services and time zones smooth each other's peaks (Round 2's single ledger saw 10×) | 69,444 → plan 70,000/s |
| us-east-1 (50%) | 1.0B/day; 11,574/s avg; × 3 | 35,000/s peak |
| eu-west-1 (30%) | 0.6B/day; 6,944/s avg; × 3 | 21,000/s peak |
| ap-southeast-1 (20%) | 0.4B/day; 4,630/s avg; × 3 | 14,000/s peak |
| Busiest database | the ledger in us-east-1: 100M/day; 1,157/s avg × 10 | 12,000/s peak |
| Average per connector at peak | 70,000 ÷ 120 databases | ≈ 580/s |
Replication. 10% of each region's events go to both other regions:
| Region | Replicated out per day | Replicated in per day | Peak replicated in (× 3) |
|---|---|---|---|
| us-east-1 | 0.1B to each of 2 = 0.2B | 0.06B + 0.04B = 0.10B | ~3,500/s ≈ 3.5 MB/s |
| eu-west-1 | 0.06B × 2 = 0.12B | 0.10B + 0.04B = 0.14B | ~4,900/s ≈ 4.9 MB/s |
| ap-southeast-1 | 0.04B × 2 = 0.08B | 0.10B + 0.06B = 0.16B | ~5,600/s ≈ 5.6 MB/s |
In total 0.4B events a day, about 400 GB/day at ~1 KB each.
Broker sizing. We size at ~1 KB per event, uncompressed (compression is headroom). Planning rule (an assumption to confirm with a load test): each broker stays below 30 MB/s of inbound writes (leader plus follower copies) at peak, and below 50 MB/s outbound with one AZ lost. Outbound counts three consumer groups per event on average, follower fetches, and MSK Replicator reads. Brokers go in multiples of three, one third per AZ.
| Region | Ingress at peak (local + replicated in) | Inbound per broker (× 3 replicas ÷ brokers) | Outbound per broker after losing an AZ | Brokers |
|---|---|---|---|---|
| us-east-1 | 35 + 3.5 = 38.5 MB/s | 115.5 ÷ 6 ≈ 19 MB/s | (3 × 38.5 + 38.5 + 7) ÷ 4 ≈ 40 MB/s | 6 |
| eu-west-1 | 21 + 4.9 = 25.9 MB/s | 77.7 ÷ 6 ≈ 13 MB/s | (3 × 25.9 + 25.9 + 4.2) ÷ 4 ≈ 27 MB/s | 6 |
| ap-southeast-1 | 14 + 5.6 = 19.6 MB/s | 58.8 ÷ 3 ≈ 20 MB/s | (3 × 19.6 + 19.6 + 2.8) ÷ 2 ≈ 41 MB/s | 3 |
Losing an AZ doesn't raise a broker's inbound load (Kafka doesn't re-create the lost replicas elsewhere), but consumers and the one remaining follower per partition now fetch from fewer brokers. Europe needs 6: with 3 brokers, losing one AZ leaves 2 carrying (77.7 + 25.9 + 4.2) ÷ 2 ≈ 54 MB/s outbound, above our line.
Broker storage (7 days × 3 replicas):
| Region | Local | Replicated in | Used | Provisioned |
|---|---|---|---|---|
| us-east-1 | 1.0B × 1 KB × 7 × 3 = 21 TB | 0.10B × 1 KB × 7 × 3 = 2.1 TB | 23.1 TB | 6 × 5 TB = 30 TB |
| eu-west-1 | 12.6 TB | 2.94 TB | 15.5 TB | 6 × 3.5 TB = 21 TB |
| ap-southeast-1 | 8.4 TB | 3.36 TB | 11.8 TB | 3 × 5 TB = 15 TB |
CDC and the ledger's WAL: 100M × ~2 KB ≈ 200 GB/day; alarms at 5 GB (warn) and 25 GB (page) on slot lag, 60 s on heartbeat lag.
The archive. 2B × 1 KB = 2 TB/day raw. Assume Parquet with compression shrinks it about 4× (to be measured): 0.5 TB/day, about 182.5 TB for a year. The last 30 days stay in S3 Standard (15 TB); older files move to S3 Glacier Instant Retrieval (millisecond reads, cheaper storage, a per-GB retrieval charge), and expire after 365 days. The archive bucket is versioned, so the lifecycle rule also sets NoncurrentVersionExpiration (7 days): without it, overwritten or deleted files would keep their old versions, and storage, forever.
The ledger's replay. One year of ledger events is 100M × 365 = 36.5B events (≈ 36.5 TB raw, ~9.1 TB in Parquet). The replay reader aggregates per account in memory and writes one upsert per account per batch, so we plan on 200,000 events/s (a planning figure):
| Rebuild | Math | Time |
|---|---|---|
| A full year from nothing | 36.5B ÷ 200,000/s = 182,500 s | ≈ 51 h ≈ 2.1 days |
| From the last monthly snapshot (≤ 31 days) | 3.1B ÷ 200,000/s = 15,500 s | ≈ 4.3 h |
| The August 1 request in R3.3 (57 days) | 5.7B ÷ 200,000/s = 28,500 s | ≈ 7.9 h |
| One account's balance at a moment | ≤ 10,000 events after its snapshot | milliseconds |
Monthly cost. us-east-1 list prices; for the other regions we assume about +10% in eu-west-1 and +20% in ap-southeast-1 (Aurora's db.r7g.xlarge list price is about 10% and 20% higher there; MSK's regional differences should be checked). The 40 service databases are paid by their teams and not counted.
| Item | Math | Monthly |
|---|---|---|
MSK brokers, kafka.m7g.2xlarge ($0.816/h) | US 6 × 0.816 × 730 = $3,574; EU 6 × 0.8976 × 730 = $3,931; AP 3 × 0.9792 × 730 = $2,144 | ≈ $9,649 |
| MSK storage | US 30,000 GB × $0.10 = $3,000; EU 21,000 × $0.11 = $2,310; AP 15,000 × $0.12 = $1,800 | ≈ $7,110 |
| MSK Connect ($0.11 per MCU-hour) | US 50 MCUs × 0.11 × 730 = $4,015; EU 40 × 0.121 × 730 = $3,533; AP 40 × 0.132 × 730 = $3,854 | ≈ $11,402 |
| MSK Replicator, per replicator | 6 × $0.30/h × 730 h | ≈ $1,314 |
| MSK Replicator, data | 400 GB/day × 30.4 ≈ 12,160 GB × $0.08 | ≈ $973 |
| Cross-region transfer | US 6,080 GB and EU 3,648 GB at $0.02/GB; AP 2,432 GB at $0.09/GB | ≈ $414 |
| S3 archive | 15 TB Standard × $0.023 + 167.5 TB Glacier IR × $0.004, plus regional uplift | ≈ $1,100 |
| Glue Schema Registry | no extra charge | $0 |
| Archivers, replay service, DLQ tooling (Fargate) | ≈ $1,500 | |
| Observability | ≈ $2,000 | |
| Total | ≈ $35,500/month |
About $0.58 per million events (35,462 ÷ 60,800M events a month), against $6.2 in Round 2. The difference is mostly that Round 2's number included the ledger's own database, and here each service pays for its own. The biggest single line is MSK Connect: 130 always-on connector MCUs.
R3.7 Trade-Offs
Outbox + CDC vs event sourcing
| Outbox + CDC (most services) | Event sourcing (the ledger) | |
|---|---|---|
| Source of truth | The service's tables | The event log |
| Events are | Announcements of state changes | The state itself |
| Point-in-time answers | Only what the tables keep | Exact, from the log plus snapshots |
| Schema burden | Tables change freely; events follow contracts | Every event version must be readable forever |
| Reads | Normal queries | Projections, eventually consistent |
| Choose when | Most services | Audit and history are the product |
Per-team relays vs a platform
| Per-team relays | Platform (our choice) | |
|---|---|---|
| Speed to start | Fast for the first team | Slower: the platform must exist |
| Correctness | 40 implementations of ordering, slots, DLQs | One, reviewed and tested |
| Cost | Hidden in 40 teams' time and incidents | Visible: a six-person team + ~$35,500/month |
| Flexibility | Anything goes | Templates; exceptions by review |
Strict vs lenient compatibility
FULL_ALL (our choice) | BACKWARD only | NONE | |
|---|---|---|---|
| Old consumers read new events | Yes | Not guaranteed | No |
| Replays through today's code | Yes, all versions | Only against the previous version, unless BACKWARD_ALL | No |
| Producer freedom | Add optional fields only | More | Total |
| Cost of a breaking change | A new event type and a migration | Coordinated deploys | Production incidents |
Closing the loop. Round 1 found the core idea: write the event where the change is written, and make every hop after that safe to repeat. Round 2 made it fast and safe for money under load. Round 3 made it everyone's: contracts, the log instead of relays, a history you can replay, regions with honest failover rules, and a platform that owns the sharp edges. Through all three rounds, the three invariants held: no lost events, no ghost events, effects once.
R3.8 Failure Modes
| Failure | What you'd see | How the design responds |
|---|---|---|
| A stuck replication slot | Heartbeat lag climbs; OldestReplicationSlotLag grows; on an RDS service, free storage falls | Page at 60 s of heartbeat lag. Restart or fix the connector. If free storage is at risk, drop the slot (or let max_slot_wal_keep_size invalidate it), then create a new one, republish from the outbox (kept 3 days) and wait for it to finish, and only then start the connector on the new slot. Consumers skip the repeats. |
| An incompatible schema deployed anyway (someone registered it by hand) | Consumers' DLQs fill with "unreadable schema" | Consumers fail closed, so nothing is applied wrongly. Roll back the producer; register a compatible version; replay the dead letters. Lock down registry write access to the CI role. |
| A replay overloads consumers | Shadow projection's database at 100% CPU; lag on its replay topic | Replays go to replay topics, never live topics, so live consumers are untouched. Lower the replay's rate; it resumes from its last position. |
| Duplicates and forks after a regional failover | A burst of skipped duplicates; some parked forks | Duplicates are harmless by design. Each fork is listed for the service's lost-tail reconciliation; nothing is silently merged. |
| MSK Replicator lag | Global topics in other regions fall behind | Consumers of replicated topics see older data; per-aggregate order still holds. Alarm on the replication latency; local processing is unaffected. |
| An archiver falls behind for longer than Kafka's retention | Gaps in the archive | Alarm on archiver lag at 24 hours, well inside Kafka's 7 days, so this should never happen. If it does, the events are gone from Kafka and older than the outboxes' 3 days, so the owning services regenerate them from their own tables with a backfill job, marked as backfilled: a last resort, which is why the alarm fires so early. |
R3.9 Runbook and Incident Response
Signals, per region OPS 8 · REL 6
| Signal | Alarm | Severity | First action |
|---|---|---|---|
| Outbox depth and oldest row's age (relay-based services) | age > 30 s | P2 | Check the relay and MSK |
| CDC heartbeat lag per database | > 60 s for 5 min | P2 | Check the connector's state and logs |
OldestReplicationSlotLag | > 5 GB warn; > 25 GB page | P3 / P1 | Slot-hang procedure below |
| Free storage (RDS services) | < 20% | P1 | Slot-hang procedure; this is the outage path |
| Oldest open transaction on any writer | > 10 min | P3 | End it; it blocks vacuum |
| DLQ inflow per consumer | ledger and audit: any; others: > 10/min | P1 / P3 | DLQ procedure below |
| Consumer lag per group | > 60 s for 5 min | P3 | Scale the consumer; check its database |
| Archiver lag | > 24 h | P2 | Fix before Kafka's 7-day retention runs out |
| Forks parked after a failover | any | P2 | Lost-tail reconciliation with the owning team |
| Reconciliation breaks (ledger vs consumers) | any | P1 | Same-day investigation |
Replication-slot hang procedure OPS 10
-
Look at the slot. Is it active? How far behind?
sqlSELECT slot_name, active, pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal FROM pg_replication_slots; -
Fix the reader first. Check the connector's state in MSK Connect and its logs. Most hangs are a failed task or MSK access; a restart resumes from the slot's position with nothing lost.
-
If the disk is at risk (RDS: free storage below 20% and falling), stop the connector, note the time of the last event the topic has from this database, and drop the slot:
SELECT pg_drop_replication_slot('ledger_outbox_slot');. Never drop an active slot, and never drop one without writing down that time. -
Recreate and republish. Create a new slot; run the republish job for outbox rows created since the noted time minus 5 minutes and wait for it to finish; then start the connector on the slot (configured not to snapshot). Republishing before streaming resumes keeps every partition in version order. Consumers skip the repeats.
-
Record it, with the time window republished.
DLQ replay procedure
- Read the reason (
GET /v1/admin/dead-letters). Producer bug or consumer bug? - Producer bug: the producing team publishes a corrective event; the original goes on the quarantine list so replays skip it. Consumer bug: fix and deploy the consumer.
- Replay (
POST /v1/admin/dead-letters/{id}/replay): the dead letter and then everything parked behind it for that aggregate, in version order. - Confirm the aggregate is unblocked and the next reconciliation run is clean.
Go deeper: CLI playbook
Plain commands, one at a time. Replace names and ARNs with real ones.
text# 1. Slot lag on the ledger's writer (bytes of WAL the slowest slot still needs) aws cloudwatch get-metric-statistics --region us-east-1 --namespace AWS/RDS --metric-name OldestReplicationSlotLag --dimensions Name=DBInstanceIdentifier,Value=ledger-use1-writer --start-time 2026-09-27T10:00:00Z --end-time 2026-09-27T11:00:00Z --period 60 --statistics Maximum # 2. Free storage on an RDS for PostgreSQL service database aws cloudwatch get-metric-statistics --region us-east-1 --namespace AWS/RDS --metric-name FreeStorageSpace --dimensions Name=DBInstanceIdentifier,Value=orders-use1 --start-time 2026-09-27T10:00:00Z --end-time 2026-09-27T11:00:00Z --period 60 --statistics Minimum # 3. Find the CDC connectors and check one's state aws kafkaconnect list-connectors --region us-east-1 --connector-name-prefix ledger- aws kafkaconnect describe-connector --region us-east-1 --connector-arn arn:aws:kafkaconnect:us-east-1:111122223333:connector/ledger-outbox-cdc/0a1b2c3d-4e5f-6a7b-8c9d-0e1f2a3b4c5d-1 # 4. Consumer lag of billing on the ledger topic aws cloudwatch get-metric-statistics --region us-east-1 --namespace AWS/Kafka --metric-name MaxOffsetLag --dimensions Name="Cluster Name",Value=backbone-use1 Name="Consumer Group",Value=billing Name=Topic,Value=ledger.entries --start-time 2026-09-27T10:00:00Z --end-time 2026-09-27T11:00:00Z --period 60 --statistics Maximum # 5. The replicators copying global topics into this region aws kafka list-replicators --region eu-west-1 # 6. The compatibility mode of one event type's schema aws glue get-schema --region us-east-1 --schema-id RegistryName=events,SchemaName=ledger.entry_posted # 7. The archive's lifecycle rules (check NoncurrentVersionExpiration is there) aws s3api get-bucket-lifecycle-configuration --region us-east-1 --bucket events-archive-use1
R3.10 Pillar Check
| Pillar | What Round 3 adds |
|---|---|
| Reliability | Regional clusters isolate failures; slot alarms, caps and republish from the outbox protect producers from their log readers; a rehearsed failover with an epoch fence and fork detection; a year of history in S3 REL 9 · REL 10 · REL 12 · REL 13 |
| Performance Efficiency | CDC instead of polling at 12K/s on one database; producers write to their own region's cluster instead of across an ocean, and only 10% of events cross regions; brokers sized per region from the math, including AZ loss; replays rate-limited on their own topics PERF 1 · PERF 4 |
| Security | Only each service's CDC connector may produce to its topics, and only CI may register schemas; schemas mark personal fields so the archive and replicas can be classified and access-controlled; encryption in transit and at rest in every region SEC 3 · SEC 7 · SEC 8 · SEC 9 |
| Cost Optimization | About $35,500 a month, $0.58 per million events; only 10% of events cross regions; the archive in Glacier Instant Retrieval after 30 days; a platform team weighed against 40 teams' duplicated effort COST 8 · COST 11 |
| Operational Excellence | Schema checks in CI and at runtime; runbooks for slot hangs, DLQ replays and regional failover; one dashboard per database and consumer group OPS 5 · OPS 6 · OPS 10 |
| Sustainability | Insert-only outboxes with partition drops instead of delete churn; compressed, tiered archive with an expiry; local topics stay local SUS 3 · SUS 4 |
R3.11 Round 3 Rubric and Follow-Ups
What an architect (L7) answer adds over L6
- Treats events as contracts: a registry, a compatibility mode chosen from who reads what (including replays), and enforcement the producer can't skip.
- Separates CDC on an outbox from CDC on business tables, and moves to CDC for a reason with numbers.
- Knows the replication slot is the sharp edge, how it differs on RDS and Aurora, and designs alarms, caps and a republish path that match the outbox's retention.
- Designs replay as a product: archive, shadow projections, snapshots, rate limits, a quarantine list.
- Uses one writer per aggregate across regions, and handles the version reuse a failover causes with an epoch fenced inside the database.
- States exactly-once per hop, and requires idempotency everywhere else.
- Weighs a platform team against duplicated effort, not only against the AWS bill.
Follow-up questions
-
"Why not have Debezium capture the ledger's
ledger_entriestable directly? The entries already are events." Answer: it's tempting for the ledger, because its table is append-only and gap-free. But the table's columns become the contract: any schema change in the ledger breaks consumers, and there's nowhere to put the envelope (schema reference, epoch, trace). An outbox row costs one insert per event and keeps the contract separate from the storage. Some teams do capture business tables for internal analytics feeds; that's a different consumer with a different contract. -
"A consumer in ap-southeast-1 needs European events within 100 ms." Answer: cross-region replication adds network time between Europe and Singapore (well over 100 ms round trip) plus replication batching. We can't promise 100 ms from a region that far away. Options: run that consumer in eu-west-1 next to the data, or relax the requirement. Say it plainly.
-
"Could we drop the per-service outbox and have services write events straight to Kafka with Kafka transactions?" Answer: no. A normal Kafka transaction makes several Kafka writes atomic with each other, not with a PostgreSQL commit. Kafka only recently added opt-in 2PC participation (KIP-939, off by default and not something we can assume on our cluster); even with it, a prepared transaction holds its row locks until the coordinator decides. The service would be back to a dual write: commit the database, then commit the Kafka transaction, and crash in between. The outbox is what ties the event to the database commit.
Interview gotchas from this round
| Gotcha | Why it's wrong |
|---|---|
| "CDC means we don't need an outbox" | CDC on business tables publishes rows, not contracts; the outbox stays and CDC replaces the relay. |
"BACKWARD compatibility is enough" | Old consumers must read new events too, and replays cross every version: FULL_ALL. |
| "A stuck slot is a consumer problem" | It fills the producer's disk on a fixed-storage database. |
| "Replay into the live topic" | Live consumers see old events as new. Replay to a dedicated topic and a shadow projection. |
| "Kafka gives us exactly-once" | Only inside Kafka transactions; not from a database, not to an email server. |
| "One global Kafka cluster" | Every produce crosses an ocean; one cluster's trouble is everyone's. |
Loop Closer: Interview Strategy for All Three Rounds
How to Run Each 60-Minute Round
| Time | Round 1 | Round 2 | Round 3 |
|---|---|---|---|
| 0–5 min | Scoping: consumers, delay, duplicates, order, existing broker | Restate Round 1 in 60 seconds | Restate Round 2 in 60 seconds |
| 5–15 min | Requirements, API, the event envelope | Scope raise → what breaks | Scope raise → what breaks |
| 15–40 min | Steps 1.0–1.6: dual write, outbox, relay, at-least-once and inboxes, SKIP LOCKED, order, growth | Steps 2.1–2.6: dispatch approaches, locking for money, poison events, outages, vacuum, audit | Steps 3.1–3.6: contracts, CDC and slots, replay, regions, exactly-once, platform |
| 40–50 min | Numbers (outbox size, relay load) and trade-offs (polling vs push, delete vs mark, outbox vs 2PC) | Numbers (13.7 TB vs kilobytes; outage drain), cost, the two comparison tables | Numbers per region, archive and replay times, cost; outbox vs event sourcing |
| 50–60 min | Failures + pillar check | Failures, gotchas + pillar check | Failures, runbook, pillar check |
For how to spend a single 45-minute round, see the 45-minute interview blueprint.
The Two Sentences That Matter Most
- Opening a round: "Before I design: who consumes these events, how fast, can they see one twice, and must they be in order per entity?"
- When the scope is raised: "Here's what breaks. I'll fix anything that could lose or invent an event first, then latency, then cost."
And the one sentence specific to this system: "The event is written in the same transaction as the change, every hop after that may repeat but never reorders an entity's events, and every consumer applies each version exactly once in its own transaction."
Well-Architected Review Sheet
| Pillar | Question you'll hear | One-sentence answer | Round | Backed by |
|---|---|---|---|---|
| Reliability | "What if the broker is down?" (REL 5) | Writes keep committing; events wait in the outbox and drain later: 3.6M rows in about 4 minutes after a 10-minute outage at peak. | 1–2 | R1.9, Step 2.4 |
| "What if a message is delivered twice?" (REL 4) | Every consumer applies an event in one transaction with an inbox or version check, so a repeat does nothing. | 1 | Step 1.3 | |
| "What if one bad event arrives?" (REL 10) | It goes to a dead-letter queue and blocks only its own entity; everything else keeps flowing. | 2 | Step 2.3 | |
| "What if a region fails?" (REL 13) | Promote the database elsewhere, bump the epoch inside it before writes open, republish the tail, and reconcile forks. | 3 | Step 3.4 | |
| Performance | "How do you get events out in milliseconds?" (PERF 1) | A coalescing poke after commit wakes the relay; at larger scale, CDC reads the log. | 2–3 | Steps 2.1, 3.2 |
| "How do you stop two transfers overdrawing an account?" (PERF 3) | Conditional updates lock the account rows in ascending order until commit. | 2 | Step 2.2 | |
| Cost | "What does it cost?" (COST 5) | About $1,060, $9,440 and $35,500 a month; $6.2 per million events with the ledger's database, $0.58 for the shared backbone. | 1–3 | R1.7, R2.6, R3.6 |
| "Build a platform or let teams build their own?" (COST 11) | A six-person platform team instead of 40 teams each running and debugging a relay. | 3 | Step 3.6 | |
| Operations | "How do you know every event arrived?" (OPS 8) | Outbox age, heartbeat lag, consumer version gaps and a daily reconciliation against the source of truth. | 1–3 | Step 2.6, R3.9 |
| "How do you stop a producer breaking consumers?" (OPS 6) | FULL_ALL schemas checked in CI, and production can only use registered versions. | 3 | Step 3.1 | |
| Security | "Who can publish to a topic?" (SEC 3) | Only the owning service's relay or connector; consumers can't write back into producers. | 2–3 | R2.10, R3.10 |
| Sustainability | "Where is the waste?" (SUS 4) | Not in the database: delete-on-success or partition drops instead of 13.7 TB a year of dead history; a compressed archive that expires. | 1–3 | Step 1.6, R3.6 |
Rubric Across Levels
| Dimension | L5 (Round 1) | L6 (Round 2) | L7 (Round 3) |
|---|---|---|---|
| Atomicity | Outbox row in the same transaction; explains the dual write. | Keeps it under contention: rows locked in order, rejections stored with keys. | Keeps it across 40 services and regional failover with an epoch. |
| Delivery | At-least-once; inbox in the consumer's transaction; offsets after commit. | Poison events contained per aggregate; outages absorbed and drained on a budget. | A per-hop exactly-once contract; republish from stored rows, never from memory. |
| Order | Buckets, key by aggregate, version check. | Order kept through DLQs, replays and parallel relays. | Order per aggregate across regions; forks detected, not merged. |
| Dispatch | Polling, and why it's enough. | Four approaches compared with real limits; a choice with numbers. | CDC on the outbox with slot alarms, heartbeats, caps and a republish path. |
| Storage | Delete on success; 274 GB/year avoided. | Vacuum, long transactions, partitions the schema can drop. | Insert-only outboxes; a tiered archive with correct lifecycle rules. |
| Proof | Version gaps alarm. | Daily per-account reconciliation against the ledger with a grace window. | Replayable history, shadow projections and a quarantine list. |
| Evolving under new scope | Builds from a dual write, one problem at a time. | Opens with "what breaks" and fixes lost and wrong events first. | Turns a pattern into a platform without weakening any invariant. |