Part 4 · 13 chapters · ~20 min

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.

42

The pressure: every team wants a hook

interviewer

“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?

notifications Needs near-real-time delivery, tolerates seconds of delay, and must never send twice.
fraud Needs every transaction, in order per account, and sometimes needs to act before the posting completes.
analytics Needs everything, eventually. Minutes of lag are fine. Volume is enormous.
loans Needs filtered events: credits to accounts belonging to borrowers with active facilities.
ledger reporting Needs a complete, ordered, replayable history for regulatory reconstruction.
customer support Needs a searchable timeline per customer, joined with non-ledger context.
non-functional, new
  1. Adding a consumer requires zero changes to the ledger service.
  2. A slow or broken consumer cannot slow or break posting.
  3. No event is ever lost, even if every consumer is down for hours.
  4. Per-account ordering is preserved, because "debited then credited" and the reverse mean different things.
  5. A new consumer can replay history from the beginning to build its own state.
43

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.

code
// 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

worked numbers
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 inversion that solves it
Stop having the ledger know its consumers. The ledger publishes what happened. Consumers decide what to do about it. Six consumers becomes the same amount of ledger code as one, and zero, because the publish is the same statement either way.

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.

44

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.

PropertyWhat it gives the design
Append-onlyWrites are sequential, so throughput is enormous and history is immutable.
Ordered"Debited then credited" is distinguishable from the reverse, within a partition.
Durable and retainedA consumer down for six hours loses nothing; it resumes from its offset.
Per-consumer positionConsumers are fully decoupled. A slow one falls behind without affecting anybody.
ReplayableA 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 consumptionMessage is removedRecord stays; the reader's offset advances
Multiple independent consumersNeeds fan-out to separate queuesNative. Each group has its own offset
ReplayImpossible once acknowledgedSeek to any offset
OrderingBest effort, or FIFO with throughput limitsTotal order per partition
Best atWork distribution: one job, one workerEvent 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

code
{
  "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
}
two decisions in that payload worth defending
Amounts as strings. JSON numbers are IEEE 754 doubles in most parsers, so 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.
45

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.

worked numbers
          /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

SettingMeaningOur value, and why
acks0 fire and forget, 1 leader only, all every in-sync replicaall. A leader failure after a leader-only ack loses the record silently.
min.insync.replicasHow many replicas must be in-sync for a write to be accepted2 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.electionMay 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.idempotenceProducer sequence numbers let the broker deduplicate its own retriestrue. Removes duplicates caused by producer-side retries, though not application-level duplicates.
retention.msHow long records are keptLong. 30 days on the entries topic so a consumer can be rebuilt; the warehouse keeps the permanent copy.
the detail that separates candidates
Nearly everyone says 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.
ISR
what acks=all actually guarantees
swipe the figure sideways, or tap expand for full screen
1/8
3 replicas
A single partition is replicated three times: one leader and two followers. All reads and writes go to the leader.
46

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.

worked numbers
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 ordering

What actually needs ordering in a bank

ScopeOrdering needed?Why
Per accountYes, strictlyA 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)UsefulA consumer reconstructing a transfer prefers to see both legs together.
Per customerRarelyOnly for cross-wallet logic, which is uncommon.
GloballyNoNo 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.

code
// 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
47

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.

KeyDistributionOrdering providedVerdict
account_idEven, with one exceptionPer account. Exactly what is required.Chosen. Note both legs of a transfer land on different partitions, which consumers must handle.
journal_idEvenBoth legs together, but no per-account orderRejected. Per-account ordering is the requirement that matters.
customer_idEvenPer customer, which implies per accountReasonable alternative. Stronger than needed, and worse tail: a corporate customer concentrates into one partition.
null (round robin)PerfectNoneFine 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:

code
// 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
}
the pattern to notice and say out loud
The hot row, the hot partition, and the hot shard are one problem wearing three costumes: a single point every transaction must touch. The fix is always split then reassemble, and a fix applied at one layer often resolves the others. Recognising that is more valuable than memorising any of the three individually.

Sizing the partition count

worked numbers
          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.
        
48

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.

code
// 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 commitKafka publishResult
✓✓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.

worked numbers
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.
        
how to raise this unprompted
Do not wait to be asked. When you draw the arrow from the ledger to Kafka, say: "This arrow is a dual write, and I need to address it before moving on." It signals that you have shipped event-driven systems, because this is the bug that teaches everyone the lesson the hard way.
dual write
and the outbox that eliminates it
swipe the figure sideways, or tap expand for full screen
1/9
two writes
The ledger service needs to do two things: commit the money to Postgres, and publish an event to Kafka.
49

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.

code
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;
code
-- 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

code
// 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)]);
  }
}
the deliberate choice inside the relay
Publish before marking, never the reverse. If the relay dies between the two, the row is still unpublished and will be sent again, producing a duplicate. If we marked first, a crash would produce a lost event. We choose the failure mode we have already defended against, which is exactly the at-least-once plus idempotent consumers argument from chapter 29.

Costs, honestly

CostDetailMitigation
Extra write per transactionOne more insert in the hot path, roughly 5% overheadAccepted. It is one sequential insert.
Added latency to publishPolling interval, typically 10 to 100 msTune the interval; or use LISTEN/NOTIFY to wake the relay immediately.
Table growthPublished rows accumulateDelete or partition-drop published rows older than a few hours. The partial index keeps queries fast regardless.
A process to operateRelay lag becomes a thing to monitor and page onAlert on oldest unpublished row age. This is a first-class SLI in Part 13.
50

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.

worked numbers
          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.
        
OutboxCDC
Application changeOne insert per transactionNone
Event shapeDesigned. You publish a business event with exactly the fields consumers needRow-shaped. Consumers receive table diffs and must interpret your schema
CouplingConsumers depend on an explicit contractConsumers depend on your table structure. A column rename breaks them
Operational surfaceA relay you ownDebezium, Kafka Connect, replication slots
Replication slot riskNoneReal. A stalled consumer prevents WAL reclamation and can fill the primary's disk
LatencyPolling intervalLower. Streams from the WAL
Deletes and schema changesExplicit, because you author the eventAwkward; DDL needs special handling
the decision, and the reason
Outbox for the money path. The event schema is a public contract consumed by six teams, and it must be decoupled from our table layout so we can refactor the ledger without breaking notifications. CDC for the analytics path, where Part 11 wants full table replication into the warehouse and the row shape is exactly what is wanted. Using both, each where it fits, is the senior answer. And the replication-slot disk risk is worth naming out loud, because it has taken down real production databases.
51

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.

worked numbers
          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 parallelism

Rebalancing, 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:

how the storm forms
  1. A consumer is slow to call poll(), typically because one message took too long to process.
  2. It misses max.poll.interval.ms and the coordinator declares it dead.
  3. A rebalance begins; all consumers pause.
  4. The pause increases lag, so when they resume each has more work queued.
  5. More work means slower polls, so another consumer is evicted, and step 3 repeats.
code
// 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

worked numbers
          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.
        
a per-consumer SLO, not a global one
Each consumer gets its own lag budget matched to its purpose. Notifications: 30 seconds, because a late text is a bad experience. Fraud: 5 seconds, because late detection means the money has already left. Analytics: 15 minutes, because nobody notices. One global threshold would either page constantly or miss the case that matters.
52

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

one: be naturally idempotent
  1. Setting a value is idempotent; incrementing is not. Prefer SET balance_cache = X over INCR.
  2. DELETE FROM cache WHERE key = ... is idempotent. Cache invalidation is naturally safe, which is why Part 3 used it.
  3. An upsert keyed by the event's own identifier is idempotent.
two: track processed event ids
  1. A table of (consumer, event_id) with a unique constraint.
  2. Insert the marker and do the work in the same transaction, so they cannot diverge.
  3. A duplicate raises a unique violation and is skipped.
three: make the side effect deduplicate
  1. Pass the event id to the SMS provider as its idempotency key, so the provider refuses the second send.
  2. Best where the side effect is external and you cannot make it transactional.
code
// 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
  });
}
the ordering subtlety that bites people
Insert the marker before doing the work, inside the transaction. If the work came first and the process died before the marker, the work would repeat. Because both are in one transaction, a failure rolls back both, and the event is redelivered and processed exactly once in effect.

Keeping the dedup store bounded

code
-- 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
53

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 modePermitsMeaning
BackwardDelete a field; add an optional fieldNew consumers can read old events. Upgrade consumers first.
ForwardAdd a field; delete an optional fieldOld consumers can read new events. Upgrade producers first.
FullAdd or delete optional fields onlyEither order works. Our choice, since we cannot coordinate six teams.
NoneAnythingEvery change is a potential outage for somebody.
rules for a money event schema
  1. Never change a field's meaning or units. Adding amount_minor beside a deprecated amount is correct; redefining amount is a catastrophe that is invisible in code review.
  2. Never remove a required field. Deprecate, wait, then remove after confirming no consumer reads it.
  3. Always add fields as optional with a default.
  4. Always carry an explicit version, so a consumer can branch rather than guess.
  5. Prefer a new event type over a breaking change to an existing one. ledger.entry.posted.v2 alongside v1 lets consumers migrate on their own schedule.

Format choice

FormatSizeSchema enforcementFit
JSONLargestNone by defaultEasy to debug, easy to break. Acceptable with a registry validating on publish.
AvroCompactStrong, with a registry and built-in evolution rulesChosen. Designed for exactly this problem, and the Kafka ecosystem's default.
ProtobufCompactStrong, via field numbersExcellent, and the right pick if the company already uses it for gRPC. Part 14 revisits this.
54

Sketch v4: event-driven core

The topic layout

TopicKeyPartitionsRetentionConsumers
ledger.entry.postedaccount_id4830 daysnotifications, fraud, loans, support
ledger.journal.completedjournal_id2430 daysreconciliation, reporting
ledger.account.changedaccount_id12compactedeveryone needing current account state
ledger.dlqoriginal key690 dayshumans

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

ChangeDriven byCost accepted
Kafka as the integration layerSix consumers, and more arrivingA cluster to operate; eventual consistency for consumers
Transactional outboxThe dual-write problemOne extra insert per transaction, a relay to run, and duplicates by design
Partition by accountPer-account ordering requirementBoth legs of a transfer land on different partitions
Idempotent consumersAt-least-once deliveryA dedup store per consumer, partitioned for cleanup
Avro plus a registry, full compatibilitySix independent deployment schedulesSchema discipline, and a registry in the path
CDC for analytics onlyThe warehouse wants table shape, not business eventsReplication slot monitoring, with real disk risk if it stalls
how to close round four
"v4 inverts the dependency: the ledger publishes facts and knows none of its consumers, so the seventh team costs me nothing. The two things I want to be precise about are that I solved the dual write with an outbox rather than publishing after commit, and that the relay produces duplicates on purpose because losing an event is worse and consumers are idempotent anyway. What is still missing is products: this ledger only knows how to move money that already exists."

Which is the opening the interviewer takes next, because a bank that cannot lend is not much of a bank.

architecture v4
one publish, six independent consumers
swipe the figure sideways, or tap expand for full screen
1/8
sole writer
The ledger service is still the sole writer, exactly as decided in round one.