Change Streams & the Transactional Outbox
The Dual-Write That Broke Search Consistency
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
Part 0. Start here
The problem: one change, four systems
An order service handles 500 order changes a second (our example's number). After each commit it writes the search index, then the cache, then asks the notifier to email the customer: four writes to four systems, one after another. Deploys, scale-downs and crashes kill a request between its commit and its last write. Suppose that happens to 1 request in 100,000 (an assumption, to have a number):
- 500 changes a second × 86,400 seconds = 43,200,000 changes a day
- 43,200,000 ÷ 100,000 = 432 orders a day whose index, cache or email silently disagrees with the database
Nothing anywhere records those 432. The database says SHIPPED, the search page says PAID, and no log line, alarm or queue knows the two differ.
The team "fixes" it by feeding the index from the table's change stream (DynamoDB Streams, the table's own log of changes) with a Lambda function on default settings. A month later one malformed record makes the function throw. Lambda retries the batch, in order, and everything behind it on that shard waits. AWS's own documentation says that with the default settings "a bad record can block processing on the affected shard for up to one day". A DynamoDB stream keeps records for only 24 hours, and the records queued behind the bad one are only a little younger than it, so many of them expire before the function can catch up. The first fix lost 432 changes a day; the second can lose a day's worth in one go.
The order row commits. Where should the index, the cache and the notifier learn about it from, so that no crash, anywhere, can make one of them miss it? And what new problems does your answer create?
The big picture
Synthesizing vector architecture diagram...
What to notice: the only write the service makes is one transaction. A mover copies what committed into a topic, and every consumer decides for itself, by version, whether an event is new, old or early.
What you'll be able to do after this page
- Explain the three ways writing to two systems goes wrong, and when a direct write is still allowed (Part 1).
- Write the outbox transaction and say what it guarantees and what it doesn't (Part 2).
- Walk through a polling relay, explain why its delivery is at least once, and apply events by version (Part 3).
- Keep per-order order across every hop, including resizes, rebalances and multi-event transactions (Part 4).
- Explain log-based change data capture on PostgreSQL and MySQL, and trace one change through every layer (Part 5).
- Handle a record that always fails without blocking everyone behind it, and repair from the source of truth (Part 6).
- Refill a gap after a lost position, and bootstrap a consumer that was down longer than retention (Part 7).
- Change an event's schema safely, and make deletes stick (Part 8).
- Run the same design on DynamoDB Streams, Lambda, Pipes and Kinesis Data Streams for DynamoDB (Part 9).
- Choose between dual writes, an outbox, change data capture, event sourcing and two-phase commit on equal terms (Part 10).
- Map all of it to AWS services, and name the look-alikes that are not change streams (Part 11).
You may have arrived from a step that relies on this: step 1.1 of the transactional outbox loop (the order committed but the publish timed out), step 2.4 of the email loop (a parser writing DynamoDB and the index), step R1.6 of the ride-sharing loop (a status hint fed from the stream, not a second write), or step 3.1 of the ad-click aggregation loop (note the stream position, then scan). This page is the "why" behind all four.
Part 1. Why writing twice fails
At 09:05:00, before our service had an outbox, order o7 was paid. The service committed version 2, updated the search index, and then the pod was killed by a deploy. The cache still says PLACED, the customer never gets "payment received", and nothing records that anything went wrong. This Part shows why no amount of care in the request fixes that.
Trace: the pod dies after the index write
A dual write means the application writes two (or more) systems one after the other, with nothing tying them together. Here is the counterfactual timeline B1 to B3, before the outbox existed. Every trace on this page was produced by running a private reference implementation of the example, not worked out by hand.
Synthesizing vector architecture diagram...
What to notice: the database and the index moved on; the cache and the notifier never heard. Nothing in any of the four systems says so.
A minute later two app servers change o7 at almost the same moment. Server A commits version 3 (a new address) at 09:06:00.000; server B commits version 4 (SHIPPED) at .010, after A's commit released the row lock. B happens to be quicker on the way out:
Synthesizing vector architecture diagram...
What to notice: the database serialized the two commits, but nothing serialized the two index writes. The slower request overwrote the newer document. The cache writes crossed the same way (B at .040, A at .050), so the cache also ends at version 3.
Three ways two writes disagree
| Shape | What happens | In our trace |
|---|---|---|
| Lost | The first write commits; the process dies before the second | B2: the cache and the email never happen |
| Ghost | The second write happens; the first fails or rolls back | If the service had emailed before committing, a failed commit would leave "payment received" for a payment that never happened |
| Reordered | Two writers commit in one order and write the copy in the other | B3: the index and cache hold v3 while the database holds v4 |
"Publish first, then commit." Does that fix B2?
Snapshot S1: four systems, four answers
Synthesizing vector architecture diagram...
After B3: the database says SHIPPED v4; the index and the cache say PAID v3; one email is missing. Each copy is wrong in its own way, and none of them knows.
One write, then followers
The mental model for the rest of the page: the change and its announcement are one write; everything else is a follower of that write. A follower reads what committed, in order, and may read it more than once. It never needs the writer to be alive.
The one allowed dual write. A direct write is fine as a latency hint, when two things are true: (a) a change stream delivers every change anyway, and (b) both paths apply by version or monotonically (for example "set if larger"), so the stream heals whatever the hint missed or reordered. The leaderboard loop (steps 1.3 and 2.6) writes scores directly for speed and re-applies from the stream using the current stored value; the Drive loop (step 2.3) sends a fast-path nudge and keeps the stream as the reliable path. A direct write with no stream behind it is a dual write, however fast.
Three rules the design never breaks
| Rule | What breaks without it |
|---|---|
| The change and its announcement are one write: an outbox row in the same transaction, or the database's own log | Lost and ghost events (B2), with nothing that records them |
| Every event carries the aggregate's version, and every consumer applies by version: skip what it has, apply exactly the next one, park anything further ahead | Duplicates double-count and late events overwrite newer state (B3) |
| Every consumer has a way back to the truth: a repair path that re-reads the source, and a periodic compare | One bad record or one lost position leaves a copy wrong forever |
The example we follow
Real systems carry millions of orders, which is too many to watch. So one story runs through the whole page: follow orders o7, o8 and o9, with o7 as the thread from placement to erasure. An aggregate is the unit whose changes must stay in order (here, one order); its version goes up by one with every change.
| Setting | Our example | At real scale |
|---|---|---|
| Source of truth | One PostgreSQL-like database: tables orders (order_id, status, total, version) and outbox | Aurora or RDS PostgreSQL or MySQL; DynamoDB in Part 9 |
| Outbox row | seq (identity), event_id, aggregate_id, aggregate_version, event_type, schema_version, payload, writer_epoch, created_at | Debezium's outbox router expects id, aggregatetype, aggregateid, type, payload by default |
| Orders | o7 (customer c1), o8, o9; each change bumps version in the same statement | millions |
| Relay (Parts 3 and 4) | Polls every 1 s; claims every free bucket that has rows (up to 4 buckets; bucket = order number mod 4: o7 → 3, o8 → 0, o9 → 1) with FOR UPDATE SKIP LOCKED, publishes each bucket's rows in seq order, deletes them and commits | The outbox loop: 16, then 64 buckets; a 1 s sweep plus an in-process signal |
| Broker topic | orders.events, 2 partitions, key = order_id. The example fixes the key-to-partition mapping with a lookup table: o7 → P1, o8 → P0, o9 → P1. In Part 4's resize side trace only, 4 partitions: o7 → P3, o8 → P0, o9 → P1 | Kafka's default partitioner: a hash of the key modulo the partition count (murmur2 in the Java client). Kinesis: an MD5 hash of the partition key onto shard hash ranges |
| Broker retention | 24 hours, the same as DynamoDB Streams, so one number serves both tracks | Kafka's default is 7 days (log.retention.hours 168); Kinesis 24 hours by default, up to 365 days; DynamoDB Streams 24 hours, fixed |
| Consumers | Indexer (search documents), cache updater (key order:o7), notifier (emails the customer on placed, paid, shipped and delivered, and appends each email to the customer's activity feed, a gapless numbered list c1#1, c1#2, … whose entries store src = the event's event_id) | OpenSearch, Valkey, an email provider |
| Consumer rule | Apply an event only if its version is exactly the next for that key; skip versions already applied; park anything further ahead | The outbox loop's step 1.5 |
| Consumer group | A member's partitions move to another member 45 s after it stops heartbeating | Kafka's consumer session.timeout.ms defaults to 45 s |
| Error boundary | Retry 3 times (after 0.1 s, 1 s, 5 s), then dead-letter into the consumer's own dead_letters table, block that key for this consumer, commit the offset | The outbox loop's step 2.3 |
| Change data capture (from Part 5) | One connector on one logical slot, started with snapshot mode no_data. The outbox becomes insert-only; old rows go by dropping daily partitions after 3 days, never by DELETE. The connector writes a heartbeat row every 10 s | Debezium on MSK Connect; 3 to 7 days of outbox in the loops |
| Slot cap | In the example, the cap on retained WAL is reached 6 hours into a stall (a stated outcome) | Derived in Part 5 |
| Periodic compare (Part 6) | A job compares each order in a read replica with the index and writes repairs; the replica lags up to 40 s in event 28 | Nightly or hourly in the loops |
| Clock | Day and wall time, D1 10:00:00 onward |
The story in nine beats: the four-writes failure (this Part), one transaction (Part 2), the relay dies after publishing (Part 3), two relays, a resize and a rebalance (Part 4), the log becomes the relay (Part 5), a poison record (Part 6), the slot is lost and the index is down 30 hours (Part 7), a new field and an erasure (Part 8), and the same story on DynamoDB (Part 9). Each Part shows only its own slice of events; the full event table is in Part 12.
This answers the first question of the drill The Dual-Write That Broke Search Consistency: when the commit succeeds and the publish fails, the database and every copy drift apart silently, and no retry in the request can close that window; the fix is to make the event part of the commit. Its second question (log tailing versus polling every 500 ms) is answered in Part 5.
What to remember from Part 1
- Two systems written one after the other will disagree after some crash, and nothing records it.
- The fix is one write: the change and its announcement commit together.
- A direct write is allowed only as a hint, with a versioned, stream-fed path as the guarantee.
Part 2. The outbox: one transaction
The announcement has to commit with the change or not at all. The simplest place to put it is the same database, in the same transaction, as a row in a table built for the purpose.
One transaction
Event 4 at 10:01:00.0: o7 is paid. Here is the whole transaction:
sqlBEGIN; UPDATE orders SET status = 'PAID', version = version + 1 WHERE order_id = 'o7' RETURNING version; -- 2, and o7's row lock is held until COMMIT INSERT INTO outbox (event_id, aggregate_id, aggregate_version, event_type, schema_version, payload, writer_epoch, bucket) VALUES (gen_random_uuid(), 'o7', 2, 'order.paid', 1, '{"order_id": "o7", "status": "PAID", "total": 100}', (SELECT epoch FROM writer_epoch), 3); -- seq 102 comes from the identity COMMIT;
The table:
sqlCREATE TABLE outbox ( seq bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, event_id uuid NOT NULL UNIQUE, aggregate_id text NOT NULL, aggregate_version bigint NOT NULL, event_type text NOT NULL, schema_version int NOT NULL, payload jsonb NOT NULL, writer_epoch int NOT NULL, bucket smallint NOT NULL, created_at timestamptz NOT NULL DEFAULT now() );
Why does the UPDATE of the order come before the INSERT into the outbox?
The table and the envelope
| Field | Why it's there |
|---|---|
seq | An identity column: a number given at insert, used by the relay to publish each bucket in insert order. It is not commit order (Part 3) |
event_id | A unique id for the event; consumers and page 04's idempotency keys use it, and derived records store it as src |
aggregate_id | The key: the broker partitions by it, so one order's events stay together |
aggregate_version | The order's version after this change; consumers apply by it |
event_type, schema_version | What happened, and which schema the payload follows (Part 8) |
payload | What consumers need, validated against its schema and size limit before the insert, so a bad or oversized event fails the business transaction instead of reaching the stream (about 1 MB is Kafka's default maximum message size) |
writer_epoch | Copied from a one-row table in this database; bumped when a failover promotes a new writer. How it is bumped and fenced is in Leases, Fencing Tokens & Distributed Locks (Part 7); how consumers use it is in Part 4 |
bucket | Order number mod 4: the unit a relay claims (Part 4) |
What it guarantees, and what it doesn't
| A crash … | What exists afterwards | What happens next |
|---|---|---|
Before COMMIT | Neither the order change nor the outbox row | The client sees an error or a timeout and retries with its idempotency key (page 04) |
After COMMIT, before the reply | Both | The event will be published; the client's retry is recognised as a repeat (page 04) |
So the outbox guarantees exactly one thing: the event exists if and only if the change committed. It does not publish anything, does not keep order across different orders, and does not stop a consumer from seeing the event twice. The outbox is a promise to publish, not a publish.
Events you mean, not rows that changed
An outbox row is an event the producer means: order.paid, with a payload designed for consumers and a version. The alternative, reading raw row changes from the business tables (Part 5), gives consumers the table's columns as they are today; a column rename becomes a breaking change for every consumer. Part 5 compares the two on equal terms.
The same idea works on a phone: the mobile chat loop (step 1.1) writes a message and its SENDING row in one local transaction, and a sender drains them.
What to remember from Part 2
- The outbox row commits with the change or not at all.
- Update the aggregate first: its row lock puts versions in commit order.
- The outbox is a promise to publish, not a publish.
Part 3. The relay and at-least-once delivery
The outbox holds rows; the index, the cache and the notifier read a topic. Something must move rows from one to the other, and that something can die at the worst moment. At 10:01:00.6 it does.
Trace: the relay dies after publishing
The first minute went smoothly. At 10:00:00.0 o7 was placed (outbox seq 101); at 10:00:00.6 relay R1 claimed bucket 3, the only bucket with rows, published 101 to P1@0 (partition 1, offset 0), got the acknowledgement, deleted the row and committed; at 10:00:00.7 all three consumers applied o7@1, the customer got "order received", and the feed got c1#1.
At 10:01:00.0 o7 is paid (seq 102) and at 10:01:00.1 o8 is placed (seq 103). From the reference implementation:
Synthesizing vector architecture diagram...
What to notice: the topic now holds each of the two events twice (P0@0 and P0@1, P1@1 and P1@2). Nothing was lost, and nothing downstream changed twice, because every consumer compared versions.
Snapshot S2 (after event 6, before R1 restarts):
| Lane | State |
|---|---|
| Database | orders: o7 v2 PAID, o8 v1 PLACED. outbox: rows 102 (o7 v2) and 103 (o8 v1) still there |
| Mover | R1 dead; its claim and its DELETE rolled back with its transaction |
| Topic | P0: @0 o8 v1 (seq 103). P1: @0 o7 v1 (seq 101), @1 o7 v2 (seq 102) |
| Consumers | Indexer, cache and notifier all at o7 → 2, o8 → 1. Emails: "order received" (o7), "order received" (o8), "payment received" (o7). Feed: c1#1, c1#2 |
The topic already has what the outbox still holds. When R1 comes back, it can't know that, so it publishes both rows again.
R1 died after the broker acknowledged but before its DELETE committed. Which order of "publish" and "delete" would have avoided the duplicate?
The relay loop
textRELAY LOOP (every 1 s, and again at once while work was found) BEGIN buckets = up to 4 bucket rows that have outbox rows, claimed with FOR UPDATE SKIP LOCKED -- the claim: page 03, Part 9 if none: COMMIT, sleep until the next tick or a signal for each claimed bucket b: rows = SELECT * FROM outbox WHERE bucket = b ORDER BY seq send each row to orders.events, key = aggregate_id wait for every acknowledgement -- timeout: ROLLBACK, back off, try again DELETE FROM outbox WHERE seq IN (the published seqs) COMMIT
The relay holds a transaction open while it publishes, which is normally a bad idea; here it locks only bucket rows and outbox rows, never an order, and the publish has a hard timeout. The claim mechanism itself (SKIP LOCKED: take rows nobody holds, skip the rest) belongs to Leases, Fencing Tokens & Distributed Locks (Part 9).
Consumers apply by version
Duplicates are certain. "The broker removes duplicates" isn't true across a relay crash (the idempotent producer only removes its own retries within one session), and count = count + 1 turns every duplicate into a wrong number. The rule every consumer on this page follows:
textON EVENT (key K, version v) -- where last_version is stored: page 04 last = last_version[K] if v <= last: skip -- a duplicate or an old event elif v == last + 1: in ONE transaction: apply the effect, set last_version[K] = v then apply parked events for K while each is exactly the next else: park it under K -- something before it is missing; alarm if it stays commit the offset only after the effect -- our example commits after each batch event types this consumer ignores still advance last_version[K]
| Case | Example | Action |
|---|---|---|
v ≤ last | Event 8: o7 v2 arrives again when last is 2 | Skip; still commit the offset |
v = last + 1 | Event 6: o7 v2 when last is 1 | Apply and store v in one transaction; then unpark |
v > last + 1 | Part 4: o7 v4 arrives when last is 2 | Park; never apply, never drop |
Effects are set by version, never added: "status is PAID as of version 2", not "one more payment". An email can't be un-sent, so the notifier also stores an idempotency key per email (o7:2); where that key lives and for how long is page 04's (Idempotency & Effectively-Once Processing).
A new owner replays the old owner's work
Events 10 to 12 (Part 4) publish o7 v3 and v4 to P1@4 and P1@5. The notifier runs as two instances: N1 reads P1, N2 reads P0. At 10:03:00.7, N1 is killed halfway through v4:
| Time | Who | Record | What happens | State after |
|---|---|---|---|---|
| 10:03:00.7 | N1 | P1@4, o7 v3 | Exactly next (last 2): apply; address_changed sends no email | last_version[o7] = 3 |
| 10:03:00.7 | N1 | P1@5, o7 v4 | Exactly next: records email key o7:4, emails "shipped", appends c1#3 with src = 107's event_id, then killed | Feed last_seq still 2; last_version[o7] still 3; P1's committed offset still 4 |
| 10:03:45.7 | Group | 45 s after N1's last heartbeat, P1 moves to N2 | ||
| 10:03:45.7 | N2 | P1@4, o7 v3 | 3 ≤ 3: skip | |
| 10:03:45.7 | N2 | P1@5, o7 v4 | Exactly next. Email key o7:4 exists: no email. Appending c1#3 fails "already exists"; N2 reads it: its src is this event's event_id, so this is its own group's earlier attempt: it only bumps last_seq to 3. Stores last_version[o7] = 4, commits offset 6 | c1#3 once, one "shipped" email |
The feed is a derived numbered record: entry n exists only after n − 1. A naive new owner that takes "the next free number" would write c1#4 for the same event, and the customer's feed would show "shipped" twice. The src check tells "my own replay" apart from "someone else's entry":
textAPPEND DERIVED ENTRY (customer c, event e) n = last_seq[c] + 1 write entry c#n with src = e.event_id, only if c#n doesn't exist if it existed: if entry(c#n).src == e.event_id: last_seq[c] = n -- my own replay: finish the bookkeeping else: -- someone else took n check I still own the partition (group generation or lease: page 03); if not, stop re-read last_seq[c]; write at the first free number after n; never retry n else: last_seq[c] = n last_seq never expires
Snapshot S3 (after event 14):
| Lane | State |
|---|---|
| Database | o7 v4 SHIPPED, o8 v2 PAID, o9 v1 PLACED. Outbox empty |
| Mover | R1 and R2 running; at 10:03:00.6 R2's claim skipped bucket 3 (R1 held it) and found no other bucket with rows |
| Topic | P0: @0 o8 v1, @1 o8 v1 (repeat), @2 o8 v2. P1: @0 o7 v1, @1 o7 v2, @2 o7 v2 (repeat), @3 o9 v1, @4 o7 v3, @5 o7 v4 |
| Consumers | All three at o7 → 4, o8 → 2, o9 → 1. N2 owns P0 and P1. Customer c1: c1#1 to c1#3, last_seq 3; "order received", "payment received" and "shipped" each sent once |
The cursor trap: numbers are not commit order
A cheaper relay design keeps no claims and deletes nothing: it remembers "published up to seq N" and reads rows above N. Event 9 shows why that fails. Two transactions start at 10:02:00: T1 places o9, T2 pays o8.
Synthesizing vector architecture diagram...
What to notice: seq is handed out when a row is inserted, not when it commits. The cursor moved past 104 while 104 was still invisible, and no later read looks below the cursor.
| Time | What happens | Cursor |
|---|---|---|
| 10:02:00.000 | T1 inserts, gets seq 104 | 103 |
| 10:02:00.001 | T2 inserts, gets seq 105 | 103 |
| 10:02:00.002 | T2 commits | 103 |
| 10:02:00.003 | The cursor relay reads seq > 103: sees only 105, publishes it | 105 |
| 10:02:00.010 | T1 commits seq 104 | 105 |
| 10:02:01.003, 10:02:02.003, … | Reads seq > 105: nothing | 105 |
o9's "order placed" is never published, silently. Our claim-and-delete relay doesn't have this problem: it has no cursor, and any row still in the table is unpublished by definition, so at 10:02:00.6 it published both (105 to P0@2, 104 to P1@3). The other fixes are to read in commit order, which is what a log reader does (Part 5), or to read only up to the oldest transaction still in flight. The payment loop (step 3.6) has its sealer read the change stream in commit order for exactly this reason: "entry IDs can commit out of order, so batching by ID could skip a late-committing entry".
A related rule for the relay loop: keep claiming until a claim comes back empty, and treat an empty claim as "nothing visible right now", not as "done forever". A row that commits a moment later is still in the table at the next tick.
Waking the relay
| How the relay wakes | Delay after a commit | What it costs |
|---|---|---|
| Poll every 1 s | Up to 1 s | One claim query per relay per second, forever, even when idle |
| An in-process signal after each commit, plus a 1 s sweep | About the time to claim and publish | Works only for commits made by the same service; the sweep still runs to catch missed signals |
PostgreSQL LISTEN/NOTIFY | About the same | A notification reaches only sessions listening at that moment, so a restarting relay misses some and the sweep stays; payloads must be shorter than 8,000 bytes by default |
| Read the database's log (Part 5) | Decode and publish time | No queries against the outbox; a connector and a replication slot to run and watch |
When the broker is down
The outbox absorbs a broker outage: rows simply wait. The outbox loop (step 2.4) sizes it: MSK is down for 10 minutes while 6,000 events a second keep arriving, so 6,000 × 600 = 3,600,000 rows pile up. When MSK returns, the relays drain at most 20,000 events a second (a cap, because every claim, read and delete runs on the same writer as live transactions: a shared limit), while 6,000 a second keep arriving:
drain time = backlog ÷ (drain rate − arrival rate) = 3,600,000 ÷ (20,000 − 6,000) = 3,600,000 ÷ 14,000 ≈ 257 s, about 4.3 minutes.
Each bucket drains oldest first, so the oldest events are published within seconds of the broker's return and the newest backlog events wait the longest.
Delete churn. Every event is an insert and then a delete: at 6,000 events a second that is 6,000 dead rows a second for vacuum to clean. If anything holds vacuum back, the relay's claim and read queries slow down. The outbox loop (steps 1.6 and 2.5) handles it with short transactions and partitioned outboxes whose empty partitions are dropped; that detail stays there.
What to remember from Part 3
- A crash between publish and mark gives duplicates, so delivery is at least once.
- Consumers apply by version, never by adding: exactly next, skip the old, park the new.
- A reader that remembers "up to
seqN" misses rows that commit out of number order.
Part 4. Order per aggregate
Consumers need one order's events in version order: "shipped" before "address changed" makes the cache wrong and the customer confused. Order is not a property of Kafka or of the database; it holds only if every hop keeps it. This Part breaks it one hop at a time.
Trace: two relays split one order
At 10:02:30 a second relay, R2, starts beside R1. At 10:03:00.0 o7's address changes (v3, seq 106) and at 10:03:00.2 it ships (v4, seq 107). Suppose the relays claimed rows instead of buckets: at 10:03:00.6 R1 claims row 106 and R2 claims row 107, and R2 happens to publish first. From the reference implementation:
| Record | Indexer (last 2) | Cache updater without the version check |
|---|---|---|
P1@4 = o7 v4 (from R2) | Gap after 2: park v4 | Applies: SHIPPED@4 |
P1@5 = o7 v3 (from R1) | Exactly next: apply v3, then unpark v4 and apply it | Applies: PAID@3 |
| Final | SHIPPED v4, correct | PAID v3, while the database says SHIPPED v4 |
The version check saved the indexer; the cache without it is now as wrong as the dual write in Part 1. The real design never lets this happen in the first place: relays claim whole buckets, so one order's rows have one publisher at a time. In the main line, R1 claimed bucket 3 at 10:03:00.6, R2's claim skipped it and found no other bucket with rows, and R1 published 106 to P1@4 then 107 to P1@5.
Every hop must keep order
Synthesizing vector architecture diagram...
What to notice: the first four links keep order; the last one detects when any of them didn't. Remove any link and order can break; remove the check and nobody notices.
| Hop | What keeps o7's order | In events 10 to 12 |
|---|---|---|
| Database | The row lock (Part 2) | v3 commits at 10:03:00.0, v4 at 10:03:00.2 |
| Relay | One publisher per bucket, publishing in seq order | R1 holds bucket 3; R2 skips it |
| Topic | Key = order_id: one partition per order, delivered in offset order | P1@4 = v3, P1@5 = v4 |
| Consumer group | One consumer per partition at a time | The indexer reads P1 in offset order |
| Consumer | The exactly-next check | v3 applied at last 2, v4 at last 3 |
The cost: parallelism is capped by the bucket count and the partition count, and a very busy order delays the other orders in its bucket.
Gaps are checked number by number
A consumer that reads a batch must check every version, not the first and last of the batch. The mobile news feed loop (step 2.4) states it for feed numbers: every number from last + 1 onward must be present, with no hole anywhere in the batch.
textAPPLY BATCH (records of one partition, in offset order) for each record r (key K, version v): if K is blocked: park r; continue -- Part 6 if v <= last_version[K]: skip; continue if v != last_version[K] + 1: park r; continue -- a hole before v apply r and set last_version[K] = v in one transaction while parked[K] holds last_version[K] + 1: apply it and store its version -- the hole just closed commit the partition's offset after the batch
Adding partitions
More throughput needs more partitions. The trap: the partition is chosen when a record is produced. Our example grows orders.events from 2 to 4 partitions while R1 is publishing o7 v3 and v4 (a side trace; the main line keeps 2 partitions). The example forces R1's producer to refresh its metadata between the two sends; real producers pick up new partitions at their next metadata refresh (in Kafka, metadata.max.age.ms, 5 minutes by default, or sooner after an error).
Synthesizing vector architecture diagram...
What to notice: one order's events now live in two partitions with no order between them. The newer event arrived 1.9 s before the older one, and only the version check put them back in order.
After the resize, o7 v4 arrived in P3 before v3 arrived in P1. What does each consumer do, and what would you change about how the resize is done?
Rebalances move partitions between consumers of a group, and the new owner replays from the last committed offset: event 14 in Part 3 was one. The version check and the src check make the replay harmless. DynamoDB Streams and Kinesis shards also split, and there a child shard must be read only after its parent (Part 9).
One transaction, several events
A transfer in the wallet loop (step 2.2) debits one account and credits another in one transaction, writing two outbox rows keyed by the two accounts. The keys hash to different partitions, so a consumer can see the debit without the credit for a while, even though they committed together.
| Option | How | Cost |
|---|---|---|
| One event for the whole transaction | One outbox row transfer.completed, keyed by the transfer, carrying both legs | Consumers that care about one account must unpack it, and per-account order has to come from somewhere else |
| One ordering scope | Key both rows by a scope that contains both accounts | Wider scopes cut parallelism; a busy scope becomes a hot partition |
| A transaction id and a count | Each event carries the transaction's id and how many events it has; a consumer that needs both waits for both. Debezium can emit transaction boundary events with an event count (provide.transaction.metadata, off by default) | Buffering, and an alarm when a transaction stays incomplete |
The rule: choose the aggregate as the smallest scope whose events must be seen together. The outbox loop (step 1.5) keys ledger events by account; the payment loop (step 3.6) reads the change stream in commit order so a transaction's entries are never split by a batch boundary.
Epochs: when the writer itself changes
After a failover promotes a new database, the version numbers of an order can repeat for different events (the lost tail of the old writer). The outbox loop (step 3.4) copies a writer_epoch into every event and has consumers compare (epoch, version):
| Consumer sees | Meaning | Action |
|---|---|---|
| Same epoch, version = last + 1 | Normal | Apply |
| Same epoch, version ≤ last | Duplicate | Skip |
| Higher epoch, version = last + 1 | Failover with nothing lost for this order | Apply; store the new epoch |
| Higher epoch, version ≤ last | A fork: we applied events the new writer never had | Park the order; reconcile |
| Lower epoch than stored | A write from the fenced old writer | Park it; never apply silently |
| Version > last + 1 | Missing events | Park behind the gap |
How the epoch is bumped inside the promoted database before writes open is in Leases, Fencing Tokens & Distributed Locks (Part 7).
What to remember from Part 4
- Order holds per key, only while every hop keeps it: one publisher, one partition, one consumer at a time.
- Adding partitions moves keys at produce time: the version check is what saves you.
- A gap means park and wait, never skip or apply.
Part 5. Reading the database's own log
By 11:00 the relays work, but look at what they cost at scale: every relay runs a claim query every second against the busiest database in the company, every event is an insert and a delete for vacuum to clean, every event waits up to a poll interval, and every team builds its own relay. The outbox loop (step 3.2) hits this at 2 billion events a day. The database already writes every committed change to its write-ahead log (WAL) for durability; a reader can follow that log instead of querying tables. That is log-based change data capture (CDC).
Why read the log
| Polling relay (every 500 ms) | Log tailing (CDC) | |
|---|---|---|
| Load on the database | A claim query every 500 ms per relay, forever, plus a read and a delete per batch; dead rows for vacuum | Reads WAL the database writes anyway; decoding runs on the writer; almost no query load |
| Delay | Up to the poll interval, unless signalled | As soon as the commit is in the log, plus decode and publish time |
| Order | Per bucket, only with whole-bucket claims; a cursor design misses late commits (Part 3) | Commit order, from one decoder |
| Duplicates | A crash between publish and delete | A crash between publish and the slot's confirmed position; the position is saved only at checkpoints |
| What you run | The relay, inside the service | A connector platform and a slot per database |
| How it breaks | Table bloat, slow claims, vacuum falling behind | A stalled slot pins WAL on the writer; a cap turns that into a gap (Part 7) |
| Throughput ceiling | Relays scale with buckets, but share the writer | One task per connector: its rate is the database's CDC ceiling (the outbox loop plans on about 25,000 small events a second per connector, a figure it load-tests, and splits databases beyond that) |
This answers the second question of the drill The Dual-Write That Broke Search Consistency: polling every 500 ms adds a query and table churn to the busiest database forever, and up to half a second of delay; tailing the WAL costs almost no queries and keeps commit order, and in exchange you run a connector and watch a slot. At one service and modest rates the relay is the better deal; at many services or high rates, CDC is.
The outbox table or the business tables?
CDC can read the outbox table (the same events as before, now moved by the log) or the business tables themselves.
| CDC on the outbox table | CDC on business tables | |
|---|---|---|
| Contract | Events the producer designed; stable while the tables change | The table's columns as they are today; a column rename is a contract change |
| Meaning | order.paid | "a row changed; status is now PAID": consumers infer the meaning |
| Version | aggregate_version in the envelope | Only if the table has a version column |
| Extra write | One insert per change (and its WAL) | None |
| A change touching three tables | One event can describe it | Three row events, one per table |
| Deletes | An explicit event with a version (Part 8) | A delete event, plus a tombstone record for compacted topics |
From 11:00 our story tails the outbox table: consumers keep the same events and versions, and only the mover changes.
PostgreSQL: the replication slot
Logical decoding needs wal_level = logical. A publication lists the tables to capture:
sqlCREATE PUBLICATION outbox_pub FOR TABLE outbox, cdc_heartbeat;
A logical replication slot records how far its reader has confirmed. The server keeps every WAL segment from that position on, whether or not the reader is connected. The Write-Ahead Log page (Part 8) gives the general rule: a log segment may go only when its data is durable elsewhere and no reader still needs it, so the slowest reader sets the log's size. This page owns what happens next.
Synthesizing vector architecture diagram...
What to notice: the connector moves the slot forward only after the broker acknowledged. Everything between the slot's position and now stays on the writer's disk, up to the cap, if one is set.
| Slot fact | What it means for us |
|---|---|
| A slot is crash-safe and keeps WAL (and the catalog rows needed to decode it) even with no connection | A dead connector doesn't lose its place; it pins the writer's disk instead |
| "A logical slot will emit each change just once in normal operation"; its position is persisted only at checkpoints | After a database crash the slot can return to an earlier position and resend recent changes: at least once |
| Concurrent transactions are decoded in commit order, each whole | No cursor trap (Part 3): a late-committing row comes out after the rows that committed before it |
max_slot_wal_keep_size (default −1: no limit) caps the WAL a slot may hold | Past it, the slot's wal_status becomes unreserved and then lost: the reader can't continue. A disk risk becomes a guaranteed gap (Part 7) |
| Heartbeats | If the captured tables are quiet but others are busy, the connector has nothing to confirm and the slot never moves. Debezium writes a heartbeat row on a timer (heartbeat.interval.ms with heartbeat.action.query) into a captured table, giving it something to confirm |
catalog_xmin | An idle slot also holds back cleanup of the system catalogs, not only WAL |
REPLICA IDENTITY | Decides which old values an update or delete carries; FULL sends the whole old row, needed for before images on business tables, not for an insert-only outbox |
| Failover slots (text only) | A new primary lacks the slot unless failover slots are set up (the slot created with failover, synchronised to the standby); otherwise a failover is a lost position (Part 7) |
| One task, on the writer | Debezium's PostgreSQL connector runs a single task; Aurora PostgreSQL doesn't support logical decoding from readers, so it runs against the writer |
To watch it:
sqlSELECT slot_name, active, wal_status, pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) AS retained_bytes FROM pg_replication_slots;
On Aurora and RDS the same numbers arrive as CloudWatch metrics: OldestReplicationSlotLag, ReplicationSlotDiskUsage and TransactionLogsDiskUsage.
The outbox is quiet, but the orders table is busy with updates that write no outbox rows (say, a batch job re-pricing items). Why does the slot's retained WAL still grow?
When the slot stalls
A stalled reader keeps WAL growing at the database's WAL rate:
retained WAL = WAL rate × stall time
At real scale, from the outbox loop: about 2 KB of WAL per event × 12,000 events a second at peak (its CDC round) = 24,000 KB/s = 24 MB/s. A cap of 200 GB then lasts 200,000 MB ÷ 24 MB/s ≈ 8,333 s ≈ 2.3 hours of stall at peak. In our tiny example we simply set the cap so it is reached 6 hours into a stall.
Two shared limits decide where the cap goes:
- Every holder counts. Retained WAL runs from the oldest position anything needs: every logical slot, every physical replica's slot,
wal_keep_size, an archiver that is behind. The slowest one decides. - The tables share the disk. They keep growing through the stall. Set the cap below the free disk minus the tables' growth over the longest stall you plan to ride out; the outbox loop pairs a 200 GB cap with at least 400 GB free.
Before you set a cap, plan the republish (Part 7): past the cap, the writer is safe and the reader has a guaranteed gap. Check that your engine version's parameter group exposes max_slot_wal_keep_size (it appears in Aurora PostgreSQL's parameter list; the outbox loop checks per engine version).
MySQL: the binlog is purged by time (text only)
MySQL's change log is the binary log in row format (binlog_format=ROW, the default). A reader stores its own position (a GTID set, or a file and an offset). The server doesn't keep binlogs for a reader: they expire by age (binlog_expire_logs_seconds, 30 days by default; on RDS for MySQL the binlog retention hours setting, which is NULL by default, meaning binlogs are not retained, with a maximum of 168 hours). The same stall ends the opposite way:
| Same stall | PostgreSQL logical slot | MySQL binlog reader |
|---|---|---|
| Who keeps the log | The server, until the slot confirms (up to the cap) | The server, for its retention period only |
| Before the limit | The writer's disk fills | Nothing: the reader is just behind |
| After the limit | Slot invalidated: the reader can't continue; the writer is safe | Binlogs expire by age; only the file a connected reader has open (and later ones) is spared, so a stopped or disconnected reader loses its place; the writer is safe |
| Who is harmed with no limit set | The writer: a full disk stops the database | Nobody's disk; the reader still loses its place at the retention time |
| Recovery | New slot, republish, resume (Part 7) | Re-snapshot and republish, then resume from a new position |
Other databases expose the same idea under other names (SQL Server's CDC, Oracle's redo mining, MongoDB change streams); no loop on this site uses them.
Cutting over from the relay (event 19)
The order of the cut-over matters: open the new reader before the old one stops.
Synthesizing vector architecture diagram...
What to notice: from 11:00:00 every commit is covered by the slot. Had the relays stopped first, a row committed between their stop and the slot's creation would reach neither path; with the slot first, a row committed during the overlap is published by both, a harmless duplicate.
no_data means the connector takes no snapshot and streams from the slot's creation point; the relays had already confirmed every earlier row, so nothing is missing. From now on the outbox is insert-only: the connector never deletes, and old rows go by dropping daily partitions after 3 days. Keeping them is what makes Part 7's republish possible.
End to end through the layers
Event 20 (o8 v3, seq 111) through every layer:
Synthesizing vector architecture diagram...
What to notice: two positions are saved, the slot's (step 4) and the consumer's offset (step 7), and each is saved after the work it covers. That is why a crash repeats work and never skips it.
| Crash after step | What survives | What happens next | Result |
|---|---|---|---|
| 1, before commit | Nothing | The client retries (page 04) | Nothing lost |
| 2 | The commit, in the WAL | The connector reads it when it comes back | Nothing lost |
| 4, before the slot position is saved | The record in the broker | The connector restarts from the slot's position and publishes it again | A duplicate |
| Database crash with the slot's position unsaved | The WAL | The slot returns to its last checkpointed position and resends | A duplicate |
| 5 | An acknowledged record, if acks=all really waited for enough replicas (the WAL page, Part 10, explains acks=all with min.insync.replicas) | Nothing lost | |
| 6, before the consumer's transaction commits | The record and the old offset | The consumer reads it again and applies it | Nothing lost |
| 6, after its transaction commits, before step 7 | The effect and last_version | The consumer reads it again; version ≤ last: skip | A skipped duplicate |
Gaps come only from losing a position: a slot invalidated or dropped, a failover to a primary without the slot, a connector restarted from the wrong offset, a consumer offset reset to "latest", or records trimmed by retention. Part 7 handles each.
Lag is a sum of the dependent hops: commit, decode, publish and acknowledge, the consumer's poll, the apply. In the tiny example each hop takes 0.1 s: o8 v3 commits at 11:10:00.0, is published at 11:10:00.1 and applied at 11:10:00.2. At real scale, measure it end to end with the heartbeat.
Each hop is at least once; exactly-once per hop, and Kafka transactions, are covered by the Idempotency & Effectively-Once Processing loop primitive (page 04).
What to remember from Part 5
- A change stream reads the commit log, so it adds almost no query load and keeps commit order.
- A cap turns a full disk into a guaranteed gap: plan the republish before you set it.
- A MySQL binlog is purged by time: a stalled reader loses its place instead of filling the disk.
Part 6. Poison records and repair
At 10:04:30.0 a producer bug stores an address for o9 that the indexer's parser can't handle. The event is valid as far as the outbox knows, it is published like any other, and the indexer fails on it every time. A record that always fails is called a poison record. Retrying it forever blocks everything behind it on the partition; catching the error and moving on loses the change silently. This Part does neither.
Trace: o9's bad address
At 10:04:30.0 o9 v2 order.paid (seq 108) commits with the bad address; at 10:04:30.2 o7 v5 (seq 109) commits. R1 publishes both at 10:04:30.6, to P1@6 and P1@7, so o7 v5 sits right behind the poison.
Synthesizing vector architecture diagram...
What to notice: the partition waited 6.1 s (three retries, after 0.1 s, 1 s and 5 s), not forever. After that only o9 is stopped, for this consumer; o7 flows past it.
| Time | Event | Indexer | Cache updater | Notifier |
|---|---|---|---|---|
| 10:04:30.0 | o9 v2 (seq 108) with the bad address | |||
| 10:04:30.2 | o7 v5 (seq 109) | |||
| 10:04:30.6 | 108 → P1@6, 109 → P1@7 | |||
| 10:04:30.7 | o9 v2 fails; retries at 30.8 and 31.8 fail | Applies o9@2 and o7@5 | o9@2: emails "payment received"; o7@5: no email | |
| 10:04:36.8 | Fourth attempt fails: dead-letter o9 v2, block o9, then apply o7@5 | |||
| 10:20:00 | Producer fix: o9 v3 (seq 110) rewrites the address → P1@8 at 10:20:00.6 | o9 is blocked: park v3 | Applies o9@3 | Applies o9@3, no email |
| 10:20:30 | Operator resolves the dead letter as "superseded" | Repair worker re-reads o9 on the primary: v3. Indexes v3, unblocks o9; parked v3 ≤ 3: skip |
Only the indexer's parser chokes on the address, so only the indexer dead-letters. The error boundary is per consumer.
Snapshot S4 (after event 17):
| Lane | State |
|---|---|
| Database | o7 v5 SHIPPED, o8 v2 PAID, o9 v2 PAID (bad address). Outbox empty |
| Mover | R1 and R2 running |
| Topic | P1 has grown to @7: … @6 o9 v2 (seq 108), @7 o7 v5 (seq 109). P0 unchanged (@0 to @2) |
| Consumers | Indexer: o7 → 5, o8 → 2, o9 → 1, o9 blocked, one dead letter (o9 v2, "address parser error"). Cache and notifier: o7 → 5, o8 → 2, o9 → 2 |
The error boundary
Engine-neutral, per consumer:
- Validate the record against its schema before doing anything with it.
- Retry transient failures a few times with backoff (0.1 s, 1 s, 5 s here). A schema violation isn't transient and can skip the retries.
- Dead-letter it: in one transaction in the consumer's own database, store the record and the reason, mark the key blocked for this consumer, and commit the offset. Recording it in the same database as the offset decision means a crash can't lose it.
- Park the key's later events behind it, in version order, instead of applying them.
- Move on: every other key on the partition keeps flowing.
- Alarm on every dead letter; each needs an owner.
The indexer has dead-lettered o9 v2. How does it get o9 right again without ever applying v2?
textREPAIR (consumer, key K) -- after the cause is fixed row = read K from the PRIMARY -- never a lagging replica (event 28, Part 8) in one transaction in the consumer's database: write the effect from row's current state if row.version > last_version[K] last_version[K] = row.version mark K's dead letters resolved ("superseded by row.version") unblock K drop parked events for K with version <= row.version apply any parked event that is now exactly the next
Replays and a quarantine list. A rebuild that replays raw events from an archive would feed o9 v2 to the new code and fail again, or worse, succeed on a payload that was rejected on purpose. So replays skip a quarantine list of events resolved as invalid (the outbox loop, step 3.3), or they rebuild from the source of truth instead (Part 7).
The periodic compare. Every stream-fed copy also gets a job that compares it with the source and writes repairs: the stream is how changes normally arrive; the compare catches what slipped through anyway (a bug, a lost position nobody noticed). Its reads must come from the primary, or from a replica known to be caught up; Part 8's event 28 shows the compare job going wrong with a lagging replica.
Stop it at the producer. The best poison record is one that never commits. Validate each event against its registered schema and the size limit before the outbox insert, so a bad event fails the business transaction instead of reaching the stream (the outbox loop, step 2.3).
Managed stream readers have their own knobs for all of this (Lambda's and Pipes' bisect, maximum record age and on-failure destination); they are in Part 9.
What to remember from Part 6
- A failing record blocks its partition until you move it aside: dead-letter it and block only its key.
- Repair re-reads the source of truth; a dead letter is a pointer, not the data to replay.
- Every stream-fed copy needs a repair path and a periodic compare.
Part 7. Refilling what the stream can't give you: republish and bootstrap
At 18:00:00 slot S1 is invalidated. Two committed events, seq 112 and 113, will never come through it. No consumer can ask for them, because as far as the stream is concerned they never happened. This Part is about gaps: where they come from, how to notice them, and how to fill them in the right order.
Trace: the slot is invalidated
| Time (D1) | Event |
|---|---|
| 11:00:10 onward | The connector writes a heartbeat row every 10 s and confirms it |
| 12:00:00 | The connector dies right after its 12:00:00 heartbeat went through. S1's confirmed position is that heartbeat's commit (just after seq 111's) |
| 12:01:00 | End-to-end lag (now minus the commit time of the newest heartbeat that came through) passes 60 s |
| 12:06:00 | Lag has been above 60 s for 5 minutes: the alarm fires. The story assumes nobody acts |
| 13:00:00 | o8 v4 order.shipped (seq 112) commits; not published |
| 15:00:00 | o9 v4 order.shipped (seq 113) commits; not published |
| 18:00:00 | Retained WAL passes the cap, 6 hours after the slot's confirmed position: S1 is invalidated and its WAL removed |
Snapshot S5 (after event 23):
| Lane | State |
|---|---|
| Database | o7 v5 SHIPPED, o8 v4 SHIPPED, o9 v4 SHIPPED. Outbox (insert-only): 111 (o8 v3), 112 (o8 v4), 113 (o9 v4): in the outbox only |
| Mover | Connector dead; slot S1 invalidated (wal_status = lost) |
| Topic | P0 ends at @3 (o8 v3, seq 111); P1 ends at @8 (o9 v3, seq 110) |
| Consumers | All three at o7 → 5, o8 → 3, o9 → 3 |
Where gaps come from
| Source | Example | What is missing |
|---|---|---|
| A slot invalidated by its cap, or dropped | Event 23 | Every change from the slot's position until a new position exists |
| A failover to a primary that lacks the slot (no failover slots) | The outbox loop, step 3.4 | Changes between the old slot's position and the new slot's creation. Same procedure as below: create the slot before writes open on the new primary, republish, then start the connector |
| A consumer down longer than retention | Event 25 | Everything trimmed while it was away |
| A cursor that skipped a row | Part 3 | The late committer |
| A connector restarted from the wrong position | Offsets reset or copied from another environment | Anything in between |
Detecting a gap
textGAP CHECK (every minute) alarm if heartbeat lag > 60 s for 5 minutes alarm if any consumer has held a parked event longer than a threshold t = newest seq the topic holds from this database -- read the topic first o = newest committed seq in the outbox -- then the outbox if o > t: wait briefly, then read both again once -- an in-flight publish is not a gap if o is still > t and the connector's position isn't moving: declare a gap
The re-check matters: an empty or short read that races a writer isn't evidence of anything. The mobile news feed loop (step 2.4) re-checks an empty read the same way before declaring a gap.
Republish: refill from the stored source, then resume
textREPUBLISH (after a lost position) 1. create the new slot FIRST -- every later commit is covered from here 2. newest = the newest event the topic holds from this database 3. from = newest.committed_at - margin -- can reach back only as far as rows are kept 4. publish every outbox row committed at or after `from`, in seq order, key = aggregate_id 5. wait for every acknowledgement 6. start the connector on the new slot -- it re-sends what committed after step 1 7. consumers skip what they already have
Event 24, from the reference implementation:
Synthesizing vector architecture diagram...
What to notice: the newest event the topic held was seq 111 (committed 11:10:00); a 10-minute margin reaches back to 11:00:00, the cut-over, so rows 111 to 114 are republished. Rows before 11:00 were deleted by the relay and aren't needed: the relay confirmed every one of them. And 114 comes twice, once from the republish and once from the connector.
That margin works only because the outbox has been insert-only since the cut-over. A relay that deletes on publish keeps nothing to republish from, and rows are kept only 3 days: the outbox's retention must be longer than your worst time to detect and repair a gap.
Snapshot S6 (after event 24):
| Lane | State |
|---|---|
| Database | o7 v5, o8 v5 DELIVERED, o9 v4. Outbox: 111 to 114. Slot S2 active |
| Mover | Connector on S2 |
| Topic | P0: … @3 o8 v3 (111), @4 o8 v3 (111 again), @5 o8 v4 (112), @6 o8 v5 (114), @7 o8 v5 (114 again). P1: … @9 o9 v4 (113) |
| Consumers | From 19:00:05.1 the indexer skips P0@4, applies o8 v4 and v5, skips P0@7 (written at 19:00:05.2), applies o9 v4; so do the others. All three at o7 → 5, o8 → 5, o9 → 4. The notifier emails o8 "shipped", then "delivered", and o9 "shipped" |
Why republish before starting the connector on the new slot, if consumers skip duplicates anyway?
Trace: the index is down for 30 hours
| Time | Event |
|---|---|
| D2 00:00 | The indexer goes down. Its index is intact, but stops changing |
| D2 03:00:00 | o9 v5 order.delivered (seq 115) → P1@10; the cache and the notifier apply it |
| D3 06:00:00 | The indexer returns. Retention (24 hours) has removed every record older than D2 06:00:00, P1@10 included. Its committed P1 offset (10) is below the partition's first remaining offset (11): it can't resume. It must rebuild |
A consumer that was away longer than retention can never catch up from the stream alone. It needs a snapshot of the source plus the stream, joined at the right point:
Synthesizing vector architecture diagram...
What to notice: in the right order, every change after the noted position is in the stream the rebuild applies, so v6 is seen twice and applied once. In the wrong order, v6 committed after the scan read o8 and before the position was noted, so it is in neither.
| Time | Step (right order) | Result |
|---|---|---|
| 06:00:00 | Note the end offsets: P0 8, P1 11 | |
| 06:00:01 | The scan reads o7 on the primary | v5 SHIPPED |
| 06:00:05 | o8 v6 order.returned (seq 116) commits → P0@8 at 06:00:05.1 | |
| 06:00:10 | The scan reads o8 | v6 RETURNED |
| 06:00:19 | The scan reads o9 | v5 DELIVERED |
| 06:00:20 | Apply the stream from P0 8, P1 11 | P0@8 o8 v6: 6 ≤ 6, skip |
Snapshot S7 (after event 26):
| Lane | State |
|---|---|
| Database | o7 v5 SHIPPED, o8 v6 RETURNED, o9 v5 DELIVERED. Outbox: 111 to 116 |
| Mover | Connector on S2 |
| Topic | P0: only @8 o8 v6 (seq 116) remains; older records trimmed. P1: empty, next offset 11 |
| Consumers | Rebuilt indexer: o7 → 5, o8 → 6, o9 → 5, documents SHIPPED, RETURNED, DELIVERED. Cache and notifier at the same versions |
Scan first and note the position after, or the other way round?
A replica is a valid snapshot source only once its replay position is at or past the noted stream position. A replica that lags reads the past, and the stream from the noted position won't contain what it's missing. The same mistake returns in event 28 (Part 8), where the compare job reads a replica 40 s behind. The ad-click loop (step 3.1) and the mobile news feed loop (step 2.3: build the head after reading last_seq) follow the same order.
The tools the source gives you
Debezium's PostgreSQL connector has snapshot modes that decide what it does when it starts:
| Mode | What it does |
|---|---|
initial (the default) | Snapshot the captured tables, then stream |
no_data | No snapshot: stream from the stored position, or, if none is stored, from the point the slot was created (our event 19) |
when_needed | Snapshot when there is no stored position, or when the stored position is no longer available on the server |
when_needed quietly re-snapshots after a lost position: every row the outbox still holds is read again, the same effect as a republish but unannounced. Know which mode you run before an incident, not during one.
Shadow rebuild and switch. When the live copy must keep serving while it is rebuilt, build a shadow copy (a new index) by the same procedure, let it catch up to the live position, then switch readers to it and drop the old one. The outbox loop (step 3.3) rebuilds a projection this way, skipping its quarantine list.
What to remember from Part 7
- A lost position is a guaranteed gap: open the new position, refill from the stored source in order, then resume.
- Note the position first, then snapshot, then apply the stream keeping only newer versions.
- Keep the outbox longer than your worst detection-plus-repair time.
Part 8. Schema changes and deletes
Events outlive the code that wrote them. At 10:04:30.2 o7 v5 arrives with schema_version 2 and a new field, currency, and two of the three consumers are still running last week's code. Later, o7 is erased, and a repair job that doesn't know it tries to bring it back. Both are the same problem: what a consumer does with an event it didn't expect.
Trace: currencymeets old code
o7 v5 order.items_updated (seq 109) adds an optional field, currency. Event 17, per consumer:
| Consumer | What it does with o7 v5 |
|---|---|
| Indexer | Indexes v5 with currency (at 10:04:36.8, behind the poison record of Part 6) |
| Cache updater (old code) | Ignores the unknown field, writes status and total: SHIPPED@5 |
| Notifier (old code) | items_updated is not an email type: advances o7 to 5, no email |
Nothing broke, because adding an optional field is a safe change. Most changes are not.
Compatibility in plain words
| Word | What it means | Why we need it |
|---|---|---|
| Backward compatible | New consumer code can read old events | Replays and rebuilds feed old events to new code |
| Forward compatible | Old consumer code can read new events | Producers usually deploy before every consumer has |
| Full | Both | Either side may deploy first |
| Transitive | Checked against every earlier version, not just the previous one (AWS Glue Schema Registry's BACKWARD_ALL, FORWARD_ALL, FULL_ALL) | A replay may feed events from any version ever published |
| Change | Safe? | Why |
|---|---|---|
| Add an optional field (with a default) | Safe | Old code ignores it; new code uses the default on old events |
| Add a required field | Breaks backward | Old events don't have it |
| Remove a field old code requires | Breaks forward | Old code finds it missing |
| Rename a field | Breaks both | A rename is a remove plus an add |
| Change a field's type or unit (units to cents) | Breaks, even when a registry can't tell | The same name now means something else |
The tools: a schema registry keeps each event type's versions and a compatibility mode; CI rejects an incompatible change before it ships, and producers can publish only registered versions. With change data capture on business tables (Part 5), the table is the schema, so a column rename is a contract change for every consumer.
The producer wants to rename total to total_minor (the amount in cents). What do you do, step by step?
textDECODE (event) -- fail closed reader = the registered reader for (event_type, schema_version) if none: dead-letter "unknown schema version"; block the key fields = reader.decode(payload) if a required field is missing or has the wrong type: dead-letter with the reason; block the key -- never default to 0 ignore unknown optional fields
A delete must travel
A deleted order leaves nothing behind to publish. "Consumers will notice it's gone" is wrong: a consumer sees events, not absences. So a delete travels as an event with a version, like any change:
| Source | How the delete travels |
|---|---|
| Outbox | An explicit event, order.erased with the next version, written in the same transaction that deletes the row. (Deleting outbox rows is not an event: Debezium's outbox router filters out DELETEs on the outbox table) |
| CDC on a business table | A delete event, followed by a tombstone record (a null value for the key) so compacted topics can drop the key; Debezium does this by default (tombstones.on.delete is true) |
| DynamoDB Streams | A REMOVE record, including deletes made by TTL, which carry a userIdentity of type Service with principal dynamodb.amazonaws.com (Part 9) |
Trace: o7is erased, and a repair comes too late
| Time (D3) | Event | Indexer | Cache |
|---|---|---|---|
| 06:59:50 | The Part 6 compare job reads o7 from a read replica lagging 40 s: it sees the state as of 06:59:10, v5 SHIPPED | ||
| 07:00:00 | o7 erased: one transaction deletes the orders row and inserts order.erased v6 (seq 117) → P1@11 at 07:00:00.1 | ||
| 07:00:00.2 | Deletes the document; keeps a tombstone (o7, 6) for 7 days | Writes order:o7 = GONE@6 | |
| 07:00:20 | The compare job reads the index: o7 missing | ||
| 07:00:30 | The compare job writes its "repair": o7 v5 SHIPPED | Tombstone 6 ≥ 5: dropped | SET NX finds GONE@6: no-op |
Snapshot S8 (after event 28):
| Lane | State |
|---|---|
| Database | o8 v6 RETURNED, o9 v5 DELIVERED; o7 gone. Outbox: 111 to 117 |
| Mover | Connector on S2 |
| Topic | P0: @8 o8 v6. P1: @11 o7 v6 order.erased |
| Consumers | Indexer: no o7 document, tombstone (o7, 6); o8 → 6, o9 → 5. Cache: order:o7 = GONE@6. Notifier: o7 → 6, no email for an erasure. The repair write was dropped everywhere |
At 07:00:30 a repair job writes o7 v5, which it read from a lagging replica. What stops it re-creating the document?
Caches store "gone". Deleting order:o7 would leave room for a late refill to put v5 back; storing GONE@6 means refills and repairs, which write with SET … NX (only if absent), can't. Compacted topics (Kafka's cleanup.policy=compact, which keeps the latest value per key) keep a delete's tombstone for delete.retention.ms, one day by default, then drop it: a consumer away longer than that never sees the delete and must rebuild (Part 7). Erasure events carry no personal data: order.erased holds the id and the version, nothing else, so the log isn't where the data lives on.
What to remember from Part 8
- Add optional fields; a breaking change is a new event type, run alongside the old.
- A delete is an event with a version, not an absence: remember it longer than any late write can arrive.
- Caches store "gone", so a late refill or repair can't bring the value back.
Part 9. The same story on DynamoDB
Our orders table now lives in DynamoDB, and the search index is fed from its stream by a Lambda function. The indexer was down for 30 hours. From where can it resume, and how do you rebuild it without losing a change?
This Part replays events 15 to 17 and 25 to 26, plus one shard split, on DynamoDB. The replay keeps the same orders and versions but its own times.
DynamoDB Streams in one table
| Fact | What it means for a consumer |
|---|---|
| Each change appears exactly once in the stream, in the order of that item's changes | Duplicates come from the reader's retries, not from the stream |
| Order is guaranteed per item only, "not across an entire partition", because one item collection (items sharing a partition key) can be split across partitions for heat | Two items with the same partition key can come out of order |
| Records are kept 24 hours, then may be trimmed at any time | A consumer away longer loses records |
Shard iterators: TRIM_HORIZON, LATEST, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER | There is no timestamp start |
| Each open shard belongs to one partition; shards split, and a parent must be read before its children | Readers must follow shard lineage |
| At most two readers per shard (one for global tables) | Every Lambda mapping, Pipe, OpenSearch Ingestion pipeline and KCL worker on the stream counts |
| A write that changes nothing writes no record | Versions move only on real changes |
TTL deletes appear as REMOVE records | Expiry looks like a delete to consumers |
| The stream is asynchronous and doesn't affect the table's performance; its view type can't be edited once enabled | Pick NEW_AND_OLD_IMAGES if consumers may need before images |
Trace: the poison batch, bisected
The table's stream uses NEW_AND_OLD_IMAGES. The Lambda mapping reads with batch size 4, bisect on (split a failing batch in two), 2 retries, a maximum record age of 3,600 s, and an SQS on-failure destination. Shard sh-1 holds four changes in a row (stream sequence numbers 100 to 103): o8 v2 (10:04:29.5), o9 v2 with the bad address (10:04:30.0), o7 v5 (10:04:30.2) and o8 v3 (10:04:30.4). The function applies each record by version and throws on the bad address. From the reference implementation (our model: split on each error; a single failing record is retried twice, then its metadata goes to the destination; AWS documents that splitting doesn't consume the retry quota):
Synthesizing vector architecture diagram...
What to notice: 7 invocations for 4 records. o8 v2 was delivered 3 times and o9 v2 5 times; the version check made every repeat harmless. The shard stopped for one bad record's retries, not for a day.
The destination receives metadata, not the record. For SQS and SNS destinations Lambda sends only a description of the failed batch; an S3 destination also gets the records. Trimmed to the fields that matter:
json{ "requestContext": { "condition": "RetryAttemptsExhausted" }, "DDBStreamBatchInfo": { "shardId": "sh-1", "startSequenceNumber": "101", "endSequenceNumber": "101", "batchSize": 1 } }
So the repair worker re-reads the item from the table, never trusts the message to hold it. At 10:04:40 it reads o9: the table itself still holds the bad address (v2), because the bug was in the producer, so the repair fails and stays open with an alarm. At 10:20:00 the producer's fix writes v3; the stream delivers it and the function parks it (its last o9 is 1, a gap at 2). At 10:20:30 the repair re-reads o9 (v3), writes the index and sets the version to 3; the parked v3 is skipped. The same repair as Part 6.
Trace: a shard splits
At 10:30 shard sh-1 is closed and split into two children. From the reference implementation:
Synthesizing vector architecture diagram...
What to notice: o7 v5 is in sh-1, and anything after it will be in sh-3. Reading sh-3 before sh-1 is finished could apply a newer change of o7 or o9 before an older one. DynamoDB's rule is that a parent is processed before its children; the DynamoDB Streams Kinesis Adapter for KCL follows it for you, and any reader you write must too.
The same rule has a sharp edge. With Lambda's default settings (no bisect, retry until the record expires), sh-1 would have been stuck on o9 v2 until D2 10:04:30, and a reader that keeps the parent-first rule would hold sh-2 and sh-3 back just as long. One bad record would have delayed every later change to these orders by up to a day.
Lambda and Pipes settings
| Setting | Lambda default | EventBridge Pipes default | Our replay |
|---|---|---|---|
| Batch size | 100 (up to 10,000) | 1 to 10,000 | 4 |
| Split a failing batch | BisectBatchOnFunctionError: off | OnPartialBatchItemFailure: set AUTOMATIC_BISECT to turn it on | On |
| Retries | MaximumRetryAttempts −1: until the record expires (up to 10,000) | −1: until the record expires | 2 |
| Maximum record age | MaximumRecordAgeInSeconds −1: until expiry (up to 604,800) | −1: infinite | 3,600 s |
| Where failed records go | On-failure destination: SQS or SNS (metadata only) or S3 (metadata and records) | DeadLetterConfig | SQS |
| Report per-record failures | ReportBatchItemFailures | partial batch responses | |
| Parallel batches per shard | ParallelizationFactor 1 to 10; order per item is kept | 1 to 10 | 1 |
| Starting position | TRIM_HORIZON or LATEST; LATEST can miss events while the mapping is being created or updated | The same, with the same warning | TRIM_HORIZON |
Two more facts belong with these. Lambda polls each shard 4 times a second and delivers at least once. Its concurrency for a stream is open shards × ParallelizationFactor, drawn from the account's Regional pool (1,000 concurrent executions by default), shared with every other function; throttling counts as an error before invocation and is retried until the record expires or passes its maximum age. Alarm on the mapping's IteratorAge metric (in the Lambda namespace), the Lambda-side equivalent of heartbeat lag. (Lambda can also aggregate over tumbling windows; windows belong to the Event Time, Watermarks & Checkpoints loop primitive.)
Rebuilding inside 24 hours
The 30-hour outage on DynamoDB, from the reference implementation:
| Time | Step | Result |
|---|---|---|
| D3 06:00:00 | The old mapping's saved position points into trimmed data: every record older than D2 06:00:00 is gone, o9 v5 (D2 03:00) among them | Can't resume |
| D3 06:00:00 | Create a new mapping at TRIM_HORIZON; request an export as of 06:00:00 | |
| D3 06:00:05 | o8 v6 changes; the mapping writes it at 06:00:05.25 (only if newer) | Index o8 = v6 |
| D3 06:10:30 | The export (assumed ready at 06:10) loads into the new (shadow) index, only if newer | o7 v5 written to the new index; o8 v5 dropped (6 is newer); o9 v5 written |
| after | The rebuilt index | o7 v5, o8 v6, o9 v5: equal to the table |
Both writers carry whole items (each stream record has the new image, each export row is the item), so "only if newer by version" is enough. Had we loaded the export first and started the stream afterwards, the stream position would have to be one from before the export, and if the export and load took 24 hours or more, that position would be gone. AWS states this limit for OpenSearch Ingestion's DynamoDB source: "If ingestion from an initial snapshot of a large table takes 24 hours or more, there will be some initial data loss."
Kinesis Data Streams for DynamoDB
DynamoDB can also send the same changes into a Kinesis data stream. It looks like "the same thing, kept longer", and it isn't quite:
| DynamoDB Streams | Kinesis Data Streams for DynamoDB | |
|---|---|---|
| Retention | 24 hours, fixed | 24 hours by default, up to 365 days (extra charge above 24 hours) |
| Start positions | TRIM_HORIZON, LATEST, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER; no timestamp | The same, plus AT_TIMESTAMP |
| Order | Per item, guaranteed | Records "might appear in a different order than when the item changes occurred": order by your own version (ApproximateCreationDateTime, milliseconds by default and microseconds if configured, can tie) |
| Duplicates | Each change once | The same change "might also appear more than once" |
| Readers | At most two per shard (one for global tables) | Shared: 2 MB/s and 5 GetRecords calls a second per shard across all consumers; or enhanced fan-out, 2 MB/s per consumer per shard |
| Cost | Stream read requests, per DynamoDB's pricing | Change data capture units of 1 KB, plus the Kinesis stream itself |
| Where it lives | Owned by the table | One Kinesis stream per table, same account and Region |
From the reference implementation, o7 v3 and v4 arrive as v4, v3, v4:
| Arrives | Consumer (last 2) |
|---|---|
| v4 | Gap after 2: park |
| v3 | Exactly next: apply |
| (parked v4) | Unpark: apply |
| v4 again | 4 ≤ 4: skip |
The maps loop (step 2.6) relies on exactly this: its session index is fed from Kinesis Data Streams for DynamoDB, and route_version absorbs the duplicates and reordering.
Global tables
With global tables, each Region's replica has its own stream, and that stream carries every change, including the ones replicated in from other Regions. A consumer in every Region would therefore send every email once per Region. Side effects must act in one Region only (the item's owning Region), while each Region may still build its own read models from its own stream, as the leaderboard loop (step 3.4) does.
Do you need an outbox on DynamoDB?
Usually not. The stream is the table's commit log, kept by DynamoDB, with nothing to poll. A separate event item written in the same transaction costs more: a transactional write takes 2 write units per item of up to 1 KB, so a 1 KB update that costs 1 unit on its own costs 2 × 2 = 4 units as a transaction with an event item. The leaderboard loop (step 2.6) makes the same call: the outbox would close the gap too, but costs twice the write units per item of a plain write, a read on every poll, and up to one poll interval per repair. Write a semantic event item only when consumers need a meaning the item change doesn't carry.
What to remember from Part 9
- DynamoDB Streams: per item, 24 hours, no timestamps, parent shard before child.
- Defaults on Lambda and Pipes block a shard for up to a day: set retries, age and a destination, and repair by re-reading.
- On a 24-hour stream, the snapshot load must finish inside the retention, or keep reading while you load.
Part 10. Choosing
Students (and interviewers) propose other answers: "just write both and fix it nightly", "use two-phase commit", "use event sourcing". Each has a place. The fair way to compare them is the same change under each, then the same questions for every approach.
The same change, three ways
o7 v4 (SHIPPED) under three designs, from the reference implementation. The hourly compare and the 2-minute coordinator outage are assumptions.
| Approach | What happens | How long a copy is wrong or blocked |
|---|---|---|
| Dual write + an hourly compare | B3: the index holds v3 from 09:06:00.030; the compare at 10:00:00 fixes it | 3,240 s, about 54 minutes wrong |
| Outbox + relay | v4 commits at 10:03:00.2; consumers apply it at 10:03:00.7; repeats are skipped by version | 0.5 s behind, never wrong |
| Two-phase commit (a database and a broker that can join) | The database prepares at 10:03:00.2, holding o7's row lock; the coordinator dies and is back at 10:05:00.2 | o7's row is locked for 120 s: the next change to o7 waits that long |
On equal terms
| Approach | Atomicity | Order | Delay | Load on the source | What breaks |
|---|---|---|---|---|---|
| Dual write | None | Whatever order the writes land in | None added | Each copy written by the request | Lost, ghost and reordered copies, silently |
| Dual write + reconciliation | None; the compare repairs later | Wrong until the next compare | Up to the compare interval to be right | Each request's writes plus the compare's reads | Wrong for a whole interval; the compare job needs its own correct source. Its safe form is a hint plus a stream (Part 1) |
| Outbox + polling relay | Event and change commit together | Per aggregate, with whole-bucket claims and the version check | Up to the poll interval (less with a signal) | Claim queries, reads and deletes; dead rows for vacuum | Table bloat; a relay per team; duplicates to absorb |
| Outbox + log tailing | Event and change commit together | Commit order from one decoder; per aggregate downstream | Decode and publish time | Reading WAL it writes anyway; one extra insert per change | A stalled slot pins WAL; a cap turns it into a gap to republish; one task per connector |
| CDC on business tables | Every committed row change is captured | Commit order; per key downstream | Decode and publish time | No extra writes | Consumers depend on table columns; a rename breaks them; meaning must be inferred |
| The table's own change stream (DynamoDB Streams) | Every committed item change is captured | Per item | Seconds or less, plus the reader's polling | Nothing: kept by the database | 24-hour retention; two readers per shard; a poison record blocks a shard on default settings |
Three more, in sentences, with their costs next to their benefits:
- Log first (write the event to the broker; the database is built from it): atomic by construction, because there is only one write. But the application loses read-your-writes on its own database, and the broker becomes the source of truth, which needs its own retention and backup story.
- Event sourcing (the event log is the source of truth; state is derived): removes the dual write and gives a full history. It costs projection lag, schemas that must stay readable forever, and snapshots to bound replay time. See Primitive #23: Event sourcing and CQRS.
- Two-phase commit (XA): most brokers, and every search index and email provider, can't join; where one can, availability is tied to a coordinator, a prepared participant holds its locks until the coordinator decides, and retries still need to be idempotent. See Primitive #10: Two-phase commit and sagas.
Why not event sourcing, and why not two-phase commit?
Choose this when:
- Outbox + polling relay: one service, modest rates, no connector platform yet (the outbox loop's Rounds 1 and 2).
- Outbox + log tailing: many services or high rates on one database, and a team to run connectors and watch slots (the outbox loop's Round 3).
- CDC on business tables: copying data (a search index, a lake) where consumers want rows, not meanings, and you accept the coupling.
- The table's own stream: DynamoDB tables; add a semantic event item only when the item change doesn't carry the meaning.
- A direct write: only as a hint, with one of the above behind it.
- Sagas (a chain of local transactions with compensations) are how a change spanning services completes; they ride on an outbox at every step (Primitive #10).
What to remember from Part 10
- Pick the source of truth first; everything else follows it.
- CDC trades query load and delay for a connector and a pinned (or trimmed) log.
- 2PC couples availability to a coordinator; event sourcing and log-first make the log the truth, with their own costs.
Part 11. On AWS
Everything on this page is on AWS under different names. Each service reads some commit log, and each has its own retention, order and duplicate rules.
Managed services that use it
| Service | What it reads or provides | Order and duplicates | Retention and limits | What AWS documents |
|---|---|---|---|---|
| Amazon DynamoDB (Streams) | The table's own change log | Each change once; order per item | 24 hours; two readers per shard (one for global tables) | Iterators TRIM_HORIZON, LATEST, AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER; shards split, parent before child; no record for a no-op write; no effect on table performance |
| Amazon Kinesis Data Streams (as Kinesis Data Streams for DynamoDB) | The same item changes, into a Kinesis stream | May be reordered and repeated: order by your version | 24 hours to 365 days; shared or enhanced fan-out reads | ApproximateCreationDateTime to order and de-duplicate; one stream per table, same account and Region; billed in 1 KB change data capture units |
| AWS Lambda (event source mapping) | Polls stream shards and invokes a function | At least once; ParallelizationFactor 1 to 10 keeps per-item order | Concurrency from the Regional pool (1,000 by default) | 4 polls a second per shard; defaults retry until the record expires ("up to one day"); bisect doesn't consume the retry quota; SQS and SNS destinations get metadata only, S3 also the records; LATEST can miss events during create or update |
| Amazon EventBridge Pipes | Moves records from a DynamoDB stream, Kinesis, MSK, Kafka, SQS or Amazon MQ to a target; it does not read a database | Keeps a source's order where the source has one | Same stream limits as its source | DynamoDB source: TRIM_HORIZON or LATEST only; AUTOMATIC_BISECT; MaximumRetryAttempts and MaximumRecordAgeInSeconds default to infinite, so a bad record can block a shard until it expires, as with Lambda |
| Amazon OpenSearch Service (OpenSearch Ingestion's DynamoDB source) | Exports the table (point-in-time recovery), then reads the stream into an index with external versions | Versions make repeats harmless | Counts as a reader of the stream's shards | "If ingestion from an initial snapshot of a large table takes 24 hours or more, there will be some initial data loss" |
| Amazon MSK with MSK Connect | Managed Kafka, and managed Kafka Connect that runs "connectors developed by 3rd parties like Debezium for streaming change logs from databases" | Per partition; at least once | Kafka retention (7 days by default); one task per Debezium PostgreSQL connector | Debezium: outbox router columns and key, outbox DELETEs filtered, id header for de-duplication; snapshot modes. Schemas can be governed with AWS Glue Schema Registry |
| AWS DMS | Ongoing replication (CDC) from database logs to targets including Kinesis and Kafka | One record per changed row "regardless of transactions"; "Kinesis Data Streams don't support deduplication" | Uses a logical replication slot on PostgreSQL sources (test_decoding or pglogical) | Partition key by table by default for record-to-record mapping (so one table can land on one shard), by primary key when ParallelApply* is used; REPLICA IDENTITY FULL for before images |
| Amazon Aurora and Amazon RDS | PostgreSQL logical replication slots; MySQL binlogs | Commit order from the log | The slot pins WAL on the writer (RDS: the instance's storage; Aurora: the cluster volume); RDS for MySQL binlog retention up to 168 hours | Aurora PostgreSQL: no logical decoding from readers; metrics OldestReplicationSlotLag, ReplicationSlotDiskUsage, TransactionLogsDiskUsage. Aurora zero-ETL integrations (to Amazon Redshift and the SageMaker lakehouse) are a "fully managed data pipeline" giving "near real-time" copies; MySQL sources "rely on MySQL binary logging (binlog)"; a Global Database failover makes the integration inactive |
Running it yourself
| Option | What it is | Sizing and notes |
|---|---|---|
| Amazon EC2 (or containers) running Debezium on Kafka Connect, self-run Kafka, a custom relay, or KCL with the DynamoDB Streams Kinesis Adapter | The same mechanics as the managed versions; you own the workers, offsets, upgrades and checkpoints | Kafka: 7-day default retention, cleanup.policy=compact for latest-per-key topics, delete.retention.ms one day by default. The adapter processes shards in lineage order. KCL's lease table is page 03's topic |
| Sizing in words | The shared budgets | Free disk on the writer ≥ WAL rate × the longest stall you ride out, on top of the tables' own growth (Part 5); catch-up time = backlog ÷ (drain rate − arrival rate) (Part 3); readers per shard and Lambda concurrency are shared by every consumer |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like a change stream | Why it isn't |
|---|---|---|
| Amazon SQS, Amazon SNS or an EventBridge event bus, published to by the application | "The event goes to a queue" | Publishing from the app after the commit is a dual write. These are destinations and movers, not a record of committed changes |
| DynamoDB global tables replication | "Changes flow to other Regions" | Replication you don't consume; each replica's own stream is the change stream |
| Aurora Global Database, read replicas, RDS Multi-AZ | "They replay the log" | Physical replication for failover and reads, not change events you can route; and a lagging replica is not a valid snapshot for a newer stream position |
| RDS event subscriptions | "Notifies me of database events" | Instance events (failovers, backups, configuration changes), not data changes |
| AWS CloudTrail | "A log of changes" | An audit log of API calls, not of your data's committed changes |
| Amazon Data Firehose | "Streams data" | A delivery service into stores: it reads from a stream, it doesn't capture changes |
| DynamoDB TTL | "Deletes appear in the stream" | Background cleanup, typically within a few days of expiry: not a deadline and not an event you control |
What to remember from Part 11
- Streams, Kinesis for DynamoDB, Debezium on MSK Connect, DMS, OpenSearch Ingestion and zero-ETL all read a commit log; each has its own retention, order and duplicate rules.
- Lambda's and Pipes' defaults block a shard on one bad record for up to a day.
- A queue or event bus you publish to from the app is a dual write, not a change stream.
Part 12. What you've learned
Back to the 432 a day
Our service lost 432 changes a day to crashes between writes, and its first fix could stall a shard for a day. Here is how each piece closed a hole, snapshot by snapshot:
- One write. The index, the cache and the notifier no longer hear from the request; they follow the commit (S1, Parts 1 and 2). A crash can no longer leave a copy behind silently: the event exists if and only if the change committed.
- At least once, applied by version. The relay died after publishing and republished (S2); a notifier died mid-event and its successor recognised its own feed entry by
src(S3, Part 3). Nothing happened twice. - Order per order. Whole-bucket claims, one partition per key and the exactly-next check kept
o7's versions in order, even across a resize (Part 4). - The log as the mover. CDC replaced the relays' queries with a slot, and the heartbeat kept it moving (Part 5).
- Poison without a stall.
o9's bad address stopped onlyo9, for only the indexer, and was repaired from the primary (S4, Part 6), where Lambda's defaults would have blocked the shard for up to a day (Part 9). - Gaps refilled in order. The invalidated slot's lost events were republished from the outbox before the connector resumed (S5, S6); the indexer, away for 30 hours, rebuilt by noting the position first (S7, Part 7).
- Deletes that stick.
o7's erasure travelled as a versioned event, and a late repair from a lagging replica was dropped by the tombstone and theGONEmarker (S8, Part 8).
The whole story, event by event
| # | Time | Event |
|---|---|---|
| B1 | D1 09:00:00 | Before the outbox: o7 v1 placed; commit, index, cache and email all succeed |
| B2 | 09:05:00 | o7 v2 paid: commit, index; pod killed: no cache write, no "payment received" |
| B3 | 09:06:00.000 to .050 | A commits v3, B commits v4; index writes at .020 (B) and .030 (A), cache at .040 (B) and .050 (A): both copies end at v3 |
| 1 | 10:00:00.0 | o7 v1 placed: order and outbox seq 101 in one transaction |
| 2 | 10:00:00.6 | R1 claims bucket 3, publishes 101 → P1@0, deletes, commits |
| 3 | 10:00:00.7 | Consumers apply o7@1; "order received"; c1#1 |
| 4 | 10:01:00.0, .1 | o7 v2 paid (seq 102); o8 v1 placed (seq 103) |
| 5 | 10:01:00.6 | R1 claims buckets 0 and 3, publishes 103 → P0@0 and 102 → P1@1, acknowledged; killed before its DELETE commits; rows stay |
| 6 | 10:01:00.7 | Consumers apply o7@2 and o8@1; "payment received"; c1#2 |
| 7 | 10:01:05.0 | R1 restarts and republishes 103 → P0@1, 102 → P1@2; deletes; commits |
| 8 | 10:01:05.1 | All three consumers skip both repeats by version |
| 9 | 10:02:00.000 to .010 | o9 v1 gets seq 104, o8 v2 gets seq 105; 105 commits first. Published at 10:02:00.6: 105 → P0@2, 104 → P1@3 |
| 10 | 10:02:30 | R2 starts |
| 11 | 10:03:00.0, .2 | o7 v3 address changed (seq 106); o7 v4 shipped (seq 107) |
| 12 | 10:03:00.6 | R1 claims bucket 3: 106 → P1@4, 107 → P1@5; R2 finds nothing |
| 13 | 10:03:00.7 | Indexer and cache at o7@4. N1 applies v3, then for v4 emails "shipped", appends c1#3, and is killed before its bookkeeping |
| 14 | 10:03:45.7 | P1 moves to N2: v3 skipped; v4 applied with no second email; c1#3 recognised by src; last_seq 3 |
| 15 | 10:04:30.0 | o9 v2 (seq 108) with an address the indexer can't parse |
| 16 | 10:04:30.2 | o7 v5 (seq 109), schema_version 2 with currency; both published at 10:04:30.6 → P1@6, P1@7 |
| 17 | 10:04:30.7 to 36.8 | Indexer fails at 30.7, 30.8, 31.8 and 36.8; dead-letters o9 v2, blocks o9, applies o7@5. Cache and notifier apply both |
| 18 | 10:20:00 to 10:20:30 | o9 v3 (seq 110) → P1@8, parked by the indexer; the repair re-reads o9 v3 on the primary, unblocks; parked v3 skipped |
| 19 | 11:00:00 | Slot S1 created; connector starts (no_data); relays stop; outbox insert-only |
| 20 | 11:10:00 | o8 v3 (seq 111) → P0@3 through the connector |
| 21 | 12:00:00 | The connector dies after its 12:00:00 heartbeat; the lag alarm fires at 12:06:00 |
| 22 | 13:00, 15:00 | o8 v4 (seq 112) and o9 v4 (seq 113) commit; not published |
| 23 | 18:00:00 | Retained WAL passes the cap: S1 invalidated |
| 24 | 19:00:00 to :05 | Slot S2 created; o8 v5 (seq 114) commits at 19:00:02; 111 to 114 republished and acknowledged at 19:00:05; the connector starts on S2 and sends 114 again; repeats skipped |
| 25 | D2 00:00 to D3 06:00 | The indexer is down; o9 v5 (seq 115, D2 03:00) → P1@10 is trimmed before it returns: it can't resume |
| 26 | D3 06:00:00 to :20 | Rebuild: offsets noted (P0 8, P1 11), scan of the primary, o8 v6 (seq 116) commits at 06:00:05 and is read by the scan; the stream's v6 is skipped |
| 27 | D3 07:00:00 | o7 erased: order.erased v6 (seq 117) → P1@11; tombstone (o7, 6) and GONE@6 |
| 28 | D3 07:00:30 | The compare job's "repair" from a replica 40 s behind (v5) is dropped by the tombstone and by SET NX |
The cheat card
| Topic | Remember |
|---|---|
| The one rule | The change and its announcement are one write; everything else follows it |
| Outbox transaction | Update the aggregate (row lock, version + 1), insert the outbox row, commit |
| Relay | Claim whole buckets with SKIP LOCKED, publish in seq order, wait for acks, delete, commit; at least once |
| Consumer | Skip v ≤ last, apply v = last + 1 with its version in one transaction, park v > last + 1 |
| Derived numbered records | Store src; same src on a conflict = your own replay; last_seq never expires |
| Order | Row lock → bucket → partition → one consumer → exactly-next check; adding partitions moves keys at produce time |
| Cursor relays | seq is assigned at insert, not commit: claim-and-delete or read in commit order |
| Slot | Pins WAL from its confirmed position; heartbeats keep it moving; a cap turns a full disk into a gap |
| MySQL binlog | Expires by time; a stalled reader loses its place |
| Poison | Retry a few times, dead-letter, block the key, park its later events, repair by re-reading the primary |
| Gap | New position first, republish from the outbox in order, wait for acks, then resume |
| Bootstrap | Note the position, snapshot at or after it, apply the stream keeping only newer versions |
| Schema | Add optional fields; breaking change = new event type; consumers fail closed |
| Deletes | A versioned event; tombstones kept longer than any late write; caches store "gone"; refills SET NX |
| DynamoDB Streams | Per item, once, 24 hours, no timestamp start, parent shard first, two readers per shard |
| Lambda and Pipes defaults | Retry until expiry: set bisect, retries, maximum age and a destination |
| Formulas | Retained WAL = WAL rate × stall time; drain time = backlog ÷ (drain rate − arrival rate) |
Failure checklist
- Does any service write the database and then another system? A direct write with no stream behind it is a dual write, however fast.
- Does every event carry the aggregate's version, and does every consumer apply by version, never by adding?
- Do relays claim whole buckets (or does one reader decode in commit order), so one aggregate has one publisher?
- Does any relay remember "published up to
seqN"? - Does every consumer check versions contiguously, and park rather than skip a gap?
- Do derived numbered records store
src, and doeslast_seqoutlive what it numbers? - Is heartbeat lag alarmed, and does every captured database have a heartbeat table in its publication?
- Is
max_slot_wal_keep_sizeset below free disk minus table growth, with a tested republish runbook for when it fires? - Is the outbox kept longer than the worst time to detect and repair a gap?
- Does every consumer have an error boundary, a dead-letter owner, a repair path that reads the primary, and a periodic compare?
- Are Lambda mappings and Pipes on DynamoDB streams configured with bisect, limited retries, a maximum record age, a destination and
TRIM_HORIZON, withIteratorAgealarmed? - Do deletes travel as versioned events, with tombstones kept longer than any late or repair write can arrive?
Think-first drills
Drill 1. A consumer's last_version for o5 is 7. A batch arrives with o5 v9, v8, v8, v10, v12. What is applied, skipped and parked after each record, and what is last_version at the end?
Drill 2. A slot stalls at a WAL rate of 5 MB/s with max_slot_wal_keep_size set to 100 GB. When is it invalidated, and what must the recovery do, in which order?
Drill 3. A search index fed from a DynamoDB table's stream was down for 3 days. List the recovery steps in order, say which start position you use and why, and what limits how long the snapshot load may take.
Interview questions
| Question | Model answer |
|---|---|
| What's wrong with writing the database and then publishing to Kafka? | Two independent writes: a crash between them loses the event (lost), publishing first can announce a change that then fails to commit (ghost), and two writers can land in opposite orders (reordered). Retries in the request can't help because they run in the process that dies. Make the event part of the commit: an outbox row, or the database's own log. |
| Walk me through a transactional outbox and what happens when the relay crashes. | One transaction updates the aggregate (taking its row lock, bumping its version) and inserts an outbox row with the version. A relay claims whole buckets with SKIP LOCKED, publishes their rows in seq order keyed by aggregate, waits for acks, deletes and commits. A crash after publishing and before the commit republishes the rows, so delivery is at least once and consumers skip by version. |
| CDC or a polling relay: when each? | A relay is simplest for one service at modest rates, but costs queries, delete churn and a poll interval of delay on the busiest database. CDC reads the WAL the database writes anyway, adds almost no query load and keeps commit order, but you run connectors, watch slots (a stalled slot pins WAL; a cap turns that into a gap to republish), and one connector task sets the ceiling. |
| How do you keep per-entity order end to end, and what breaks it? | Every hop: the row lock orders versions; one publisher per aggregate (bucket claims or one decoder); key by aggregate so it stays in one partition; one consumer per partition; and an exactly-next version check as the safety net, parking gaps. It breaks with row-level claims, a partition count change (keys move at produce time), multi-event transactions across keys, and a failover that reuses versions (hence the writer epoch). |
| A consumer has been down longer than the stream's retention. How do you rebuild without losing a change? | Note the stream position first, then read a snapshot of the source at or after it (the primary, an export, or a replica that has caught up past the position), then apply the stream from the noted position keeping only newer versions. On a 24-hour stream, keep consuming while the snapshot loads, or finish the load well inside 24 hours. |
| A transfer writes a debit and a credit event in one transaction. Can a consumer see only one? What do you do? | Yes: keyed by different accounts, they land in different partitions and arrive independently. Either publish one event for the whole transfer, key both by one ordering scope, or carry a transaction id and event count so a consumer that needs both waits for both (with a timeout alarm). Choose the aggregate as the smallest scope whose events must be seen together. |
Where to go next
- Write-Ahead Log, fsync & Group Commit: the log a change stream reads, why a reader pins it (Part 8) and what
acks=allpromises (Part 10). - Leases, Fencing Tokens & Distributed Locks: the
SKIP LOCKEDclaim (Part 9) and howwriter_epochis bumped across a failover (Part 7). - The Idempotency & Effectively-Once Processing loop primitive (coming): where consumers store what they applied, notification idempotency keys, and exactly-once per hop.
- The Replication, Quorums & Read-Your-Writes loop primitive (coming): replica lag, and why a lagging replica can't be a snapshot. The Sharding, Hot Keys & Rebalancing and Multi-Region Failover loop primitives (coming) pick up partition choice and cross-Region event flow.
- Primitives: #05 Message queues vs event streams, #21 Isolation levels (what the outbox transaction's row lock does), #10 Two-phase commit and sagas, #23 Event sourcing and CQRS.
- Drill: The Dual-Write That Broke Search Consistency, answered in Parts 1 and 5.
- Loops: the transactional outbox (steps 1.0 to 1.6, 2.1 to 2.4, 3.1 to 3.6), payments (2.1, 2.3, 2.5, 3.6), wallet (2.2, 2.5), hotel reservations (2.1, 2.2, 2.4, 3.1, 3.5, R2.8), matching engine (3.4), proximity (2.1, 2.6, R2.8), ride-sharing (R1.6, 2.4, 2.6, 3.6), maps (2.6, R2.8), nearby friends (2.5), email (2.2, 2.4), YouTube (R1.6, R2.8), Drive (1.4, 2.3, 2.6, 3.1), ad-click aggregation (3.1), mobile stock trading (1.2, R2.5), mobile news feed (2.2, 2.3, 2.4), news feed (1.4, 3.1, 3.3), leaderboard (1.3, 2.6, 3.4), chat (R1.6), URL shortener (3.3), S3-like storage (3.1), web crawler (3.6), message queue (2.2), notification system (R1.11, 3.5) and mobile chat (1.1).