Why is Pivot fast?
Although database performance comes from many different optimization strategies and architectural decisions, most fall into a couple of main categories. To fully maximize performance, pivot tries to utilize optimizations in each one of these categories:
It utilizes the hardware well.
Section titled “It utilizes the hardware well.”Utilizing the CPU cache
Section titled “Utilizing the CPU cache”Modern analytical engines are often bound by memory access rather than computation. A simple integer comparison may take roughly 1 CPU cycle (<1 ns), while fetching data from DRAM can take on the order of 70–100 ns, or hundreds of CPU cycles.
This makes effective use of the CPU cache hierarchy critical for analytical workloads. Data in the CPU’s L1 cache can typically be accessed in ~1 ns, L2 in ~3–5 ns, and L3 in ~10–20 ns, compared with ~70–100 ns for DRAM. Keeping data close to the CPU can therefore improve a query’s performance by one to two orders of magnitude.
Pivot’s execution strategy is inspired by Morsel-Driven Parallelism and is designed around this principle. It processes data in small batches (i.e. morsels) sized to fit within the CPU’s L1 cache, then performs as much of the query pipeline as possible on each morsel while its data remains hot in the cache.
By processing data in a cache-aware manner, Pivot significantly reduces the number of accesses to RAM, thereby reducing the time the CPU spends stalled waiting for memory:
Utilizing memory throughput
Section titled “Utilizing memory throughput”Over the past few years, memory bandwidth has been growing significantly faster than single-core CPU performance. As a result, many engines that were designed around the compute-to-bandwidth ratios of earlier hardware can’t fully take advantage of modern DRAM speeds.
To make full use of this growing memory bandwidth, Pivot combines many techniques, including:
- Aggressive software prefetching, so data is requested before it’s needed.
- Avoiding instruction dependency chains (like pointers inside pointers) when accessing memory, so the CPU can issue many memory requests at once instead of waiting for each one to return before starting the next.
- Loops that mix different kinds of work, so the CPU’s execution units stay busy in parallel while memory requests are still in flight (for example, shifting bits while multiplying numbers and fetching memory all in parallel).
Thanks to these and other optimizations, Pivot can use a larger share of the machine’s memory bandwidth, significantly speeding up memory-bound queries:
It processes less data
Section titled “It processes less data”The fastest way to process data is to avoid processing it in the first place. Pivot uses several strategies to minimize the amount of data that needs to be read and processed for each query. Some of these include:
Soft ordering data
Section titled “Soft ordering data”Parquet stores min/max statistics for each column in every row group. When evaluating a query, query engines can use these statistics to determine that certain row groups cannot contain matching rows and skip them entirely.
For example, consider the following query:
SELECT * FROM sales WHERE user_id=1377;Any row group whose user_id range does not include 1377 can be skipped. A row group with a minimum user_id of 10 and a maximum of 200, for example, cannot possibly contain the requested user and therefore does not need to be read.
While row group pruning is common among engines that query Parquet, Pivot takes greater advantage of it by allowing users to soft-sort their tables:
CREATE TABLE sales ( user_id BIGINT, amount DOUBLE, ts TIMESTAMP) WITH ( sort_by = 'user_id');Soft sorting organizes rows within generated Parquet files to make row group statistics more selective for columns commonly used by a workload. Pivot also prioritizes compacting overlapping files, helping keep similar values clustered together and further improving pruning efficiency.
For example, a sales table can be configured to soft-sort data by user_id. Without sorting, values may be scattered across row groups, producing heavily overlapping min/max ranges:
| Row Group | Min User ID | Max User ID |
|---|---|---|
| 1 | 10 | 2000 |
| 2 | 5 | 1500 |
| 3 | 1270 | 1800 |
| 4 | 300 | 1350 |
A query filtering for user_id = 1377 would need to read row groups 1, 2, and 3 because all three ranges could contain the requested value.
After soft-sorting by user_id, the ranges become much more selective:
| Row Group | Min User ID | Max User ID |
|---|---|---|
| 1 | 5 | 500 |
| 2 | 501 | 1000 |
| 3 | 1001 | 1500 |
| 4 | 1501 | 2000 |
Now, the same query only needs to read row group 3, allowing Pivot to prune the other three entirely.
Soft sorting alone, however, does not guarantee that data remains well clustered over time. As new Parquet files are written, their value ranges can begin to overlap:
| File | Min User ID | Max User ID |
|---|---|---|
| A | 1 | 1000 |
| B | 2 | 1300 |
| C | 400 | 1400 |
| D | 900 | 1700 |
Even if each file is internally sorted, a query for user_id = 1200 would need to inspect files B, C, and D.
Pivot’s compaction strategy prioritizes files with overlapping ranges. By compacting and re-sorting these files, the data can be reorganized into less-overlapping ranges:
| File | Min User ID | Max User ID |
|---|---|---|
| A | 1 | 500 |
| B | 501 | 1000 |
| C | 1001 | 1500 |
| D | 1501 | 2000 |
The same query for user_id = 1200 can now prune three of the four files and only read file C.
Background compaction does this in bulk. Files whose sort-key ranges overlap are rewritten together, six at a time, once six such files have accumulated: sorting six files at once and cutting the result back into full-size files narrows each file’s range about sixfold per rewrite, so a row reaches its final place in a few rewrites however the data arrived. Files whose ranges do not overlap are never rewritten, so a table whose sort key arrives in order, such as a timestamp, costs no compaction work beyond merging small files. COMPACT table FINAL finishes the job on demand, rewriting until no two files’ ranges overlap.
Together, soft sorting and overlap-aware compaction keep similar values clustered as the table evolves, improving pruning at both the file and row-group level.
This is conceptually similar to ClickHouse’s sparse primary-key index / table order by definition, but applied to Parquet files and open table formats.
Late materialization
Section titled “Late materialization”Queries often return many columns even though only a small subset is needed to determine which rows belong in the final result. Reading and decoding all selected columns upfront can therefore waste significant I/O and CPU on rows that will eventually be discarded.
Pivot uses late materialization to delay reading columns until they are actually needed.
For example, consider a query that returns detailed information about the 10 most recently active users in US:
SELECT user_id, name, email, country, city, company, job_title, profile_image_url, bio, last_activityFROM usersWHERE country = 'US'ORDER BY last_activity DESCLIMIT 10;Although the query selects many columns, only country and last_activity are needed to determine which 10 rows should be returned. Instead of immediately reading and decoding every selected column, Pivot can first process the columns required for filtering and Top-K selection.
Once the 10 matching rows have been identified, the remaining columns (such as name, email, city, company, job_title, profile_image_url, and bio) can be materialized only for the surviving rows.
For wide tables or highly selective queries, this can significantly reduce the amount of data read, decoded, and processed. For example, on ClickBench Q23 (a query of a similar shape) running on a c8g.4xlarge instance with late materialization on and off:
Parquet-optimized pruning
Section titled “Parquet-optimized pruning”Pivot is built specifically to perform well on Parquet. After laying the data out for efficient pruning and filtering, Pivot tries to leverage a couple of Parquet-specific strategies to prune data as much as possible.
For example, in a query such as:
SELECT * FROM events WHERE country = 'IL';Pivot can eliminate work at several levels:
- Dictionary-based pruning - if the country column is dictionary-encoded, Pivot can first check whether “IL” exists in the dictionary of that column. If it does not, every data page referencing that dictionary can be skipped without decompressing or decoding it.
- Page-level pruning - even when a row group contains matching rows, many of its pages may contain none of the rows that need to be read. Pivot uses the selected row positions to identify these pages and skips their decompression and decoding entirely.
- Selective decoding - for pages that do need to be read, Pivot passes the surviving row positions down into the Parquet reader. Where the encoding allows, values belonging to unwanted rows can be skipped without being decoded at all.
- Bloom-filter pruning (coming soon) - Pivot will also use Parquet Bloom filters to quickly determine that a row group cannot contain a requested value, allowing it to be skipped entirely.
It works well with the OS
Section titled “It works well with the OS”Linux is an amazing operating system. It can run a wide variety of workloads with great performance and stability. However, like many low-level systems, it is difficult to build a “generalist” system or algorithm that performs optimally across very different workloads.
Because of this, some databases have explored building specialized operating systems (e.g. DBOS) or bypassing parts of the OS entirely, giving the database more direct control over how resources are managed.
Rather than forcing Pivot users to move to an unfamiliar or less mature operating system, Pivot tries to get the best of both worlds: taking as much control as possible over scheduling, I/O, and memory management from the OS, while still benefiting from the stability, ecosystem, and extensive set of libraries that Linux has to offer.
Memory: custom memory management
Section titled “Memory: custom memory management”Instead of relying on a general-purpose allocator, Pivot uses its own block-based allocator for most of its memory needs. Every operation in Pivot (decompression, decoding, and so on) works over a collection of fixed-size, non-contiguous 2 MB blocks rather than a single contiguous buffer.
This predictable and simple memory model provides several advantages over a traditional general-purpose allocator:
Fewer page faults - Since all of Pivot’s buffers are in a constant size – they can be mmaped and faulted in once at startup, then kept for the lifetime of the process. Acquiring a buffer is simply a pop from a per-worker free list, and releasing one is a push. There are no size classes to select, no fragmentation to manage, and no mmap or munmap calls on the allocation path like other allocators have:
Huge pages can be utilized. Because every buffer is exactly 2 MB, Pivot asks Linux to back the buffer ring with transparent 2 MB huge pages instead of ordinary 4 KB pages. When a query’s working set spans more pages than the TLB can hold, each miss costs a multi-level walk through the page tables. With huge pages, one TLB entry covers 512 times as much memory, so operators that jump around large structures, such as the hash tables behind joins and aggregations, hit in the dTLB far more often and walk the page tables far less:
Scheduling: a thread per core architecture with custom internal scheduling:
Section titled “Scheduling: a thread per core architecture with custom internal scheduling:”In the world of database performance, context switches are one of the worst enemies of efficient execution. When multiple tasks or threads compete for the same CPU core, constantly switching between them adds overhead and, more importantly, disrupts the CPU caches that each task relies on. The result is more time spent switching between work and less time actually doing it.
To minimize this overhead, Pivot uses a thread-per-core architecture, designed so that context switches almost never occur during query execution. Each CPU core runs a dedicated Pivot “worker” responsible for scheduling, executing, and orchestrating the work assigned to that core.
To make the most of the CPU caches, each worker also prefers to schedule operations whose data is already “hot” in cache over operations that would require bringing new data in. For example, after decoding a batch of data, a worker will prefer to immediately run a filter over that batch rather than start decoding a new one.
When a worker runs out of work, it can also steal work from a neighboring worker. This allows idle cores to help busy ones, keeping CPU utilization high while still preserving the cache locality of the thread-per-core model:
When multiple queries run in parallel, Pivot workers schedule work across them in a way that balances CPU cache locality with fairness. Workers try to balance between executing operations whose data is already hot in cache, and ensuring that no single query consumes all available CPU time and causes others to starve.
IO: direct IO with IO uring:
Section titled “IO: direct IO with IO uring:”As a complement to its thread-per-core architecture, Pivot uses io_uring for I/O. This allows each worker to efficiently dispatch and manage many concurrent I/O operations while reducing the syscall overhead associated with reads and writes.
Pivot also avoids relying on the operating system’s page cache for caching disk data, as some other databases such as ClickHouse do by default. Instead, it uses direct I/O together with its own dedicated caching system.
Managing the cache directly gives Pivot more control over what memory is used for. Rather than having disk pages cached independently by the operating system, Pivot can make eviction decisions across different types of cached data, for example choosing between compressed data read from disk and decompressed or otherwise processed data that is more expensive to reconstruct.
Because these decisions are made by the database itself, Pivot can prioritize cached objects based on their actual value to query execution, rather than relying on the more general-purpose caching policies of the operating system.
It’s locally optimized with LLMs
Section titled “It’s locally optimized with LLMs”Pivots initial codebase is built entirely without LLMs, to ensure our initial architecture is at the highest level, exactly according to our intentions.
But LLMs are incredibly powerful at finding “local” optimizations given the right starting conditions, and will oftentimes find improvements in places humans wouldn’t. By ensuring our initial architecture removes “noise”, for example page-faults or context-switching, LLMs can reason more easily about what is going on and find very interesting improvements.
To find local optimizations with LLMs, we give LLMs access to servers such as CherryServers (which allows all perf events to the level of seeing memory bandwith) and AWS servers. We give the LLM a series of benchmarks, and ask it to use perf and other tools to find single improvements
that yield above x% in given benchmarks.
If the LLM succeeds (spoiler: it usually does!) we go over the general idea of what allowed the improvement; most times, the idea behind the improvement is sound, while the implementation is not. LLMs are powerful in this regard; they can (much quicker than humans) confirm whether an idea is worth delving into, even if their solution is messy.
As a real example in Pivot, Pivot used to spread row groups across workers for decoding, one row group per worker. This meant that on large machines, the last row group being decoded held up all the other workers running the query. The LLM realized this was as a bottleneck, and implemented a fix that allowed parallelization in decoding a single row group by using mutexes and some other “unsavory” code. This sped up benchmarks by over 5% on large machines, which led us to understand that our current lack of parallelization in decoding a single row group was indeed a real problem that could also affect workloads out in the wild. The implementation however was nothing close to what we wanted, and could create serious issues. We then designed the “sound” way to do this - create another operator with a stealable channel of “decode jobs” - which both fit well within our architecture and actually sped us up even more than the LLMs implementation.
As per our contribution guide, code is never merged without a human deeply understanding it; LLMs could simply “cheat” and benchmax, which of course does not help Pivot in the long run.
