Low-Latency Software · Technology
12Queues, Ring Buffers and Shared Memory
In 2011 the engineers of a retail exchange published the numbers that had made them redesign their system. Passing an event through a three-stage pipeline built on the standard Java blocking queue took 32 757 nanoseconds per hop on average, and its 99.99th percentile was over four milliseconds; through the ring buffer they built instead, 52 nanoseconds on average and under 8 192 at the 99.99th percentile (Thompson et al., 2011). The work in each stage was the same; the queues were the latency. A trading process is a pipeline of threads (feed handler, book, strategy, gateway, logger), and how data moves between them decides much of its tail. This chapter builds the structures that move it: the single-producer single-consumer ring, the broadcast ring of the “disruptor” pattern, the sequence lock for data where only the latest value matters, and the same structures across processes in shared memory, with what to do when a consumer cannot keep up.
12.1 The single-producer single-consumer ring
The producer writes the slot at its position and then publishes the new position with a release store; the consumer loads the producer’s position with acquire, reads every slot before it, and publishes its own position the same way, handing the slots back (chapter 11). Both operations are wait-free. Lamport proved such a queue correct without locks in 1983; what decades of practice added are the details that make it fast. The two positions live on separate cache lines, so that the producer’s writes do not invalidate the consumer’s position (chapter 3). And each side keeps a private copy of the other’s position and refreshes it only when the ring looks full or empty: while the ring is neither, a message costs one cache line transfer (the slot) instead of three.
Figure 12.2 measures a round trip between two pinned threads: one sends an eight-byte message through a ring, the other sends it back through a second ring. About at the median: four cache-line transfers between cores. The same round trip through a queue protected by a mutex and a condition variable costs about , over a hundred times more, because the waiting thread goes to sleep in the kernel and must be woken. One-way, the ring streams about 250 million eight-byte messages a second on this laptop.
bench_rings.py.12.2 Many readers: the disruptor pattern
Market data is the natural case: one feed handler, several consumers (the book builder, a second strategy, the recorder, a risk monitor), each needing every message. Two policies differ in who waits. With back-pressure (next section) the producer waits for the slowest consumer, which is right for an order path, where nothing may be lost. For market data the producer must not wait: the exchange will not slow down. The build’s broadcast ring therefore never blocks its producer; each slot carries the sequence number of the message it holds, so a reader that has been lapped sees a later sequence than it expected and knows it lost messages, instead of reading a mixture.
12.3 Sequence locks for last-value data
The Linux kernel documentation describes the mechanism as having lockless readers with read-only retry loops and no writer starvation, for data rarely written relative to how often it is read. In trading the pattern suits the top of the book or a reference price, where readers want the latest consistent value, not every intermediate one. Two details make it correct in C++ as well as on x86: the value must be copied with atomic operations (relaxed ones suffice), since a copy that races with the writer’s non-atomic stores is a data race and undefined behaviour (Boehm, 2012); and the reader’s second counter read must be ordered after the copy by an acquire fence. The build’s test runs a writer in a tight loop against a reader that checks two hundred thousand pairs for tearing.
12.4 Shared memory between processes
Processes isolate failures: a strategy that crashes need not take the feed handler with it, and components written in different languages can share a feed. The ring does not care: Figure 12.2 shows the same round trip between two processes as between two threads. What changes is the contract between the processes: the layout must be specified to the byte and versioned, since two binaries built at different times read it; every field must be of fixed size and alignment; and nothing in the region may be a pointer, which means nothing in the other process. The build’s layout is a 192-byte header (magic number, version, slot size, capacity, then the two positions on their own lines) followed by the slots, and its C++ and Rust implementations both reproduce the same byte image of a ring, written independently from the specification by a script.
12.5 Back-pressure, slow consumers and conflation
A queue’s size is a statement about bursts. One Quant Book 4, chapter 8, gives the M/M/1 answer for a stage with Poisson arrivals at rate and exponential service at rate : with utilisation , the stationary number in the system exceeds with probability . A ring that must overflow less than once in a billion messages at therefore needs with , about 197 slots, 256 as a power of two. Market data is burstier than Poisson: the sizing starts from the M/M/1 figure and is checked against recorded bursts (chapter 18).
Proof. The consumer’s backlog grows at messages a second; the producer overwrites the consumer’s next unread slot when the backlog reaches . ∎
A ring of 65 536 slots, a burst at 2 million messages a second and a monitor that reads 1.5 million gives the monitor 131 milliseconds; an opening burst lasting longer laps it. The design answers are to give such consumers conflated views (a sequence lock per instrument instead of every message), to size their rings for measured bursts, or to let them fall behind by design and resynchronise from a snapshot, as the feed handler of chapter 18 does from the exchange.
12.6 Tutorial: rings in threads, processes and two languages
Goal. Build the rings, check the shared layout in C++ and Rust against one fixture, and measure round trips and throughput. End state: Figure 12.2, a throughput figure, and green tests in both languages.
Write and publish. The producer copies the message into its slot and publishes its position with a release store; it reads the consumer’s position only when the ring looks full.
bool try_write(const void* msg, std::uint32_t len) { // producer thread only if (len > h_.slot_size - 4) return false; const std::uint64_t w = word(base_, 64).load(std::memory_order_relaxed); if (w - read_cache_ == h_.capacity) { // looks full: refresh the cached read position once read_cache_ = word(base_, 128).load(std::memory_order_acquire); if (w - read_cache_ == h_.capacity) return false; } std::uint8_t* slot = base_ + kHeader + (w & mask_) * h_.slot_size; std::memcpy(slot, &len, 4); std::memcpy(slot + 4, msg, len); word(base_, 64).store(w + 1, std::memory_order_release); // publish return true; } int try_read(void* out, std::uint32_t cap) { // consumer thread only; length, or -1 when empty const std::uint64_t r = word(base_, 128).load(std::memory_order_relaxed); if (r == write_cache_) { write_cache_ = word(base_, 64).load(std::memory_order_acquire); if (r == write_cache_) return -1; } const std::uint8_t* slot = base_ + kHeader + (r & mask_) * h_.slot_size; std::uint32_t len; std::memcpy(&len, slot, 4); if (len > cap) len = cap; std::memcpy(out, slot + 4, len); word(base_, 128).store(r + 1, std::memory_order_release); // hand the slot back return static_cast<int>(len); }Listing 12.1. The SPSC ring’s producer and consumer. code/firm/ring/cpp/firm_ring.hpp The same in Rust, over the same bytes: a region of atomic words, written through their interior mutability.
pub fn try_write(&mut self, msg: &[u8]) -> bool { if msg.len() > self.slot - 4 { return false; } let w = self.r.atomic(64).load(Ordering::Relaxed); if w - self.read_cache == self.cap { self.read_cache = self.r.atomic(128).load(Ordering::Acquire); if w - self.read_cache == self.cap { return false; } } let at = HEADER + ((w & (self.cap - 1)) as usize) * self.slot; // SAFETY: the producer owns slot w until the release store below; the consumer does not read it before. unsafe { let p = self.r.raw().add(at); std::ptr::copy_nonoverlapping((msg.len() as u32).to_le_bytes().as_ptr(), p, 4); std::ptr::copy_nonoverlapping(msg.as_ptr(), p.add(4), msg.len()); } self.r.atomic(64).store(w + 1, Ordering::Release); true }Listing 12.2. The Rust producer, writing the layout the C++ reader expects. code/firm/ring/rust/src/lib.rs A sequence lock whose value is copied as relaxed atomic words.
void store(const T& v) { // single writer const std::uint64_t s = seq_.load(std::memory_order_relaxed); seq_.store(s + 1, std::memory_order_relaxed); std::atomic_thread_fence(std::memory_order_release); std::uint64_t w[kWords]; std::memcpy(w, &v, sizeof v); for (std::size_t i = 0; i < kWords; ++i) words_[i].store(w[i], std::memory_order_relaxed); seq_.store(s + 2, std::memory_order_release); } T load(std::uint64_t* retries = nullptr) const { for (;;) { const std::uint64_t s1 = seq_.load(std::memory_order_acquire); if (s1 & 1) { if (retries) ++*retries; continue; } std::uint64_t w[kWords]; for (std::size_t i = 0; i < kWords; ++i) w[i] = words_[i].load(std::memory_order_relaxed); std::atomic_thread_fence(std::memory_order_acquire); if (seq_.load(std::memory_order_relaxed) == s1) { T v; std::memcpy(&v, w, sizeof v); return v; } if (retries) ++*retries; } }Listing 12.3. Writer and reader of the sequence lock. code/firm/ring/cpp/firm_ring.hpp - Check the layout.
make_ring_fixture.pywrites the byte image of a ring from the specification; both test suites build the same ring and compare, then consume the fixture. - Measure with
python bench_rings.py: round trips between two threads on two pairs of CPUs, between two processes, and through a locked queue, then one-way throughput.
What to change next. Remove the cached copies of the other side’s position and measure the throughput again; put the two positions on the same cache line and measure the round trip.
12.7 Build: the firm’s rings
Purpose. The transport of every Part IV component: the feed handler publishes normalised events on a broadcast ring (chapter 18), the strategy engine reads them and hands orders to the gateway through an SPSC ring (chapters 20–21), the binary logger drains a ring (chapter 23), and the top of each book is published with a sequence lock.
Interface. C++20 firm::ring: format(base, capacity, slot_size), region_size, Spsc(base) with try_write(msg, len) and try_read(out, cap); Broadcast(base) with write, read(pos, out, len) returning Ok, Empty or Overrun, and head(); SeqLock<T> with store and load; Segment(name, bytes, create). Rust firm_ring: Region, format, Spsc::attach with try_write and try_read, SeqLock2.
Rules. The byte layout above, little-endian, versioned by a magic number and a version; capacities are powers of two; positions on separate cache lines; the SPSC ring reports full and empty, the broadcast ring never waits and reports overruns; no allocation after construction.
Acceptance tests. code/firm/ring/: the fixture’s bytes reproduced and consumed in C++ and Rust; 200 000 messages in order across two threads in both; broadcast overrun detection and resynchronisation; no torn read in 200 000 sequence-lock reads under a writer in a tight loop; 1 000 messages from a child process over shared memory.
Stretch. A broadcast ring with a declared slowest-consumer policy (wait, drop, conflate); a ring in a file on a huge-page filesystem.
Sources and further reading
- M. Thompson, D. Farley, M. Barker, P. Gee and A. Stewart, Disruptor: high performance alternative to bounded queues for exchanging data between concurrent threads, 2011.
- L. Lamport, “Specifying concurrent program modules”, ACM TOPLAS 5(2), 1983.
- H.-J. Boehm, “Can seqlocks get along with programming language memory models?”, MSPC 2012; Linux kernel documentation, “Sequence counters and sequential locks”.
- Linux man page
shm_open(3).
12.8 Exercises
Exercise 12.1 ★
A ring has 1 024 slots; the producer’s position is 70 001 and the consumer’s 69 500. How many messages are waiting, which slot will the producer write next, and how many more can it write before the ring is full?
Exercise 12.2 ★
At 250 million messages a second, how long does a 64-bit position take to wrap? And a 32-bit one?
Exercise 12.3 ★
A stage runs at . What ring size makes an M/M/1 overflow rarer than one in a million?
Exercise 12.4 ★★
Why does each side of the SPSC ring keep a private copy of the other’s position? What does removing the copies cost per message?
Exercise 12.5 ★★
A risk monitor reads the broadcast ring at 400 000 messages a second. A burst of 3 million messages a second lasts 80 milliseconds on a ring of 131 072 slots. Is the monitor lapped, and if so when?
Exercise 12.6 ★★
Why must the sequence lock’s value be copied with atomic operations in C++, although the reader discards any copy taken during a write?
Exercise 12.7 ★★★
Coding. Put the producer’s and consumer’s positions on the same cache line (edit the header offsets in a copy of firm_ring.hpp) and rerun the round-trip and throughput benchmarks. What changes, and why?
Exercise 12.8 ★★★
Find the flaw. “Our shared-memory ring stores a pointer to the head message and a std::string for the instrument name in its header.”
12.9 Problem: Three Strategies, One Feed
Problem 12.1
Weekend problem — sizing a broadcast ring
A feed handler publishes normalised events on a broadcast ring of 65 536 slots of 64 bytes. Three consumers read it: a book builder at up to 4 million events a second, a strategy at 2.5 million, and a monitor at 1.5 million. On an ordinary day the feed averages 300 000 events a second; at the open, bursts reach 2 million a second for up to half a second.
Part I — The ring.
- How many bytes does the ring occupy, and does it fit the second-level cache of chapter 3?
- What is each consumer’s utilisation on an ordinary day?
- Which consumers can fall behind during a burst?
- Why does the producer never wait for them?
Part II — The burst.
- How long before the monitor is lapped (Proposition 12.6)?
- Is it lapped during a half-second burst?
- What ring size would carry the monitor through the burst?
- How much memory would that cost?
Part III — The alternatives.
- The monitor needs only the latest price of each of 5 000 instruments. What does a sequence lock per instrument cost in memory, and what share of updates does the monitor skip during the burst?
- What does the monitor do when it detects an overrun on the ring?
- Why is back-pressure wrong for this producer and right for the order path?
- What does the M/M/1 model say about the strategy’s ring at 300 000 events a second, and why is it optimistic?
Part IV — The verdict.
- State the named result: the time before the monitor is lapped in a burst, and the share of updates a conflating monitor skips.
- What does the measured round trip say about the cost of adding a consumer?
- Why is the ring’s layout versioned?
- What would the published comparison of a blocking queue and a ring predict for a blocking queue here?
- Which placement rule of chapter 4 applies to the feed handler and the book builder?
- How would you record the burst sizes that the ring must carry?
- What must the strategy do after a gap in its input?
- In one sentence: what is a ring buffer’s size a statement about?
12.10 Interview questions
Interview question 12.1 ★ developer
Implement a single-producer single-consumer ring buffer. Which memory orderings does it need?
Interview question 12.2 ★★ developer
Why is a ring buffer faster than a queue protected by a mutex and a condition variable?
Interview question 12.3 ★★ developer
What is a sequence lock, when would you use one, and what are its pitfalls?
Interview question 12.4 ★★ developer
One market-data feed, three consumers of different speeds. Design the distribution.
Interview question 12.5 ★★ developer
What must a data structure in shared memory between two processes avoid?
Interview question 12.6 ★★★ developer, researcher
How would you size the queue in front of a stage whose input is bursty? What does queueing theory give you, and what does it miss?