In a vectorized engine, how does the batch size affect CPU cache behaviour, and why not use very large batches?
answer
- it should still be in cache when the next operator runs
- too small and you are back to per-call overhead
- too big and every operator boundary hits RAM
- multiply by columns and by threads
- a wide projection shrinks the usable batch
basics
~20 sBatches are sized so that the vectors an operator touches stay resident in L1 or L2 cache between operators. Too small and per-call overhead dominates; too large and intermediates spill to slower memory, so each operator re-reads its input from RAM.
solid answer
~50 sThe point of a batch is to amortize per-call overhead while keeping the working set inside fast cache. If a batch is a handful of values, you are back to the per-call dispatch cost that vectorization exists to remove. If a batch is a million values, the column vectors and every intermediate an operator produces no longer fit in L1 or L2; the filter writes its output to RAM, and the next operator reads it back from RAM — you have effectively materialized each step through main memory and thrown away the locality that made the chain fast. The sweet spot is a batch whose *live vectors across the operator chain* fit comfortably in the per-core cache, which in practice lands around a thousand to a few thousand values. Wide rows, many projected columns and per-thread copies all shrink that number, and total memory is batch size times columns times concurrent threads.
code
text · 8 linesworking set per batch (single pipeline):
batch_size x projected_columns x bytes_per_value
2048 x 4 cols x 8 B = 64 KB -> L1/L2 resident
2048 x 40 cols x 8 B = 640 KB -> spills past L2
65536 x 40 cols x 8 B = 20.9 MB -> RAM round trip per operator
multiply again by concurrent threads for node footprintgo deeper
Know the shape of the idea: engines process a chunk of rows at a time, and the chunk is sized to fit in fast CPU cache rather than being as large as possible.
Be able to argue both failure directions — per-call overhead at tiny batches, memory round trips at huge ones — and to compute a working set as batch size times columns times value width.
Recognize when a profile points at working-set size rather than arithmetic, and reach for the real fixes first: narrow the projection, avoid needless materialization, confirm the query is CPU-bound at all.
Own the node-level footprint story. Batch size multiplied by columns and by concurrent pipelines is what turns a memory limit into a concurrency limit, and that shapes how you size and isolate compute.
## What batch size is trading off Batch-at-a-time execution exists to sit between two bad extremes. Row-at-a-time pays an operator call, a type dispatch and a branch per value. Full-column-at-a-time — materializing an entire column's worth of intermediate results per operator — amortizes all of that away, but produces intermediates far larger than any cache, so every operator boundary becomes a round trip to main memory. The batch is the compromise: large enough that per-call overhead is negligible when divided across it, small enough that the values stay in fast cache from one operator to the next. ## The cache arithmetic Modern cores have a small L1 data cache (tens of kilobytes), a mid-sized private L2 (hundreds of kilobytes to a few megabytes), and a larger shared L3 that is several times slower than L1. The quantity that must fit is not one vector but **all the vectors live at once** in the operator chain: the input columns being scanned, the intermediate results (comparison masks, computed expressions, selection vectors), and the output columns being built. A rough calculation makes the shape obvious. Take a batch of 2048 values, eight projected columns of 8 bytes each, plus a couple of intermediates: ``` 2048 values x 8 columns x 8 bytes = 131,072 bytes (~128 KB) + masks and selection vector = a few KB ``` That already exceeds a typical L1 and is living in L2. Multiply the batch by a hundred and the same working set is measured in megabytes, sits in L3 or RAM, and each operator now streams its input in from memory rather than reading it out of the cache the previous operator just left it in. The instruction count is identical; the memory stalls are not, and stalls are the dominant cost in a scan-heavy engine. Going the other way is equally bad in a different currency. With a batch of 8 values, the fixed cost of entering the operator, resolving the type and setting up the loop is divided across 8 values instead of 2000, and the loop is too short to amortize its own prologue or to let the prefetcher get going. You have paid for the complexity of a vectorized engine and are getting interpreter-like throughput. ## Other pressures on the number Cache is the primary constraint but not the only one: - **Concurrency and memory footprint.** Each concurrent thread or pipeline holds its own batches, and per-core private caches are not shared. Total resident memory scales with batch size × columns × threads, and on a highly concurrent warehouse node that multiplication is what causes memory pressure, not the batch alone. - **Column width and count.** A projection of forty columns, or wide decimal and long string values, blows the same batch size past cache far sooner than a projection of four narrow integers. This is one reason `SELECT *` hurts an execution engine as well as the scan. - **Time to first row.** Bigger batches delay the first output, which matters for interactive queries with `LIMIT`. - **Selectivity.** After a very selective filter, a large batch may yield only a handful of surviving rows, so downstream operators run with poor lane utilisation regardless of nominal batch size. ## Why the number rarely needs tuning Engines pick a default in the region of one to a few thousand values because that is where the curve is flat: performance rises steeply from tiny batches, plateaus across a wide band, and degrades slowly once the working set leaves cache. Being anywhere in the plateau is fine, which is why this is a *mechanism* question rather than a tuning knob you are expected to twiddle. If you find yourself hand-tuning it, the more likely problems are an over-wide projection, an operator that materializes far more than it should, or a workload that is not CPU-bound at all. ## Connecting it to the plan A useful diagnostic instinct: when an engine's profile shows high cycles-per-instruction or a large share of time in memory stalls during a filter-and-project pipeline, the cause is usually the working set, not the arithmetic. Narrowing the projection reduces bytes per batch directly, which both shrinks the working set and reduces the bytes read from storage — the same fix helping at two levels. ## How to answer Say the goal in one line — keep the live vectors of the operator chain in per-core cache so the next operator reads them from cache, not RAM — then give both failure directions, then note the multiplication by columns and threads that determines the real footprint.
- If bigger batches amortize overhead better, why does throughput fall off past a certain size?Because the amortization curve flattens quickly while the cache curve falls off a cliff. Once the live vectors exceed the private cache, each operator writes its output to RAM and the next reads it back, so you pay a memory round trip per operator boundary. The saved call overhead is a rounding error next to that.
- How does a wide projection interact with batch size?Directly: the working set is batch size times the number of projected columns times value width. Selecting forty columns instead of four makes the same nominal batch ten times larger in bytes, pushing it out of cache. That is an execution-level reason to project only what you need, on top of the storage-level reason of reading fewer bytes.
saying these in an interview costs you the question
- Says larger batches are always faster because overhead is amortized
- Ignores that the working set is batch size times columns times threads
- Thinks batch size affects only memory usage, not cache locality
- Treats batch size as a routine tuning knob for slow queries
- Confuses batch size with the storage block or row-group size