Round four: “six teams want this data”
Every other team in the company needs to know when money moves. Answering that by letting them all call the ledger, or by having the ledger call all of them, is how a core banking system becomes impossible to change. This round introduces the log as the integration primitive, goes deep enough into Kafka to explain what a partition physically is and what the ISR actually guarantees, and then solves the problem that makes naive event publishing unsafe for money: the dual write.
The pressure: every team wants a hook
“Notifications wants to text the customer. Fraud wants to score every transaction. Analytics wants everything. Loans wants to know when a borrower gets paid. Four more teams are asking next quarter. How do you serve them?”
The real question underneath: how do you let the company build on the ledger without every new consumer becoming a change to the ledger?
- Adding a consumer requires zero changes to the ledger service.
- A slow or broken consumer cannot slow or break posting.
- No event is ever lost, even if every consumer is down for hours.
- Per-account ordering is preserved, because "debited then credited" and the reverse mean different things.
- A new consumer can replay history from the beginning to build its own state.
Why not synchronous calls to six services
Dispose of the obvious approach properly, because being able to say precisely why it fails is what justifies the infrastructure that replaces it.
// the design that seems reasonable and is not await db.postTransfer(...); await notifications.send(...); // +40 ms await fraud.score(...); // +60 ms await analytics.record(...); // +30 ms await loans.notify(...); // +50 ms return ok; // 180 ms of other people's latency
The three fatal properties
1. latency adds
transfer p99 = ledger p99 + Σ(every consumer's p99)
our 300 ms budget is gone by the fourth callee
2. availability multiplies
6 dependencies at 99.9% each
= 0.9996 = 99.4% → 52 hours of downtime a year
the ledger becomes less available than its least available consumer
3. partial failure has no good answer
money posted, notification failed. now what?
roll back a committed posting? no.
fail the request after taking the money? no.
ignore it? then the event is lost.And a fourth, which is organisational and in practice the most damaging: every new consumer means a code change, a review, and a deployment of the ledger. The most safety-critical service in the company becomes the one that changes most often, purely for reasons that have nothing to do with money.
The one case that must stay synchronous
Fraud, sometimes. If the business requires blocking a transaction before it posts, no amount of asynchrony helps: you need an in-line decision. The resolution is to split fraud into two paths, a fast synchronous check on the posting path with a hard timeout and a fail-open default, and a rich asynchronous analysis that consumes events. Part 9 builds exactly that.
The log as the integration primitive
The abstraction is older and simpler than the products built on it: an append-only, ordered, durable, replayable sequence of records, where each consumer tracks its own position independently.
| Property | What it gives the design |
|---|---|
| Append-only | Writes are sequential, so throughput is enormous and history is immutable. |
| Ordered | "Debited then credited" is distinguishable from the reverse, within a partition. |
| Durable and retained | A consumer down for six hours loses nothing; it resumes from its offset. |
| Per-consumer position | Consumers are fully decoupled. A slow one falls behind without affecting anybody. |
| Replayable | A new consumer, or a consumer with a fixed bug, rebuilds its state from offset zero. |
Why this is not a queue
A distinction worth being crisp about, because interviewers probe it:
| Queue (SQS, RabbitMQ) | Log (Kafka, Kinesis) | |
|---|---|---|
| On consumption | Message is removed | Record stays; the reader's offset advances |
| Multiple independent consumers | Needs fan-out to separate queues | Native. Each group has its own offset |
| Replay | Impossible once acknowledged | Seek to any offset |
| Ordering | Best effort, or FIFO with throughput limits | Total order per partition |
| Best at | Work distribution: one job, one worker | Event distribution: one fact, many interested parties |
Our problem is one fact with many interested parties, and we want replay. That is a log. Part 6 will use an actual queue for outbound payout work, which is the other shape, and using both for their respective strengths is the correct answer rather than a compromise.
The event, designed
{
"event_id": "01JBQ2X...", // ULID: unique AND sortable
"event_type": "ledger.entry.posted",
"version": 2, // schema version, always present
"occurred_at": "2026-03-14T09:12:03.441Z",
// the facts. a consumer must not need to call us back to understand this.
"journal_id": "j-8f1a",
"entry_id": 998812,
"account_id": "acct-771",
"amount": "-50000", // string: JSON numbers lose precision
"currency": "NGN",
"kind": "transfer",
"balance_after": "75075", // saves every consumer a query
// provenance, so a consumer can trace and deduplicate
"correlation_id": "req-4c2e",
"shard": 7
}9007199254740993 silently becomes 9007199254740992. A
string preserves the exact integer. Include balance_after. It
is denormalised and it means six consumers do not each make a callback to the
ledger for context, which would rebuild the coupling we are removing.Kafka internals: partitions, offsets, ISR
Enough internals to justify the configuration choices rather than copying them.
What a partition physically is
A partition is a directory on a broker's disk, containing segment files. Each segment is an append-only sequence of records, with index files mapping offsets to byte positions.
/var/lib/kafka/ledger.entry.posted-7/
00000000000000000000.log ← records, appended
00000000000000000000.index ← offset → byte position
00000000000000000000.timeindex
00000000000524288000.log ← next segment, rolled at 1 GB
a write is an append to the active segment: a sequential write.
a read is a seek via the index, then a sequential scan.
this is why Kafka is fast. it is not clever, it is sequential.Two further mechanisms carry most of the performance: the page cache,
since recently written records are served from memory without Kafka managing a
cache itself, and zero-copy transfer via
sendfile, which moves bytes from page cache to socket without
passing through the application at all.
Replication and the in-sync replica set
| Setting | Meaning | Our value, and why |
|---|---|---|
acks | 0 fire and forget, 1 leader only, all every in-sync replica | all. A leader failure after a leader-only ack loses the record silently. |
min.insync.replicas | How many replicas must be in-sync for a write to be accepted | 2 with replication.factor=3. Without this, acks=all is nearly worthless: if followers fall out of the ISR, "all" can mean "the leader alone". |
unclean.leader.election | May an out-of-sync replica become leader to restore availability? | false. Electing a stale leader discards committed records. For money we take the outage. |
enable.idempotence | Producer sequence numbers let the broker deduplicate its own retries | true. Removes duplicates caused by producer-side retries, though not application-level duplicates. |
retention.ms | How long records are kept | Long. 30 days on the entries topic so a consumer can be rebuilt; the warehouse keeps the permanent copy. |
acks=all. Very few add
min.insync.replicas=2 and explain that acks=all alone
degrades to acks=1 when replicas fall behind, because "all in-sync
replicas" can be a set of size one. Saying that sentence demonstrates you have
reasoned about the failure mode instead of copying a config.Ordering guarantees, and what they cost
Kafka guarantees total order within a partition and nothing across partitions. Every ordering decision follows from that one sentence.
1 partition → global total order, one consumer maximum,
~10 MB/s ceiling
N partitions → order only within each, N consumers in
parallel
so: ordering and parallelism are the same dial
the skill is choosing the smallest unit that needs orderingWhat actually needs ordering in a bank
| Scope | Ordering needed? | Why |
|---|---|---|
| Per account | Yes, strictly | A credit then a debit leaves a different balance trajectory than the reverse, and a notification sent in the wrong order confuses the customer. |
| Per journal (both legs of one transfer) | Useful | A consumer reconstructing a transfer prefers to see both legs together. |
| Per customer | Rarely | Only for cross-wallet logic, which is uncommon. |
| Globally | No | No consumer needs to know that account A's debit preceded unrelated account B's credit. Requiring it would cap us at one partition. |
So per-account ordering is the requirement, which means
account_id is the partition key. And one subtlety worth knowing:
ordering also depends on producer configuration.
// max.in.flight.requests.per.connection > 1 with retries CAN reorder: // request A fails and is retried while request B has already succeeded. // // enable.idempotence=true fixes this: the broker uses producer sequence // numbers to reject out-of-order batches, so ordering survives retries // even with up to 5 requests in flight. keep idempotence on. enable.idempotence = true max.in.flight.requests.per.connection = 5 // safe WITH idempotence retries = Integer.MAX_VALUE delivery.timeout.ms = 120000 // bound total attempt time
Partition key selection for financial events
Same reasoning as the shard key in Part 3, different consequences, and worth doing explicitly because the right answer differs.
| Key | Distribution | Ordering provided | Verdict |
|---|---|---|---|
account_id | Even, with one exception | Per account. Exactly what is required. | Chosen. Note both legs of a transfer land on different partitions, which consumers must handle. |
journal_id | Even | Both legs together, but no per-account order | Rejected. Per-account ordering is the requirement that matters. |
customer_id | Even | Per customer, which implies per account | Reasonable alternative. Stronger than needed, and worse tail: a corporate customer concentrates into one partition. |
null (round robin) | Perfect | None | Fine for the analytics topic, where order does not matter. Wrong for notifications and fraud. |
The hot partition, which is the same problem as the hot row
The fee revenue account appears on most transactions, so keying by
account_id sends a large fraction of all events to one partition.
The same structural fix applies, for the same reason:
// the fee account was already split into 64 sub-accounts in Part 3 to
// solve the hot row. that split solves the hot partition for free,
// because the sub-account id is what gets hashed.
function partitionKey(e: LedgerEvent): string {
return e.account_id; // 'fee_revenue:41' spreads; 'fee_revenue' would not
}Sizing the partition count
peak events/s = 464 transfers × ~3.2 entries ≈ 1,500 events/s
allow 4x growth = 6,000 events/s
one consumer instance handles ~500 events/s (with database work)
→ need 12 parallel consumers at peak
partitions must be ≥ desired consumers, and cannot be reduced later
→ choose 48 partitions: 4x current need, room to grow consumers
why not 1,000? each partition costs file handles, memory,
and rebalance time. more partitions also means more end-to-end latency
for replication fan-out.
The dual-write problem, stated precisely
Here is the bug that makes naive event publishing unacceptable for money. It is the most important thing in this part.
// this code is wrong, and the wrongness is not visible by reading it await db.postTransfer(...); // 1. commits to Postgres await kafka.publish(entryPosted); // 2. publishes to Kafka
Two independent systems, two separate commits, no transaction spanning them. There is a window between them, and a crash in that window leaves the two permanently disagreeing.
The four outcomes, two of which are bugs
| DB commit | Kafka publish | Result |
|---|---|---|
| ✓ | ✓ | Correct. |
| ✗ | ✗ | Correct. Nothing happened. |
| ✓ | ✗ | Money moved, nobody was told. No notification, fraud never scored it, analytics is short, loans missed a repayment. Silent and permanent. |
| ✗ | ✓ | An event for money that does not exist. The customer is notified of a transfer that never happened, and fraud scores a phantom. |
Reordering the two statements does not help, it only changes which bug you get. Wrapping them in a try/catch does not help, because the process can die between them. Retrying the publish does not help, because the retry state itself lives in the memory that just died.
the general statement
you cannot atomically commit to two systems that do not
participate in a shared transaction.
therefore: write to ONE system, and derive the other from it.
this is the entire idea behind both solutions that follow.
The transactional outbox
The fix follows directly from the general statement: write to one system. The event is committed into the database, in the same transaction as the money.
CREATE TABLE outbox ( id BIGSERIAL PRIMARY KEY, -- monotonic: preserves order aggregate_id TEXT NOT NULL, -- becomes the partition key event_type TEXT NOT NULL, payload JSONB NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), published_at TIMESTAMPTZ -- null until relayed ); -- partial index: only unpublished rows, so it stays tiny regardless of volume CREATE INDEX outbox_pending ON outbox (id) WHERE published_at IS NULL;
-- one transaction. the money and the intent to announce it are inseparable.
BEGIN;
INSERT INTO journal ...;
INSERT INTO entries ...; -- debit
INSERT INTO entries ...; -- credit
INSERT INTO outbox (aggregate_id, event_type, payload)
VALUES ($account, 'ledger.entry.posted', $json);
COMMIT;
-- either all four rows exist, or none do. the window is gone.The relay
// a separate process, running continuously. its only job is to move rows
// from the outbox to Kafka. it may crash freely; nothing is lost.
async function relay() {
while (running) {
const batch = await db.query(`
SELECT id, aggregate_id, event_type, payload FROM outbox
WHERE published_at IS NULL
ORDER BY id -- order preserved: BIGSERIAL is monotonic
LIMIT 500
FOR UPDATE SKIP LOCKED -- many relay instances, disjoint work
`);
if (!batch.length) { await sleep(50); continue; }
// publish first, mark second. a crash between them causes a
// DUPLICATE, never a loss. duplicates are handled by consumers.
await kafka.sendBatch(batch.map(toRecord));
await db.query(`UPDATE outbox SET published_at = now() WHERE id = ANY($1)`,
[batch.map(r => r.id)]);
}
}Costs, honestly
| Cost | Detail | Mitigation |
|---|---|---|
| Extra write per transaction | One more insert in the hot path, roughly 5% overhead | Accepted. It is one sequential insert. |
| Added latency to publish | Polling interval, typically 10 to 100 ms | Tune the interval; or use LISTEN/NOTIFY to wake the relay immediately. |
| Table growth | Published rows accumulate | Delete or partition-drop published rows older than a few hours. The partial index keeps queries fast regardless. |
| A process to operate | Relay lag becomes a thing to monitor and page on | Alert on oldest unpublished row age. This is a first-class SLI in Part 13. |
CDC as an alternative to the outbox
Change data capture reads the database's own replication log and turns committed changes into events. Same guarantee, obtained differently, and you should be able to argue both sides.
Postgres logical replication → Debezium → Kafka
the WAL already is an ordered, durable log of every commit.
CDC simply reads it. no application change at all.
| Outbox | CDC | |
|---|---|---|
| Application change | One insert per transaction | None |
| Event shape | Designed. You publish a business event with exactly the fields consumers need | Row-shaped. Consumers receive table diffs and must interpret your schema |
| Coupling | Consumers depend on an explicit contract | Consumers depend on your table structure. A column rename breaks them |
| Operational surface | A relay you own | Debezium, Kafka Connect, replication slots |
| Replication slot risk | None | Real. A stalled consumer prevents WAL reclamation and can fill the primary's disk |
| Latency | Polling interval | Lower. Streams from the WAL |
| Deletes and schema changes | Explicit, because you author the event | Awkward; DDL needs special handling |
Consumer groups, lag, and rebalancing storms
A consumer group is a set of instances that share a topic's partitions. Each partition is assigned to exactly one member, which is what gives both parallelism and per-partition ordering.
48 partitions, 12 consumer instances → 4 partitions each
one instance dies → its 4 are reassigned → 11 instances, ~4.4 each
more consumers than partitions → the extras sit idle
partition count is the hard ceiling on parallelismRebalancing, and why it can become a storm
When membership changes, the group must reassign partitions. Under the default eager protocol, every consumer stops consuming during the reassignment. A poorly tuned deployment turns that into a cascade:
- A consumer is slow to call
poll(), typically because one message took too long to process. - It misses
max.poll.interval.msand the coordinator declares it dead. - A rebalance begins; all consumers pause.
- The pause increases lag, so when they resume each has more work queued.
- More work means slower polls, so another consumer is evicted, and step 3 repeats.
// the settings that prevent it max.poll.records = 100 // small batches: bounded time per poll cycle max.poll.interval.ms = 300000 // generous ceiling before eviction session.timeout.ms = 45000 // heartbeat failure detection heartbeat.interval.ms = 3000 // heartbeats run on a separate thread // and the important structural one: cooperative rebalancing, which // revokes only the partitions that must move instead of stopping everyone. partition.assignment.strategy = CooperativeStickyAssignor
Lag as the primary health metric
lag = log end offset − committed consumer offset
lag in records is what tools report.
lag in seconds is what the business cares about:
time_lag ≈ record_lag ÷ consumption_rate
alert on time lag, not record lag. 50,000 records behind
at 10,000/s is 5 seconds and fine. at 100/s it is 8 minutes and an incident.
Idempotent consumers and the dedup store
Our relay deliberately produces duplicates rather than losses, and Kafka delivers at least once. So every consumer must be safe to run twice on the same event. That is a requirement on consumers, which the platform team must make easy to satisfy.
Three ways to be idempotent, in order of preference
- Setting a value is idempotent; incrementing is not. Prefer
SET balance_cache = XoverINCR. DELETE FROM cache WHERE key = ...is idempotent. Cache invalidation is naturally safe, which is why Part 3 used it.- An upsert keyed by the event's own identifier is idempotent.
- A table of
(consumer, event_id)with a unique constraint. - Insert the marker and do the work in the same transaction, so they cannot diverge.
- A duplicate raises a unique violation and is skipped.
- Pass the event id to the SMS provider as its idempotency key, so the provider refuses the second send.
- Best where the side effect is external and you cannot make it transactional.
// the standard shape: marker and effect in one transaction
async function handle(e: LedgerEvent) {
await db.transaction(async tx => {
try {
await tx.query(
`INSERT INTO processed_events (consumer, event_id) VALUES ($1, $2)`,
['notifications', e.event_id]
);
} catch (err) {
if (isUniqueViolation(err)) return; // already done. skip.
throw err;
}
await doTheWork(tx, e); // same transaction: cannot diverge
});
}Keeping the dedup store bounded
-- monthly partitions, so cleanup is a DROP rather than a mass DELETE. -- retention only needs to exceed the topic's retention. CREATE TABLE processed_events ( consumer TEXT NOT NULL, event_id TEXT NOT NULL, seen_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (consumer, event_id) ) PARTITION BY RANGE (seen_at); -- DROP TABLE processed_events_2026_01; ← instant, no vacuum pressure
Schema evolution with a registry
The event schema is now a public API consumed by six teams and, because of replay, by historical versions of those consumers. It must evolve without coordinated deployments.
| Compatibility mode | Permits | Meaning |
|---|---|---|
| Backward | Delete a field; add an optional field | New consumers can read old events. Upgrade consumers first. |
| Forward | Add a field; delete an optional field | Old consumers can read new events. Upgrade producers first. |
| Full | Add or delete optional fields only | Either order works. Our choice, since we cannot coordinate six teams. |
| None | Anything | Every change is a potential outage for somebody. |
- Never change a field's meaning or units. Adding
amount_minorbeside a deprecatedamountis correct; redefiningamountis a catastrophe that is invisible in code review. - Never remove a required field. Deprecate, wait, then remove after confirming no consumer reads it.
- Always add fields as optional with a default.
- Always carry an explicit
version, so a consumer can branch rather than guess. - Prefer a new event type over a breaking change to an existing one.
ledger.entry.posted.v2alongside v1 lets consumers migrate on their own schedule.
Format choice
| Format | Size | Schema enforcement | Fit |
|---|---|---|---|
| JSON | Largest | None by default | Easy to debug, easy to break. Acceptable with a registry validating on publish. |
| Avro | Compact | Strong, with a registry and built-in evolution rules | Chosen. Designed for exactly this problem, and the Kafka ecosystem's default. |
| Protobuf | Compact | Strong, via field numbers | Excellent, and the right pick if the company already uses it for gRPC. Part 14 revisits this. |
Sketch v4: event-driven core
The topic layout
| Topic | Key | Partitions | Retention | Consumers |
|---|---|---|---|---|
ledger.entry.posted | account_id | 48 | 30 days | notifications, fraud, loans, support |
ledger.journal.completed | journal_id | 24 | 30 days | reconciliation, reporting |
ledger.account.changed | account_id | 12 | compacted | everyone needing current account state |
ledger.dlq | original key | 6 | 90 days | humans |
The compacted topic deserves a note: log compaction retains only the latest record per key, so a new consumer can rebuild current account state by reading it once rather than replaying all history. It gives you a snapshot as a stream.
What changed, and the cost accepted
| Change | Driven by | Cost accepted |
|---|---|---|
| Kafka as the integration layer | Six consumers, and more arriving | A cluster to operate; eventual consistency for consumers |
| Transactional outbox | The dual-write problem | One extra insert per transaction, a relay to run, and duplicates by design |
| Partition by account | Per-account ordering requirement | Both legs of a transfer land on different partitions |
| Idempotent consumers | At-least-once delivery | A dedup store per consumer, partitioned for cleanup |
| Avro plus a registry, full compatibility | Six independent deployment schedules | Schema discipline, and a registry in the path |
| CDC for analytics only | The warehouse wants table shape, not business events | Replication slot monitoring, with real disk risk if it stalls |
Which is the opening the interviewer takes next, because a bank that cannot lend is not much of a bank.