Quantitative Finance · Book 15 · Technology

Research, Data and Risk Platforms

Research, Data and Risk Platforms · Technology

8Python at Scale

A researcher’s first version of a signal walks the day’s events one at a time: for each record, read its side and quantity, update an exponentially weighted average, store it. On the laptop it costs 960 nanoseconds an event, three minutes for a day of two hundred million messages. The same computation written as one array expression costs 14 nanoseconds an event, and streamed through the file in chunks of ten thousand events it never holds more than a third of a megabyte, whatever the length of the history. Nothing about the mathematics changed. What changed is how much work the interpreter does per number, how many bytes each representation of the data costs, and what is in memory at once. This chapter measures those three costs on tick-store data and turns the lessons into the chunked-pipeline tool of the research platform (firm.pyscale).

8.1 What the interpreter costs

Python is the research platform’s language because it is quick to write and has the libraries; the price is that every operation written in Python is executed by an interpreter that manipulates objects.

Definition 8.1 (Interpreter overhead)

The interpreter overhead of an operation is the work the Python interpreter does around the arithmetic itself: decoding the bytecode, looking up names and attributes, checking types, creating and destroying the objects that hold intermediate values, and counting references. For a scalar operation it is tens to hundreds of times the cost of the arithmetic.

Listing 8.1 is the first version of the signal, twice. The first loop reads the memory-mapped records of the tick store one at a time; every access to a field of rec[i] creates a record object and then a scalar object. The second converts the two columns to Python lists once and then loops over plain integers.

def loop_records(rec) -> np.ndarray:
    """Pure Python over the memory-mapped records: every field access builds objects."""
    out = np.empty(len(rec))
    y = None
    for i in range(len(rec)):
        x = (1.0 if rec[i]["side"] == 66 else -1.0) * float(rec[i]["qty"])
        y = x if y is None else (1.0 - ALPHA) * y + ALPHA * x
        out[i] = y
    return out


def loop_lists(rec) -> np.ndarray:
    """Pure Python over lists made once from the arrays: the loop pays for the interpreter."""
    sides, qtys = rec["side"].tolist(), rec["qty"].tolist()
    out = [0.0] * len(sides)
    y = None
    for i in range(len(sides)):
        x = (1.0 if sides[i] == 66 else -1.0) * qtys[i]
        y = x if y is None else (1.0 - ALPHA) * y + ALPHA * x
        out[i] = y
    return np.array(out)
Listing 8.1. The signal as a Python loop: over the mapped records, and over lists made once from the columns. code/platforms/08-python-at-scale/python/pl_pyscale.py

Figure 8.1 measures them: 960 nanoseconds an event over the records, 145 over lists. The difference is the cost of building the objects that the record access creates; what remains, 150 nanoseconds for about ten bytecode operations, is the interpreter itself. Neither is a bug; both are the price of the language.

The cost per event of one signal (the exponentially weighted average of signed volume) computed four ways on tick-store records: a Python loop over the mapped records, a Python loop over lists, one array expression, and the same expression streamed in chunks of 65 536 events. Measured on a laptop (Intel Core Ultra 7 155H) under WSL2, one thread, machine otherwise idle. Data: bench_pyscale.py.
Figure 8.1. The cost per event of one signal (the exponentially weighted average of signed volume) computed four ways on tick-store records: a Python loop over the mapped records, a Python loop over lists, one array expression, and the same expression streamed in chunks of 65 536 events. Measured on a laptop (Intel Core Ultra 7 155H) under WSL2, one thread, machine otherwise idle. Data: bench_pyscale.py.

8.2 Thinking in arrays

Definition 8.2 (Array programming)

Array programming expresses a computation as operations on whole arrays — elementwise arithmetic, reductions, cumulative sums, filters, sorts — so that the interpreter is entered once per operation and the loop over elements runs in compiled code over contiguous memory.

The only hard part is recursion. An exponentially weighted average yt=(1−α) yt−1+α xty_t = (1-\alpha)\,y_{t-1} + \alpha\,x_t looks inherently sequential, but it is a linear filter: its output is the input convolved with α(1−α)k\alpha(1-\alpha)^k, and a filter routine computes it in compiled code in one call (Listing 8.2). On the laptop that is 14 nanoseconds an event, 70 times faster than the loop over records and 11 times faster than the loop over lists, and the answers are identical to the last bit: the filter performs the same multiplications and additions in the same order.

def ewma_kernel(alpha: float):
    """y_t = (1 - alpha) y_{t-1} + alpha x_t, the recursion done by a linear filter over the
    whole chunk; the state is the last y, the filter's start for the next chunk."""
    b, a = [alpha], [1.0, -(1.0 - alpha)]

    def k(x, last):
        x = np.asarray(x, dtype=float)
        if last is None:
            last = x[0]
        y, _ = lfilter(b, a, x, zi=[(1.0 - alpha) * last])
        return y, float(y[-1])
    return k
Listing 8.2. The recursive average as a linear filter; the filter’s internal state for the first element comes from the last output of the previous chunk. code/firm/pyscale/firm_pyscale.py

Method 8.3 (From a loop to arrays)

Look for the pattern in the loop: an elementwise map (arithmetic, where); a running total (cumsum); a window (differences of cumulative sums, or a sliding view); a linear recursion (a filter); a group-by (sort, then reduce by segments); an as-of lookup (searchsorted, chapter 4). What does not fit — a book that changes with every event, a state machine with many branches — is a case for compiled code (chapter 9), not for a slower loop.

A compiler that turns the Python loop into machine code at run time is the other route, the just-in-time compilation of One Quant Book 13 (chapter 10); the series’ environment does not install one, and chapter 9 takes the explicit route of writing the kernel in C++ or Rust.

8.3 Memory: objects, arrays and copies

The second cost is space. A Python object carries a header (a type pointer and a reference count) before its value, and a list of them is an array of pointers to objects scattered in memory. An array of machine numbers is the numbers.

The same column of one million symbols (fifty distinct four-character tickers) in four representations: pandas objects (a Python string per row), a pandas categorical (integer codes and fifty strings), an Arrow string column (offsets and bytes), and a numpy fixed-width unicode array. Data: fig_pyscale.py.
Figure 8.2. The same column of one million symbols (fifty distinct four-character tickers) in four representations: pandas objects (a Python string per row), a pandas categorical (integer codes and fifty strings), an Arrow string column (offsets and bytes), and a numpy fixed-width unicode array. Data: fig_pyscale.py.

Figure 8.2 measures one column of symbols. As Python objects it takes 61 MB, a string object per row of which a few bytes are the characters; as a categorical, 1.0 MB, a small integer per row and fifty strings once; as Arrow strings, 8 MB, an offset and four bytes per row; as a numpy fixed-width unicode array, 16 MB, four characters of four bytes each. The categorical is 61 times smaller than the objects: the dictionary encoding of chapter 3, in memory. Every column of a research dataframe that holds repeated strings — symbols, venues, conditions, sides — is worth the conversion.

Remark 8.4 (Copies)

Array programming has its own trap: every intermediate expression allocates a new array. np.where(side == B, 1.0, -1.0) * qty on a day of two hundred million events allocates a boolean array, two float arrays and a product, several gigabytes of temporaries for a result of 1.6 GB. The remedies are to compute in place where the library allows it, to choose narrow types, and, above all, not to hold the whole day at once.

8.4 Chunks and out-of-core loops

Definition 8.5 (Chunked processing, peak resident memory)

Chunked processing runs a computation over a dataset one slice at a time, carrying between slices only the state the computation needs to continue exactly (a filter’s last value, a window’s tail), so that the memory it uses is set by the chunk size and not by the size of the data. The peak resident memory of a process is the largest amount of physical memory it occupied at any moment; it, not the average, decides whether a job fits a machine.

The tick store’s flat files are memory-mapped (chapter 4), so a slice of the file is a view: nothing is read until it is touched, and nothing needs to be released. Listing 8.3 streams the signal through the file and keeps only its running maximum; the filter’s state carries across chunk boundaries, so the result equals the whole-array result.

def stream_max(rec, size: int) -> float:
    """The largest |EWMA| of the history, streamed: only one chunk and the filter's state
    are ever in memory."""
    k, state, best = ewma_kernel(ALPHA), None, 0.0
    for c in chunks(rec, size):
        y, state = k(signed(c), state)
        best = max(best, float(np.abs(y).max()))
    return best
Listing 8.3. The signal streamed through the mapped file: one chunk’s temporaries and the filter’s state are all that is ever in memory. code/platforms/08-python-at-scale/python/pl_pyscale.py

Proposition 8.6 (Chunked equals whole)

If a computation over a sequence can be written as yt=f(yt−1,xt)y_t = f(y_{t-1}, x_t), or as a function of a window of the last ww inputs, then processing the sequence in chunks, with the last output (or the last w−1w-1 inputs) of each chunk passed to the next, gives the same outputs as processing it whole.

Proof. By induction over chunks: the first output of a chunk is ff applied to the carried state and its first input, which is what the whole computation applies at that position; within a chunk the recursion is the same. ∎

Peak memory of the streamed signal on a history of two million events against the chunk size, each point labelled with its time in seconds; the last point computes the whole array at once. Memory grows with the chunk; time is lowest for chunks of a hundred thousand events and within ten per cent of it from ten thousand to the whole array, once the per-chunk overhead is amortised. Measured on a laptop (Intel Core Ultra 7 155H) under WSL2, one thread, machine otherwise idle; memory counted by tracemalloc, the mapped input not included. Data: bench_pyscale.py.
Figure 8.3. Peak memory of the streamed signal on a history of two million events against the chunk size, each point labelled with its time in seconds; the last point computes the whole array at once. Memory grows with the chunk; time is lowest for chunks of a hundred thousand events and within ten per cent of it from ten thousand to the whole array, once the per-chunk overhead is amortised. Measured on a laptop (Intel Core Ultra 7 155H) under WSL2, one thread, machine otherwise idle; memory counted by tracemalloc, the mapped input not included. Data: bench_pyscale.py.

Figure 8.3 is the trade-off. Chunks of a thousand events pay the interpreter once per thousand and take 0.048 s for the two million events; chunks of ten thousand take 0.025 s with a peak of 0.31 MB, and chunks of a hundred thousand, the fastest here, 0.023 s with 2.5 MB; the whole array at once takes 0.028 s with a peak of 32 MB, 16 bytes an event. At 16 bytes an event, the series’ rule of at most 1.5 GB a process allows about 94 million events at once, a small fraction of a day of a large feed; streamed, the only limit is time.

8.5 Processes, threads and the interpreter lock

Definition 8.7 (Global interpreter lock)

The global interpreter lock of the standard Python interpreter allows one thread at a time to execute Python bytecode; threads run Python code concurrently but not in parallel, while compiled code that releases the lock (much of numpy, input and output) runs in parallel with them.

Two threads each running a pure-Python loop of five million additions take 0.20 s against 0.11 s for one: twice the work in twice the time, no parallelism. Two threads each sorting a numpy array of eight million values take 0.094 s against 0.086 s for one: the sort releases the lock, and the two run largely in parallel. Python code that must use several cores therefore uses several processes.

As of September 2026 — An interpreter without the lock

PEP 703 (status Final, consulted September 2026) added a build configuration of CPython that runs Python code without the global interpreter lock; starting with release 3.13 CPython supports this free-threaded build, and the official macOS and Windows installers optionally install it. The lock remains the default. The series’ environment (Python 3.10) has the lock.

Definition 8.8 (Process-based parallelism, serialisation cost)

Process-based parallelism runs work in several operating-system processes, each with its own interpreter and lock, so that Python code runs on several cores at once. Its serialisation cost is the time spent converting data to bytes and back to move them between processes (pickling, copying through a pipe), which shared memory or files mapped by both processes avoid.

Moving a 64 MB array to a second process and summing it there takes 0.14 s pickled through a pipe and 0.03 s through shared memory, the copy into the shared block included: the serialisation dominates the pipe, and a sum of eight million numbers is negligible beside either. The research platform’s rule follows: move work to the data, not data to the workers. Give each process a slice of a memory-mapped file, a list of Parquet partitions or a shared block, and have it return a small result.

Moving a 64 MB array to a second process (pickled through a pipe, or placed in shared memory), and running a pure-Python loop and a numpy sort on one thread and on two (each thread doing the whole task). Measured on a laptop (Intel Core Ultra 7 155H) under WSL2, machine otherwise idle. Data: bench_pyscale.py.
Figure 8.4. Moving a 64 MB array to a second process (pickled through a pipe, or placed in shared memory), and running a pure-Python loop and a numpy sort on one thread and on two (each thread doing the whole task). Measured on a laptop (Intel Core Ultra 7 155H) under WSL2, machine otherwise idle. Data: bench_pyscale.py.

8.6 Tutorial: four ways to compute one signal

Goal. Compute one signal over tick-store records by a loop, by arrays and in chunks, show that the answers agree, and measure time and memory. End state: Figures 8.1, 8.2 and 8.3.

  1. History: pl_pyscale.history(path, n) writes nn version-1 events as a tick-store flat file.
  2. Agree: agree() runs the two loops (Listing 8.1), the array expression and the chunked version on 20 000 events: the differences are zero.
  3. Measure: bench_pyscale.py times the four methods, the chunk sizes and the parallel tasks.
  4. Memory: symbol_memory() for the four representations of a column.
  5. Windows: firm_pyscale.rolling_sum_kernel(w) carries a window’s tail across chunks; check it against a whole cumulative sum.

What to change next. Stream the signal from the week’s Parquet partitions of chapter 4 instead of the flat file; split the history between two processes by halves and combine their results, and find what state crosses the boundary.

8.7 Build: the chunked pipeline

Purpose. The research platform’s way to run any per-event computation over histories larger than memory, and to measure what a computation costs in time and memory before it is scheduled (chapter 13).

Interface. chunks(array, size); run(source, kernel, state) -> (outputs, state), a kernel being kernel(chunk, state) -> (output, state); ewma_kernel(alpha), rolling_sum_kernel(window); probe(fn, *args) -> (result, seconds, peak_bytes).

Rules. Chunks are views, never copies; the carried state is the minimal one that makes the chunked result equal the whole result, and a test proves it for every kernel; memory is measured as a peak, with the method stated.

Acceptance tests. code/firm/pyscale/tests/: chunks cover the array exactly; the EWMA and rolling-sum kernels give the whole result for every chunk size, including one and the whole length; the probe reports a peak that grows with the allocation.

Stretch. A pool of worker processes each given a range of a mapped file (at most two here, by the resource rule); a kernel for a trade-sign filter that needs a state machine; peak resident memory from the operating system alongside tracemalloc’s count.

Sources and further reading

  • S. Gross, PEP 703 – Making the Global Interpreter Lock Optional in CPython (Final); Python documentation, Python support for free threading.
  • Python documentation, multiprocessing.shared_memory.

8.8 Exercises

Exercise 8.1 ★

At 960, 145 and 14 nanoseconds an event, how long does each method take for a day of two hundred million events?

Solution

Solution of Exercise 8.1.

960×2×108960 \times 2\times10^8 ns =192= 192 s, about three minutes; 145145 ns: 29 s; 1414 ns: 2.7 s.

Exercise 8.2 ★

Why does the loop over records cost about six times the loop over lists, although both run the same arithmetic?

Solution

Solution of Exercise 8.2.

Each field access on a mapped record builds a record object and then a scalar object, with their reference counting and destruction; the loop over lists reads integers that already exist. The arithmetic is the same; the objects are not.

Exercise 8.3 ★

A dataframe has a million rows and three string columns of fifty distinct values each. Estimate its string memory as objects and as categoricals from Figure 8.2.

Solution

Solution of Exercise 8.3.

About 3×61=1833 \times 61 = 183 MB as objects and 3×1.0=33 \times 1.0 = 3 MB as categoricals.

Exercise 8.4 ★★

Write the rolling sum of the last ww values as an array expression, and state the state a chunked version must carry.

Solution

Solution of Exercise 8.4.

With c0=0c_0 = 0 and ct=∑s≤txsc_t = \sum_{s\le t} x_s (a cumulative sum), the sum of the last ww values ending at tt is ct−ct−wc_t - c_{t-w} (with ct−wc_{t-w} taken as 0 when t<wt < w). A chunked version carries the last w−1w-1 values of the previous chunk (or, equivalently, the last ww cumulative sums), so that the first w−1w-1 windows of a chunk are complete.

Exercise 8.5 ★★

Two threads each ran the same pure-Python loop in twice the time of one. What would two processes have done, and what would they have cost if each needed a 64 MB input?

Solution

Solution of Exercise 8.5.

Two processes would each have their own interpreter and lock and run the two loops in parallel, in about the time of one. With a 64 MB input each, pickled through a pipe, each would first pay about 0.14 s of serialisation, more than the loop itself; through shared memory or a mapped file, about 0.03 s, or nothing if the processes map the file themselves.

Exercise 8.6 ★★

At 16 bytes an event for the whole-array computation, how many events fit a 1.5 GB budget at once, and what fraction of a day of the options feed’s planned 311 billion messages (chapter 2) is that?

Solution

Solution of Exercise 8.6.

1.5×109/16≈941.5\times10^9 / 16 \approx 94 million events, about 0.03% of a day of 311 billion messages: a large feed cannot be held whole in one process; it must be streamed.

Exercise 8.7 ★★★

Coding. Check rolling_sum_kernel(50) against a whole-array cumulative-sum version on 100 000 random values for chunk sizes 1, 7, 50, 1 000 and the whole length. What is the largest difference?

Solution

Solution of Exercise 8.7.

The differences are of the order of 10−1310^{-13} for every chunk size: rounding, since the chunked version’s cumulative sums start from each chunk’s carried tail rather than from the first value; they are exact in the mathematical sense.

Exercise 8.8 ★★★

Find the flaw. “Our feature job is slow, so we split the day into eight parts with a thread pool of eight threads, each running the Python feature loop on its part, and concatenate the results.”

Solution

Solution of Exercise 8.8.

A pure-Python loop holds the interpreter lock: eight threads run it one at a time, in about the time of one thread over the whole day, plus switching. Use eight processes, each mapping its part of the file, or rewrite the loop in arrays; and if the feature is recursive across the day, carry its state across the parts rather than restarting it at each boundary.

8.9 Problem: Three Minutes or Three Seconds

Problem 8.1

Weekend problem — a signal over a year of events

The chapter’s signal (an exponentially weighted average of signed volume), its four implementations and the measurements of bench_pyscale.py.

Part I — The interpreter.

  1. What does the interpreter do per event in the loop over records that it does not do in the array expression?
  2. How many nanoseconds an event does each of the four methods cost?
  3. What is the speed-up of the array expression over each loop?
  4. Why are the loop’s and the filter’s answers identical to the last bit?
  5. What kinds of loop cannot be written as array operations, and where do they go?

Part II — Memory.

  1. How many bytes a row does a column of symbols cost as objects, as a categorical, as Arrow strings and as fixed-width unicode?
  2. Why is the categorical smallest?
  3. What does the whole-array computation of the signal cost in memory per event, and where does it come from?
  4. How many events does a 1.5 GB budget hold at once?
  5. What is peak resident memory, and why does it, not the average, decide whether a job fits?

Part III — Chunks and cores.

  1. Why does the chunked result equal the whole result, and what state is carried?
  2. Which chunk size is fastest here, and what is its peak memory?
  3. Why are chunks of a thousand events slower?
  4. What do two threads gain on a Python loop and on a numpy sort?
  5. What does it cost to move 64 MB to another process by pickling, and by shared memory?

Part IV — The verdict.

  1. State the named result: the speed-up of the array expression over the loops, the chunk size that minimises time and its peak memory, and the largest history that fits the budget whole.
  2. How long would a year of 250 days of two hundred million events take by each method, on one core?
  3. How would you use eight cores for that year without paying the serialisation cost?
  4. Which of the chapter’s lessons would a free-threaded interpreter change, and which not?
  5. In one sentence: what makes Python fast enough for tick data?
Solution

Solution of Problem 8.1.

  1. It builds and destroys record and scalar objects, looks up names, dispatches every operation on types, and counts references, for each of the ten or so bytecode operations of an iteration.
  2. About 960 (loop over records), 145 (loop over lists), 14 (array expression) and 15 (chunked).
  3. 70 times over the records, 11 times over the lists.
  4. The filter performs the recursion’s multiplications and additions in the same order as the loop.
  5. Loops whose state changes shape with each event (a book, a state machine with many branches); they go to compiled code (chapter 9).
  6. About 61, 1, 8 and 16 bytes.
  7. Fifty distinct symbols stored once, and a one-byte code per row: dictionary encoding in memory.
  8. About 16 bytes an event: the signed quantities and the output, each 8 bytes, plus short-lived temporaries.
  9. About 94 million.
  10. The largest physical memory the process used at any moment; a job is killed at its peak, not at its average.
  11. The filter’s last output starts the next chunk’s recursion (Proposition 8.6).
  12. Chunks of 100 000 events, with a peak of 2.5 MB; chunks of 10 000 are within ten per cent at 0.31 MB.
  13. The interpreter’s overhead is paid once per chunk, two thousand times for two million events.
  14. On a Python loop, nothing (twice the work in twice the time); on a numpy sort, most of a second core, since the sort releases the lock (0.094 s for two against 0.086 s for one).
  15. About 0.14 s pickled, 0.03 s through shared memory, copy included.
  16. Named result. The array expression is 70 times faster than a Python loop over the mapped records and 11 times faster than one over lists (14 against 960 and 145 nanoseconds an event); streamed in chunks of 10 000 events it takes about the same time (0.025 s against 0.028 s whole, 0.023 s at the fastest chunk of 100 000) with a peak of 0.31 MB instead of 32 MB, so a 1.5 GB budget holds about 94 million events whole and any history streamed.
  17. 250×2×108=5×1010250 \times 2\times10^8 = 5\times10^{10} events: about 13 hours by the loop over records, 2 hours over lists, 11 minutes by arrays.
  18. Eight processes, each mapping its share of the days’ files and streaming them in chunks, each returning a small result; no data through pipes.
  19. It would let threads run Python loops in parallel, so the threaded version of exercise 8 would scale; it would not change the interpreter’s overhead per event, the memory of objects, or the value of streaming.
  20. Keeping the loops out of the interpreter and the whole history out of memory.

8.10 Interview questions

Interview question 8.1 ★ researcher, developer

Why is a Python loop over a numpy array slow, and what do you do instead?

Solution

Solution of Interview question 8.1.

Each iteration runs in the interpreter, building and checking objects for every element: hundreds of nanoseconds against a few for the arithmetic. Express it as array operations (cumulative sums, filters, searchsorted), or compile the loop.

What the interviewer is looking for: Interpreter overhead per element; the array idioms.

Interview question 8.2 ★★ researcher, developer

How would you compute an exponential moving average over a billion values in Python?

Solution

Solution of Interview question 8.2.

As a linear filter (one call in compiled code), streamed over the data in chunks from a memory-mapped file, with the last output of each chunk as the next chunk’s initial condition; memory stays at one chunk.

What the interviewer is looking for: Recursion as a filter; chunking with carried state.

Interview question 8.3 ★★ developer

What is the global interpreter lock, and when does it not matter?

Solution

Solution of Interview question 8.3.

A lock that lets one thread at a time run Python bytecode. It does not matter for input and output, or for compiled code that releases it (much of numpy), and it is sidestepped by processes; a free-threaded build of recent CPython versions removes it.

What the interviewer is looking for: Concurrency against parallelism; what releases the lock.

Interview question 8.4 ★★ researcher

Your dataframe uses 40 GB of memory. Where would you look first?

Solution

Solution of Interview question 8.4.

Object columns (strings as Python objects: convert to categoricals or Arrow strings), wide types (float64 or int64 where 32 or 8 bits suffice), copies held by intermediate variables, and whether the whole history needs to be in memory at all.

What the interviewer is looking for: Representations first, then streaming.

Interview question 8.5 ★★★ developer, researcher

Design a job that computes features over five years of tick data on one 32-core machine with 128 GB of memory.

Solution

Solution of Interview question 8.5.

Partition the work by days (or instruments) into a queue; run about 30 worker processes, each mapping or scanning its partition from the store and streaming it in chunks with carried state, writing results as files; size chunks so that 30 workers stay well within 128 GB; measure time and peak memory on one day before launching all.

What the interviewer is looking for: Processes over threads, data-local work, a measured memory budget.

Terms defined in this chapter

See all 2333 terms in the glossary