Course Content
System Design Interview
31 sections · 71 lessons
Stock Exchange: reliability, market data and follow-ups
The previous lesson built a fast, deterministic, single-threaded core. This lesson keeps it running when machines fail, gets its output to thousands of subscribers fairly, and covers the follow-up questions that close the interview.
The usual reliability toolkit — retries, replicas answering reads, asynchronous replication, health checks that fail requests over — is mostly unusable here. Retries reorder. Replicas that answer independently break determinism. Failing over mid-stream loses ordering.
Determinism supplies a different toolkit.
The sequenced event log is the whole recovery mechanism
The sequencer assigns a global sequence number to every input event — new orders, cancels, modifies, market-open and market-close events, and any administrative action — and writes them to a replicated journal before the matching engine sees them.
That log is the system's source of truth. The order book, the trades, and the market data are all derived: the pure function of a deterministic engine applied to the log. This is the same relationship as the event-sourced ledger of the digital wallet, and the raw-events archive of ad click aggregation, arriving now with a microsecond budget attached.
Journaling within the budget. Writing to a disk costs about 100 µs and is out of the question synchronously. Two approaches that fit:
- Replicate the sequence to memory on multiple machines over a low-latency network before the engine processes it. Two or three machines acknowledging over a kernel-bypass path costs on the order of 5 to 10 µs, and the event is then safe against any single machine failing. Disk writes happen asynchronously behind that.
- Write to a memory-mapped file and let the operating system flush, accepting that a simultaneous total power failure could lose the last few events, and mitigating with battery-backed or persistent memory.
The first is standard for venues, and the reason is worth stating: durability here means "survives a machine failure", and replication achieves that faster than a disk does.
Hot-warm failover
A standby matching engine consumes the identical sequenced stream and applies the identical code. Because the engine is deterministic, the standby's order book is byte-identical to the primary's at every sequence number, without any state transfer.
Failover is then a matter of deciding who is allowed to publish, not of recovering state. The standby is already correct. It stops discarding its output and starts emitting, from the sequence number after the last one the primary published.
Three details make this work:
- The sequencer, not the engine, is the ordering authority, so it must have its own high-availability arrangement — a small consensus group, or an active-standby pair with a hardware arbiter. If two sequencers ever assign numbers concurrently, ordering is destroyed and no downstream recovery is possible.
- Output deduplication downstream. During failover, both engines may briefly publish. Consumers discard by sequence number, which is why every output event carries one.
- Failover is a halt-and-resume, not a seamless handover. A brief halt of a few hundred milliseconds, with participants notified, is far preferable to a silent overlap in which two engines both match — which would create trades that never happened. Say this plainly: the correct behaviour when the venue is unsure is to stop.
End-of-day reconciliation
After the close, an independent process replays the entire day's sequenced log through a separately maintained implementation and compares its trades, positions, and final books against what the production engine published.
This is the same idea as the batch recomputation in Section 23 (Design an Ad Click Event Aggregation) and the settlement reconciliation in Section 28 (Design a Payment System): an independent implementation over the same input, failing differently. Any discrepancy is a defect, and it is found before positions are cleared and settled the next morning.
Reconciliation also runs outward: trades reported to the clearing house, positions reported by participants, and the venue's own records must all agree. A break in any of them is investigated before the next session opens.
Market data: level 1, level 2, and level 3
The engine's other output is market data. There are three depths, each an order of magnitude more expensive than the last.
Level 1 — best bid, best ask, their sizes, and the last trade. A handful of numbers per instrument, updated when the touch changes. Small, cheap, and enough for most participants.
Level 2 — aggregated depth by price level: every price with the total quantity resting there. Perhaps 20 to 50 levels per side. Updated on every book change, so far higher volume.
Level 3 — every individual order, with its identifier and position in the queue. This is the full book, and it is what allows a participant to compute their own queue position. Highest volume by a wide margin, and not offered by every venue.
The design consequence is to publish tiers separately so a level 1 subscriber does not pay the bandwidth of level 3, and so the tiers can be sized and priced independently.
Multicast, and the arithmetic that requires it
From Requirements and scale: 5,000 subscribers × 20 MB/s of level 2 updates = 100 GB/s if each gets its own connection. With multicast, the publisher sends once and the network fabric replicates to every subscriber — 20 MB/s from the publisher regardless of subscriber count.
Multicast also happens to be fairer. Every subscriber receives the same packet from the same send, so nobody is systematically served before anyone else — whereas iterating over 5,000 unicast connections means subscriber 1 receives the update microseconds before subscriber 5,000, and that ordering would be a persistent, exploitable advantage.
What multicast costs. It is unreliable — packets are not retransmitted — so the design needs sequence numbers on every message, a gap-fill service a subscriber can query over a separate reliable channel when it detects a missing sequence number, and periodic book snapshots so a subscriber joining mid-session, or one that has fallen too far behind to catch up, can resynchronise without replaying the day. Two independent multicast channels carrying the same data over different network paths, with subscribers taking whichever packet arrives first, is a common further mitigation.
Multicast also generally does not cross the public internet, so remote participants receive market data through a gateway that unicasts to them — at which point they are, by physics, behind co-located participants. That is understood and disclosed rather than hidden.
Circuit breakers and halts
When a price moves beyond a threshold in a short window, trading in that instrument pauses. Mechanically: the engine transitions the instrument to a halted state, rejects new orders for it or accepts them without matching depending on the venue's rules, publishes the halt on market data, and after the pause runs an auction to reopen — collecting orders without matching, then computing the single price that maximises executed volume, and matching everything at that price.
Two engineering points worth making. The halt state must be part of the sequenced event stream, so that a replay reproduces the halt exactly at the right sequence number; and the auction is a different matching algorithm from continuous trading, so the engine holds at least two matching modes and a defined transition between them.
Fairness and auditability, as regulators see them
The obligations that shape architecture rather than sitting beside it:
- Deterministic, explainable matching. Every trade reconstructible from the sequenced log. This is a design requirement, not a reporting one.
- Timestamp precision and clock synchronisation. Regulators specify how precisely events must be timestamped and how tightly clocks must be synchronised to a reference. The practical consequence is hardware timestamping at the network card and a disciplined time source, because software clocks drift by more than the entire latency budget.
- Equal access. Equalised cable lengths in the co-location facility, identical gateway hardware, and published rules about who may connect and how.
- Full order and trade records, retained for years, replayable on demand.
Treat the specific numeric requirements as jurisdiction-dependent and subject to change; treat the architectural implications as stable.