Design a market data distribution system
Design a system that consumes an exchange market data feed and distributes it to trading strategies on the same site.
Requirements: tick-to-strategy latency in the low tens of microseconds, deterministic behaviour under bursts, and no strategy able to slow down any other.
Be explicit about which conventional web-architecture instincts you are abandoning and why.
Solution
What changes at this latency budget. Almost every reflex from web architecture is wrong here. A hop through a load balancer is 100µs+. A TCP round trip across a datacentre exceeds the entire budget. A JSON parse per message is out of the question. Garbage collection pauses are unacceptable. Naming these explicitly is most of the answer.
Transport. Exchange feeds are typically UDP multicast: one send reaches every consumer, and one slow consumer cannot apply back-pressure to the publisher. That is a feature — with TCP, a strategy that stalls would slow the feed for everyone, which violates the isolation requirement. The price is that UDP drops, so the protocol carries sequence numbers and you need gap detection plus a recovery channel (a request-retransmit path, or a slower TCP snapshot feed to re-sync from).
In-process distribution. A lock-free single-producer/multi-consumer ring buffer in shared memory. One writer appends; each reader tracks its own position and never blocks the writer. No locks on the hot path, no allocation, fixed-size slots, and a sequence number per slot so a reader can detect it was lapped rather than reading a torn message.
Slow consumers. This is the key decision. The ring must never block the writer, so a reader that falls behind gets overwritten and detects it via the sequence number. Then either it re-syncs from a snapshot, or it is killed. Dropping data for a slow strategy is correct — stale market data has negative value, and a strategy trading on a 5ms-old book is worse than one that knows it is blind. Candidates from a web background almost always propose buffering, and explaining why buffering is wrong here is the strongest single thing you can say.
Parsing and representation. Binary protocols with fixed-offset fields; cast into a struct rather than parse. Pre-allocate everything; no allocation on the hot path. Book updates in flat arrays sized for the instrument universe.
Practicalities. Pin threads to isolated cores, disable frequency scaling, use huge pages, and kernel-bypass networking (Solarflare/ef_vi, DPDK) to skip the kernel stack entirely.
Measurement. Latency is a distribution, and the mean is the least interesting number in it — p99.9 and the maximum are what the business feels. Hardware timestamps at the NIC, histograms rather than averages, and measure in production rather than on a quiet machine.
The follow-up they will ask
A strategy reports it saw a stale book for 40ms during a burst. How would you find out where the time went?