Skip to content
Pivot is in early development and is not production ready. See Project status.

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:

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.

A table split into morsels, each run through the whole pipeline A table with columns A, B and C is split into four morsels of rows. Morsel 1 is taken and passed through a pipeline of three operators, decode, then filter, then aggregate, and on into the result, while morsels 2 to 4 wait. A bracket under the pipeline notes that every operator works on the same morsel while it is still in the L1 cache. Morsel-at-a-time execution one worker A B C Morsel 1 Morsel 2 Morsel 3 Morsel 4 Decode Filter Aggregate Result the same morsel, still in L1

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:

Pivot's estimated cache access breakdown on ClickBench A stacked bar shows the reported counter-based estimates as percentages of all loads and stores: 97.9 percent hit L1, approximately 1.09 percent are served by L2 after an L1 miss, and 1.01 percent go beyond L2. The shares sum to 100 percent, with 98.99 percent attributed to L1 or L2. L2's share is calculated as 98.99 minus 97.9, using rounded reported rates. These are estimates from hardware-counter ratios, not an exact classification of individual loads and stores. Beyond L2 can mean the shared system cache or DRAM; DRAM-only accesses were not measurable on this VM. Next to each share is the typical latency of that layer: about 1 nanosecond for L1, about 4 nanoseconds for L2, and 10 to 120 nanoseconds beyond L2. Pivot cache access breakdown Estimated share of all loads and stores 0% 100% L1 cache 97.9% ~1 ns L2 cache 1.09% ~4 ns Beyond L2 (“cold”) 1.01% ~10–120 ns
ClickBench, 43 queries, AWS c8g.4xlarge. Pivot main 253313c1, PGO; approximately 162 billion loads and stores. Percentages are approximate shares derived from Arm PMU counters.

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.

Memory bandwidth versus compute across Graviton generations A line chart across three AWS Graviton generations, normalized to Graviton2. Memory bandwidth grows from 1 times on Graviton2 (c6g, 205 gigabytes per second peak) to 1.5 times on Graviton3 (c7g, 307 gigabytes per second) and 2.6 times on Graviton4 (c8g, 538 gigabytes per second). Compute performance grows from 1 times to 1.25 times and then 1.6 times over the same generations. Memory bandwidth vs compute relative to Graviton2 1× 2× 3× Graviton2 c6g Graviton3 c7g Graviton4 c8g 205 GB/s 307 GB/s 538 GB/s Memory bandwidth 2.6× Compute 1.6×
Gains AWS reports between generations.

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:

DRAM bandwidth over a warm pass of TPC-H Three area charts, one per engine, show DRAM bandwidth from 0 to 200 gigabytes per second over a warm pass of the 22 TPC-H queries, with a dashed line at the 190 gigabytes per second measured ceiling. Each engine's pass is stretched to the same width; Pivot's takes 20.7 seconds, DuckDB's 31.2 seconds and ClickHouse's 52.0 seconds. Pivot's trace stays close to the ceiling for most of the pass. ClickHouse's and DuckDB's stay mostly below 100 gigabytes per second. DRAM bandwidth over the warm pass read + write, 5 ms samples All 22 queries back to back, each pass stretched to the same width Pivot 20.7 s DuckDB 31.2 s ClickHouse 52.0 s 190 GB/s measured ceiling Start of pass End of pass
TPC-H SF100, 22 queries, warm pass after one cold pass. AMD EPYC 9124 (16 cores), 8 × DDR5-4800, 125 GB. System-wide DRAM read and write CAS counts × 64 B, sampled every 5 ms with perf. Pivot reads Parquet files in place; ClickHouse reads native MergeTree tables and DuckDB its native database.

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:

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.

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_activity
FROM users
WHERE country = 'US'
ORDER BY last_activity DESC
LIMIT 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:

Late materialization speedup Two bar charts comparing query time with late materialization on and off. On a cold run the query takes 2.03 seconds with late materialization and 8.92 seconds without, 4.4 times faster. On a hot run it takes 33.4 milliseconds with late materialization and 181.4 milliseconds without, 5.4 times faster. Each chart uses its own time scale. Cold run Late materialization on 4.4× faster 2.03 s Late materialization off 8.92 s Hot run Late materialization on 5.4× faster 33.4 ms Late materialization off 181.4 ms

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.

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.

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:

Buffer ring versus jemalloc for 2 MB allocations Two bar charts comparing Pivot's buffer ring with jemalloc when handling 2 MB buffers. Acquiring and releasing one buffer costs 46 nanoseconds through the ring and 249 nanoseconds through jemalloc, 5.4 times cheaper. Writing to 2 GB of freshly acquired buffers for the first time takes 3.5 milliseconds with the ring and 413 milliseconds with jemalloc, 118 times faster, because the ring's memory is already faulted in. Each chart uses its own scale. Acquire + release per 2 MB buffer Pivot buffer ring 5.4× cheaper 46 ns jemalloc 249 ns First write 2 GB of fresh buffers Pivot buffer ring 118× faster 3.5 ms jemalloc 413 ms
100 rounds of acquiring 1000 × 2 MB buffers, touching every page, and releasing them. AWS c8g.4xlarge.

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:

TPC-H q21 with and without huge pages Two bar charts for TPC-H SF100 query 21 run hot with the buffer ring on huge pages and on ordinary pages. Query time is 2.47 seconds with huge pages and 3.13 seconds without, 21 percent faster. Data TLB page walks are 1.2 billion with huge pages and 5.8 billion without, 4.8 times fewer. Query time hot run Huge pages on 21% faster 2.47 s Huge pages off 3.13 s Data TLB page walks Huge pages on 4.8× fewer 1.2 billion Huge pages off 5.8 billion
TPC-H SF100 q21, hot, with transparent huge pages enabled and disabled in the kernel. AWS c8g.4xlarge.

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:

The same query on two cores, without and with work stealing Two timelines. Without stealing, worker 1 has nine batches and worker 2 has three, so worker 2 finishes early and sits idle while worker 1 works through the rest, and the query finishes when worker 1 does. With stealing, once worker 2 runs out of its own batches it steals the three batches at the oldest end of worker 1 queue. Worker 1 keeps running its newest batches, both workers finish after six batches, and the query completes a third sooner. Without stealing time Worker 1 Worker 2 idle query done With stealing time Worker 1 oldest end of its queue Worker 2 query done
Twelve equal batches on two cores. Without stealing, Worker 2 finishes its share and idles while Worker 1 works through the rest. With stealing, Worker 2 takes batches from the oldest end of Worker 1's queue, the ones Worker 1 would have reached last and that have long left its cache, so Worker 1 keeps running its newest, still-hot batches and both cores stay busy until the query is done. Stealing stays within a NUMA node, since a batch's memory lives on the node that produced it.

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.

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.

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.