System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Stock Exchange deep dive: the single-threaded engine and the latency budget


The first lesson built an order book that matches one order in a few hundred nanoseconds. This lesson answers the question every interviewer asks next — why not run it on many threads? — and then spends the 50-microsecond budget hop by hop.

Here is the contradiction, stated plainly.

For twenty-four case studies, the answer to more load has been to add machines. Stateless services behind a load balancer. Partition by key. Replicate for reads. Process asynchronously. Accept eventual consistency where it is affordable.

The matching engine does none of that. It is one thread, on one core, processing one order at a time, in a loop. That is not a legacy compromise or a simplification for teaching. It is what serious venues build, and it is the correct answer.

The first half of this lesson explains why, and — more usefully — explains what distinguishes this case from all the others, so the lesson generalises instead of reading as an exception.

Reason 1: the work does not partition

Every scale-out design in Parts II to IV rests on one property: the work divides into units that do not need to know about each other. Two users' feeds are independent. Two ad IDs' click counts are independent. Two objects' bytes are independent. Independence is what makes parallelism free.

An order book has no such division. Price-time priority is a total order over every order for that instrument. Deciding whether order 5,001 matches requires knowing the exact state left by order 5,000, which requires knowing 4,999, and so on. There is no subset of the book that can be processed without reference to the rest.

You could put a lock around the book and run several threads. Then only one thread is inside the lock at a time — which is a single-threaded program with extra steps, plus the cost of the lock.

Reason 2: the coordination costs more than the work

The arithmetic is decisive.

Matching one order against an in-memory book is roughly 200 to 500 nanoseconds — a handful of memory references, a linked-list operation, a few writes, with the working set resident in the CPU's cache.

Coordination between cores, using widely cited orders of magnitude: an uncontended atomic operation costs tens of nanoseconds; a cache line moving between cores costs on the order of 100 nanoseconds; a contended lock costs a microsecond or more once a thread is descheduled.

So the coordination for one operation is somewhere between 20% and 200% of the operation itself. Amdahl's law finishes the argument: when the section requiring serialisation is essentially the whole operation, adding threads adds overhead and no throughput. Two threads on one book are reliably slower than one.

Reason 3: one thread is fast enough, and here is the check

At 500 nanoseconds per order, one thread handles 2 million orders per second. Our peak requirement is 200,000 per second across the entire venue, and the busiest single instrument is a fraction of that.

The single thread is ten times faster than required. There is no throughput problem to solve, which means paying any correctness cost for parallelism would be buying something nobody needs.

LMAX, a trading venue, published a description of its architecture — the Disruptor pattern — reporting single-threaded business-logic throughput in the millions of operations per second on commodity hardware. Their published work is the standard public reference for this design, and the figures should be read as their measurements on their workload rather than as a universal constant.

Reason 4: determinism, which is the requirement that settles it

Even if parallelism were free and fast, it would still be wrong.

A single-threaded loop over an ordered input is a pure function: the same sequence of orders produces exactly the same trades, every time, on any machine. That gives four things this venue cannot operate without:

  • Replay for audit. A regulator asks why a trade happened at that price. Replay the sequence and the answer is the trades themselves, not an argument.
  • Failover without divergence. A standby engine consuming the same sequence reaches an identical state. Reliability without slowing down relies on this completely.
  • Testable correctness. Capture a production day, replay against a new build, and diff the outputs. A non-deterministic engine cannot be tested this way at all.
  • Fairness that is observable. Price-time priority requires an unambiguous arrival order. Concurrent processing makes "who was first" depend on scheduling, which is exactly what the rule forbids.

Multi-threaded matching gives up all four. No throughput gain compensates.

So where did the parallelism go?

It did not disappear. It moved either side of the core, and this is the part that reconciles the lesson with the rest of the course.

In front: many gateways, running in parallel, authenticate, decode, and risk-check orders. That work is independent per order, so it parallelises normally. A sequencer then imposes a single total order, stamping each order with a sequence number, and feeds the engine one stream.

Behind: market data publishers, drop-copy feeds, journals, and risk systems all consume the engine's output stream in parallel. Also independent, also normally parallel.

Across instruments: the books for different symbols are genuinely independent, so symbols are partitioned across many matching engines, each single-threaded. That is sharding by key — the same move as every other case study in this course.

The rule that generalises, and the sentence worth remembering:

Parallelise across independent units. Never parallelise within a unit whose semantics require a total order.

Every earlier case study found independent units — users, objects, ad IDs, partitions. This one has them too, at the instrument level. What is new is that inside one unit the semantics forbid concurrency, so the correct shape is a small serial core surrounded by parallel work.

SINGLE THREADPARALLEL EDGEGateway 1decode, validate, riskGateway 2Gateway 3Sequencerassigns order numbersMatching enginein-memory bookPARALLEL AGAINJournalappend before ackMarket data feedClearing + riskHot standbyreplays the sequenceordersone totally ordered streamfillssequencing is what makes the single threaddeterministic — same input order, sameoutput, every timethe standby replays the same sequence,so failover does not need shared state
The core is single-threaded on purpose: determinism is worth more than parallelism when correctness has to be replayable.

The latency budget

With the shape of the core settled, spend the budget. Fifty microseconds is not a slogan. It is a budget, and every component spends from it.

Where the microseconds go

Two configurations of the same path, from the packet arriving at the network card to the acknowledgement leaving it. Figures are illustrative and rounded, to show proportions rather than to be quoted as measurements.

Order-to-acknowledgement latency budget, by hop (microseconds)0204060Kernel network stackKernel bypass + tunedNetwork in + receive pathGateway: decode, validate, riskSequencer: sequence + journalMatching engineResponse encodeSend path + network out
Order-to-acknowledgement latency budget, by hop (microseconds)

Inbound wire time and the receive path are shown merged, as are the send path and outbound wire time — they are tuned together in practice.

Totals: about 65 µs with the operating system's network stack, about 32 µs with kernel bypass and a tuned path.

Two observations that should be said out loud.

The matching engine is the smallest component. Three to five microseconds out of thirty-two to sixty-five — under 10%. Optimising the matcher further would be optimising the wrong thing, which is exactly the kind of judgement the budget exists to produce.

The operating system is the largest single saving. Twenty microseconds of kernel network stack becomes four with bypass. That one change is roughly half the total improvement.

The techniques, and what each actually buys

Kernel bypass. The application talks to the network card directly from user space, so packets skip the kernel's network stack, its copies, and its interrupt handling. Saves roughly 8 µs per direction. Costs: a specialised library and card, application code that must handle the protocol work the kernel was doing, and reduced portability.

Busy polling instead of interrupts. A thread spins reading the receive queue rather than sleeping until woken. Removes the interrupt-delivery and wake-up cost, roughly 1 to 5 µs, and removes the variance, which matters more — a sleeping thread's wake-up time is unpredictable. Costs: a core permanently at 100% utilisation per polling thread, and the power and heat that implies. Venues pay it without hesitation.

No allocation on the hot path. Every object the engine uses is pre-allocated at start-up and reused. In a garbage-collected language this avoids collection pauses, which are catastrophic here — a 10 ms pause is 200 times the entire budget. In a manually managed language it avoids allocator locks and cache misses. This is why matching engines are written either in systems languages or in managed languages used in a deliberately allocation-free style.

Cache-friendly data layout. A main-memory reference costs about 100 ns and an L1 cache hit about 1 ns. Laying order-book data out contiguously so the hot working set stays in cache is worth hundreds of nanoseconds per order. Concretely: arrays of structures with small fixed-size fields rather than pointer-chasing object graphs.

CPU pinning and isolation. Pin the matching thread to a specific core, isolate that core from the operating system's scheduler, and disable frequency scaling and power-saving states on it. Removes context switches and the tens of microseconds a deep sleep state costs to exit.

Co-location. Participants rent rack space in the exchange's own datacentre. Light travels about one foot per nanosecond in fibre, so 100 metres of cable is roughly 0.5 µs each way. Over a metropolitan link of 20 km it is about 100 µs each way — twice the entire internal budget. This is why co-location exists, and why cable lengths inside the facility are commonly equalised between participants so that physical distance does not become an unfair advantage.