Research, Data and Risk Platforms · Technology
24Messaging and Event-Driven Architecture
A consumer of the fills topic crashes after applying a batch of messages to the positions and before recording how far it has read. When it restarts it reads the batch again and applies it again, and the firm’s position in one stock is double until someone compares it with the drop copy. Nothing was broken: the messaging system delivered every message at least once, which is exactly what it promised. This chapter is about what messaging systems promise, what they cannot, and how a consumer turns at-least-once delivery into exactly-once effects — with a partitioned log built in process, a position service run under a thousand crashes in three modes, and the outbox that lets a service change its database and announce it without losing either.
24.1 Why systems talk through messages
The platform map of chapter 1 is a set of services: trade capture, positions, risk, P&L, surveillance. They can call each other directly, each one knowing who needs what; or each can publish what happened, and whoever needs it subscribes. The second decouples them in three ways — in knowledge (the publisher does not know its consumers), in time (a consumer can be down and catch up), and in rate (a slow consumer does not slow the publisher) — and makes the flow of events itself a record that can be replayed.
Definition 24.1 (Publish–subscribe)
Publish–subscribe is a messaging pattern in which producers publish messages to named topics without addressing any consumer, and each consumer receives the messages of the topics it subscribes to.
The cost is that the call graph becomes invisible, and failures move from the call (which fails loudly) to the delivery (which can lose or repeat messages quietly). Most of this chapter is about the second.
24.2 Brokers and brokerless transports
Definition 24.2 (Message broker, brokerless messaging)
A message broker is a server that receives messages from producers, stores them, and delivers them to consumers, so that producers and consumers only know the broker. In brokerless messaging producers send directly to consumers — by unicast, multicast or shared memory — through a library in each process, with no server in the path.
The two answer different needs. A broker adds a hop and a server to run, and gives in return storage, replay, and consumers that come and go. Brokerless transports give the lowest latency and no central component: the market-data and order paths of Book 13 are brokerless, with multicast (Book 13, chapter 16) and ring buffers in shared memory, and libraries such as Aeron — "efficient reliable UDP unicast, UDP multicast, and IPC message transport" — provide the reliability on top. Between services that care more about not losing an event than about microseconds — trade capture, positions, risk, the post-trade chain — the firm uses a log.
24.3 The log as the backbone
Definition 24.3 (Partitioned log, consumer offset)
A partitioned log stores each topic as several partitions, each an append-only sequence of messages addressed by their position in it; a message’s partition is chosen from its key, so that messages with the same key keep their order. A consumer offset is the position up to which a consumer has read a partition, recorded by the consumer (committed) so that it can resume there.
Definition 24.4 (Consumer group)
A consumer group is a set of consumer processes that share a topic’s partitions, each partition read by exactly one member at a time, with one committed offset per partition for the whole group; when members join or leave, the partitions are reassigned (a rebalance).
This is the design Kreps, Narkhede and Rao described for Kafka in 2011: a message has no identifier of its own, "each message is addressed by its logical offset in the log". The chapter’s firm.eventlog builds it in process: topics, partitions by a stable hash of the key, append-only segments on disk, retention that drops whole segments, and committed offsets per group. Keying fills by account puts all of an account’s fills in one partition, in order; eight partitions shared by three consumers are assigned three, three and two.
def append(self, topic, key, value, pid=None, seq=None) -> tuple[int, int]:
"""Append to the key's partition; a (pid, seq) already stored is ignored."""
p = self.partition(topic, key)
if pid is not None:
if (pid, seq) in self.seen[(topic, p)]:
return p, -1 # duplicate send: ignored
self.seen[(topic, p)].add((pid, seq))
off = self.end(topic, p)
self.parts[topic][p].append((off, key, value))
if self.dir is not None:
self._write(topic, p, off, key, value)
return p, off
24.4 Delivery guarantees and the exactly-once myth
A consumer does two things with a batch: it processes the messages, and it commits the offset. A crash can fall between them, and the order of the two decides what the crash costs.
Definition 24.5 (At-most-once and at-least-once delivery)
Under at-most-once delivery a message may be lost but is never processed twice: the consumer commits the offset before processing. Under at-least-once delivery a message is never lost but may be processed more than once: the consumer commits after processing, and a crash in between repeats the uncommitted messages.
Definition 24.6 (Exactly-once processing)
Exactly-once processing is the guarantee that each message’s effect on the consumer’s state happens once, whatever the crashes: obtained not from the transport, which can only deliver at least once, but by making the effect idempotent (Book 13, chapter 24) or by storing the state and the offset in one atomic transaction.
The 2011 paper was explicit: "Kafka only guarantees at-least-once delivery", and an application that cares about duplicates "must add its own deduplication logic". Later versions added an idempotent producer and transactions, which give exactly-once between Kafka topics; for any other destination, the documentation says, exactly-once "generally requires cooperation with such systems". The position service is such a destination.
def run_batch(log, group, topic, p, n, apply, mode, crash=None, store=None) -> int:
"""Read up to n messages from the committed offset and process them under `mode`;
crash(stage, i) may raise Crash to simulate the process dying at that point."""
crash = crash or (lambda stage, i: None)
start = store.offset(topic, p) if mode == "atomic" else log.committed(group, topic, p)
batch = log.read(topic, p, start, n)
if not batch:
return 0
nxt = batch[-1][0] + 1
if mode == "at-most-once":
log.commit(group, topic, p, nxt) # commit first: a crash loses the rest
if mode == "atomic":
store.begin()
for i, (_off, _key, value) in enumerate(batch):
crash("process", i)
apply(value)
crash("before-commit", len(batch))
if mode == "at-least-once":
log.commit(group, topic, p, nxt) # commit after: a crash repeats it
if mode == "atomic":
store.commit(topic, p, nxt) # state and offset in one transaction
return len(batch)
The chapter’s day is 10 000 synthetic fills for ten accounts and twenty stocks, keyed by account into four partitions and consumed in batches of 100. Each day ten crashes strike at points drawn uniformly over the consumer’s steps — every message processed, and the moment between processing a batch and committing it — and a hundred days are run in each mode. For every fill we count how many times it reached the positions.
fig_eventlog.py.The first two modes are mirror images (Figure 24.2). A crash inside a batch of 100 falls on average half-way, so committing first loses about 50 fills and committing after repeats about 50; a crash just after processing and before the commit — the hook’s — repeats all 100. In the third mode the positions and the next offset are written in one SQLite transaction: a crash anywhere before the commit rolls both back, the consumer resumes from the offset that matches the state, and no fill is lost or counted twice in 1 000 crashes. The same effect comes from idempotent processing: store each fill’s identifier with the state and skip identifiers already applied.
fig_eventlog.py.24.5 Patterns: outbox, event sourcing, dead letters
The consumer’s problem has a mirror on the producer’s side. The trade-capture service must store a trade in its database and publish the event that announces it. Two separate writes cannot be atomic: write then publish, and a crash in between leaves a trade nobody hears about; publish then write, and consumers act on a trade that does not exist.
Definition 24.7 (Transactional outbox)
A transactional outbox is a table in a service’s own database into which the service writes each event in the same transaction as the change it announces; a separate relay reads the table, publishes the events in order, and marks them sent, so that an event is published if and only if its change was committed — at least once, with a key that lets the log or the consumer drop the repeats.
def write(self, trade: dict, event: dict) -> None:
"""The trade and the event that announces it, committed together or not at all."""
self.db.execute("BEGIN")
self.db.execute("INSERT INTO trades VALUES (?, ?)",
(trade["id"], json.dumps(trade)))
self.db.execute("INSERT INTO outbox (key, body, sent) VALUES (?, ?, 0)",
(event["key"], json.dumps(event)))
self.db.execute("COMMIT")
def relay(self, log: Log, topic: str, crash=None) -> int:
"""Publish unsent events in order, with the outbox number as idempotence key, then
mark each sent; a crash in between resends it later, and the log ignores it."""
n = 0
rows = self.db.execute("SELECT n, key, body FROM outbox WHERE sent = 0 "
"ORDER BY n").fetchall()
for num, key, body in rows:
log.append(topic, key, json.loads(body), pid="outbox", seq=num)
if crash:
crash("published", num)
self.db.execute("UPDATE outbox SET sent = 1 WHERE n = ?", (num,))
n += 1
return n
| design | events lost | phantom events | duplicated events |
|---|---|---|---|
| write the trade, then publish | 22 | 0 | 0 |
| publish, then write the trade | 0 | 22 | 0 |
| outbox, relay without idempotence key | 0 | 0 | 22 |
| outbox, relay with idempotence key | 0 | 0 | 0 |
Table 24.1 shows the three failure modes and the fix. The outbox moves the problem from two systems to one, and what is left — a relay that published and died before marking — is a repeat, which the idempotent producer or the consumer removes.
Definition 24.8 (Event sourcing)
Event sourcing stores a system’s state as the sequence of events that changed it, appended and never edited, and derives every current state, and every past one, by replaying them.
Chapter 21’s trade store is event-sourced, and the log is its natural home: a new downstream system subscribes and replays the topic from its start to build its own state. The price is retention — the log must keep every event, or a snapshot with the offset it covers — and discipline about changing the events’ format.
Definition 24.9 (Dead-letter queue)
A dead-letter queue is a topic to which a consumer sends a message it has failed to process a stated number of times, with the partition and offset it came from, so that one malformed message does not stop the partition behind it.
In a partitioned log a message that always fails is worse than in a queue: the consumer cannot skip it without committing past it, so every message behind it in the partition waits. The dead-letter topic lets the consumer move on and leaves a record to be repaired and replayed — which, for fills, must be done before anyone trusts the positions.
As of September 2026 — Kafka’s delivery semantics
Kafka’s documentation distinguishes at most once (messages may be lost but are never redelivered), at least once (never lost but may be redelivered) and exactly once; a consumer that saves its position before processing gets the first, after processing the second. Since version 0.11.0.0 its producer offers an idempotent delivery option, so that resending does not create duplicate entries in the log, and transactions that write to several partitions atomically; exactly-once delivery to other systems generally requires their cooperation.
24.6 Tutorial: a thousand crashes
Goal. Consume a day of fills under crashes in three modes and publish trades three ways. End state: Figure 24.2 and Table 24.1.
- Log:
firm_eventlog.Log, a topic of four partitions;pl_eventlog.publish(fills()). - One day:
day(fs, mode, crash_steps)for each mode, with the same crash steps; count applications per fill. - A hundred days:
schedules(mode), lost and duplicated fills per 1 000 crashes and the worst position error. - Outbox:
publishing(design)for the four designs. - Dead letters: put a malformed fill in a partition and consume it with
process_with_dlq.
What to change next. Replace the atomic store by idempotent processing (a table of applied fill identifiers) and check that it gives the same result; rebalance the group in the middle of the day.
24.7 Build: the event log
Purpose. Carry the firm’s events between services in order per key, durably, with consumers that can turn at-least-once delivery into exactly-once effects.
Interface. Log (create, append, read, end, retain, commit, committed), assign, run_batch, AtomicStore, Outbox, process_with_dlq.
Rules. Partition by key; never edit a message; every consumer states its mode; state and offset atomic, or processing idempotent, for any consumer whose state matters; every event published through an outbox; a message that keeps failing goes to dead letters with its origin.
Acceptance tests. code/firm/eventlog/tests/: one key in one partition and in order, segments and retention; a resend ignored; the group assignment; the three modes under a crash; the outbox relay’s resend deduplicated; a poison message to dead letters.
Stretch. Consumer-group rebalancing with a generation number that fences out a stale member; compaction by key; a network server.
Sources and further reading
- J. Kreps, N. Narkhede and J. Rao, Kafka: a Distributed Messaging System for Log Processing, NetDB 2011.
- Apache Kafka documentation, Design: message delivery semantics.
- One Quant Book 13, chapters 12, 16 and 24 (ring buffers and slow consumers, multicast, idempotent processing and sequencers).
24.8 Exercises
Exercise 24.1 ★
Why are fills keyed by account rather than sent round-robin to the partitions?
Solution
Solution of Exercise 24.1.
Order is kept only within a partition. Keying by account puts all of an account’s fills in one partition, so the position service applies them in the order they happened; round-robin would spread them over partitions read at different speeds.
Exercise 24.2 ★
A consumer commits after processing and crashes after applying 30 messages of a batch of 100. What happens when it restarts?
Solution
Solution of Exercise 24.2.
It resumes from the last committed offset, the start of the batch, and applies all 100 messages: the 30 already applied are applied twice.
Exercise 24.3 ★
Eight partitions and three consumers in a group: how are they assigned, and what happens if a fourth joins?
Solution
Solution of Exercise 24.3.
Round-robin: three, three and two partitions (0, 3, 6; 1, 4, 7; 2, 5). A fourth member triggers a rebalance to two each; each partition’s new owner resumes from the group’s committed offset, so anything processed but not committed by the old owner is processed again.
Exercise 24.4 ★★
Why does committing first lose about 50 fills per crash, and committing after duplicate about as many?
Solution
Solution of Exercise 24.4.
A crash falls uniformly within a batch of 100, on average half-way. Committing first has already moved the offset past the batch, so the unprocessed half is lost; committing after has not, so the processed half is repeated. The observed 50 323 and 48 652 per 1 000 crashes are that half-batch.
Exercise 24.5 ★★
The position service keeps its state in memory and writes a snapshot every minute. Where should it keep its offset?
Solution
Solution of Exercise 24.5.
In the snapshot, written atomically with the state it describes: on restart it loads the snapshot and resumes from the offset stored in it, replaying everything after. An offset committed elsewhere would not match the state.
Exercise 24.6 ★★
Why does the outbox relay still produce repeats, and who removes them?
Solution
Solution of Exercise 24.6.
A relay that publishes an event and dies before marking the row sent publishes it again after the restart. The idempotent producer (the log ignores a known outbox number) or the consumer (deduplicating on the event’s key) removes it.
Exercise 24.7 ★★★
Coding. Make the at-least-once position service idempotent by storing applied fill identifiers, and run schedules on it.
Solution
Solution of Exercise 24.7.
Store the set of applied fill identifiers with the positions and skip any identifier already present; with the set and the positions updated together, no fill is lost (the offset is still committed after) and none counts twice: zero and zero, as in the atomic mode.
Exercise 24.8 ★★★
Find the flaw. "We turned on exactly-once in the messaging system, so the positions cannot be double-counted any more."
Solution
Solution of Exercise 24.8.
A messaging product’s exactly-once covers what happens inside it — deduplicated writes to its log, transactions across its partitions — not the effect of a message on another system. The position database is another system: its updates and the consumer’s offset must be made atomic, or its processing idempotent, by the consumer.
24.9 Problem: Double for Eleven Minutes
Problem 24.1
Weekend problem — double for eleven minutes
The chapter’s log, fills, crash schedules and trade-capture service.
Part I — The log.
- What does publish–subscribe decouple, and what does it cost?
- When would you choose a broker, and when a brokerless transport?
- How is a message addressed in a partitioned log, and how is its partition chosen?
- What does a consumer group share, and what is a rebalance?
- What does retention drop, and what must a new subscriber do if it has dropped events it needs?
Part II — Crashes.
- What are the two steps of consuming a batch, and where can a crash fall?
- How many crashes happened in each mode, and why fewer in one?
- How many fills were lost or duplicated per 1 000 crashes in each mode?
- What is the worst end-of-day position error in each mode, and the median?
- How does the atomic mode avoid both losses and duplicates?
Part III — Producing.
- Why can the trade and its event not be written atomically without an outbox?
- What does each naive design produce under crashes?
- What does the outbox leave, and how is it removed?
- What is event sourcing, and what does it require of the log?
- What does a dead-letter queue protect?
Part IV — The verdict.
- State the named result: duplicated and lost fills per thousand crashes under at-most-once, at-least-once and idempotent-with-atomic-offset processing, and the maximum position error each produces.
- Which mode should the position service use, and which a monitoring dashboard?
- What does exactly-once in a messaging product cover?
- How would you detect the hook’s error within minutes?
- In one sentence: where does exactly-once come from?
Solution
Solution of Problem 24.1.
- Knowledge, time and rate: publishers do not know consumers, consumers can be down and catch up, a slow consumer does not slow the publisher. It costs a visible call graph, and failures move to delivery.
- A broker when storage, replay and changing consumers matter; brokerless (multicast, shared memory, a library such as Aeron) when latency matters and the endpoints are known.
- By its offset in its partition; the partition from a stable hash of the key.
- The partitions of a topic, one member each, and one committed offset per partition; a rebalance reassigns partitions when members change.
- Whole segments of old messages; a new subscriber that needs them must start from a snapshot that records the offset it covers.
- Processing and committing the offset; a crash can fall inside the processing or between the two.
- 967 in the commit-first mode, 1 000 in the others: committing first never repeats messages, so a day has fewer steps and some scheduled crash points fall after its end.
- Commit first: 50 323 lost, none duplicated; commit after: none lost, 48 652 duplicated; atomic: none of either.
- Up to 5 100 shares (median 3 500) committing first, up to 6 500 (median 3 800) committing after, zero atomic.
- The positions and the next offset are written in one transaction; a crash rolls back both, and the consumer resumes where the state ends.
- They are two systems, the database and the log, with no transaction spanning both.
- Write then publish: 22 events lost; publish then write: 22 phantom events.
- Repeats (22 without idempotence keys), removed by the log’s idempotent producer or by the consumer.
- State stored as the sequence of events that changed it; the log must retain every event or a snapshot with its offset.
- The messages behind a poison message in its partition.
- Named result. Per 1 000 crashes, committing before processing (at most once) loses 50 323 fills and duplicates none; committing after (at least once) loses none and duplicates 48 652; storing positions and offset in one transaction loses and duplicates none. The worst end-of-day position error over a hundred days is 5 100, 6 500 and 0 shares.
- The atomic (or idempotent) mode for positions; a dashboard can commit first, since it only needs the latest state.
- Its own log: deduplicated writes and transactions across its partitions, not the consumer’s database.
- Reconcile the positions with the drop copy continuously, per account and stock, and alert on any difference.
- From the consumer, which makes its effect idempotent or atomic with its offset.
24.10 Interview questions
Interview question 24.1 ★ developer
Explain at-most-once, at-least-once and exactly-once delivery.
Solution
Solution of Interview question 24.1.
At most once: the position is saved before processing, so a crash loses messages. At least once: saved after, so a crash repeats them. Exactly once: an effect, not a transport property, obtained by idempotent processing or by storing state and position atomically.
What the interviewer is looking for: Where the offset is saved.
Interview question 24.2 ★★ developer
Your consumer applies fills to a database. How do you make sure a crash never double-counts a fill?
Solution
Solution of Interview question 24.2.
Commit the offset in the same database transaction as the positions, or store applied fill identifiers with them and skip repeats; never rely on the transport alone.
What the interviewer is looking for: Atomic offset or idempotence.
Interview question 24.3 ★★ developer
How do you publish an event whenever a row is inserted in your database, without ever losing or inventing one?
Solution
Solution of Interview question 24.3.
A transactional outbox: the event is inserted in the same transaction as the row; a relay publishes unsent events in order and marks them, with the outbox number as an idempotence key so that the repeats after a relay crash are dropped.
What the interviewer is looking for: Outbox plus idempotence key.
Interview question 24.4 ★★ developer
How do you keep the messages of one account in order when a topic has many partitions and many consumers?
Solution
Solution of Interview question 24.4.
Key the messages by account, so that one account’s messages go to one partition, and let each partition be read by one consumer of the group at a time.
What the interviewer is looking for: Key to partition; one reader per partition.
Interview question 24.5 ★★★ developer
Design the event backbone that carries trades, fills and positions between the services of a trading firm.
Solution
Solution of Interview question 24.5.
A partitioned log with topics per event type, keyed by account or trade; producers through outboxes; consumers idempotent or atomic with their offsets; retention with snapshots for replay; dead letters monitored; consumer lag and reconciliations against the drop copy as the health measures.
What the interviewer is looking for: Log, outbox, idempotent consumers, reconciliation.