InventDB
All articles Storage

Concurrent B-link trees and one shared page cache

InventDB's id and field indexes are B-link trees on 4 KiB pages. Readers walk the trees without taking a latch, a writer latches only the leaf it changes, and one page cache shared by every tree holds the pages that writes are changing.

An index that many threads use at once has one hard case at its centre: a page splits in two while another thread is reading it. The traditional answer is to latch pages on the way down the tree, so readers and writers wait for each other, most often near the root, where every path begins.

InventDB's indexes use the B-link tree, published by Philip Lehman and S. Bing Yao in 1981, which lets a reader descend without latching anything. The id index of each segment and the three trees kept for every field are all B-link trees. This article explains the protocol, how a writer splits a page, and how one page cache serves every tree in the process.

Pages, keys and the trees that use them

A tree is a file of 4 KiB pages. A page splits when it holds more than 45 keys or when its contents outgrow the page. Each page carries a CRC32 checksum and is encrypted with AES-256-GCM before it is written, so a page read from disk is checked and decrypted before anything uses it.

Keys are encoded so that comparing their bytes orders them by value. Numbers and dates use an 8-byte encoding whose byte order matches numeric order. Strings are stored so that byte comparison orders them by value, and the record id is appended to every key so that equal values sort by id. A value longer than 1,024 bytes, or one whose key would exceed 1,800 bytes, is stored in the index as its SHA-256 hash, and a query that meets such an entry tests its condition against the record itself.

The same format serves the id index of each segment, which maps a record id to the record's position in its file, and the three trees of every field: the sorted value tree, the counts tree of distinct values with their counts and sums, and the id-to-value tree.

A B-link tree adds two fields to every page. The right link points to the page's sibling to the right on the same level. The high key is the largest key the page may hold. The protocol keeps one rule: any key larger than a page's high key lives somewhere to the right of that page and can be reached by following right links.

A reader uses that rule instead of latches. It loads the root's position from an atomic variable, and at each page it compares the key it wants with the page's high key. If the key is larger and the page has a right link, the reader moves right. Otherwise it searches the page and either descends to a child or, at a leaf, has its answer.

The rule earns its place when a page splits under a reader. Suppose a reader looking for key k30 has read the parent and is about to read leaf A, and in that moment a writer splits A, moving k23 and everything above it into a new page B. The reader reads A, finds that k30 is above A's new high key k22, follows A's right link to B and finds its key there. It never waited, and it never needed to know that a split happened.

A leaf split: new leaf B takes the upper keys with A's old right link and high key, and a reader looking for k30 at leaf A moves right to B Parent gains separator k22 last Reader seeking k30 Leaf A k00 ... k22 high key now k22 right link New leaf B k23 ... k45 high key k90, from A right link Leaf C k91 ... unchanged 1 B is written with A's old right link and high key; nothing points to it yet. 2 A is rewritten with high key k22 and a right link to B. 3 A reader that reaches A looking for k30 sees k30 above k22 and moves right.
The split writes the new page before it rewrites the old one, so a reader sees either the old leaf with every key or the new pair joined by a right link.

Right links also make batch lookups cheaper. A batch sorts its keys and answers each one from the current leaf when it can, following up to four right links before it falls back to a fresh descent from the root, so a run of nearby keys costs a handful of page reads.

A writer latches one page

A single insert starts like a read. The writer descends without latches and remembers the parent page it passed at each level. At the leaf it takes that page's write latch, a reader-writer lock looked up by the page's offset, and checks the high key again under the latch, moving right if a split has moved its key's place since it looked. Then it inserts the key. Most inserts end there, holding exactly one latch for the time it takes to change one page.

If the leaf now holds more than 45 keys, or its contents no longer fit, the writer splits it in this order:

  1. It allocates a new page and, in its own copy of the leaf, moves the upper half of the keys into the new page.
  2. It gives the new page the old page's right link and high key, and writes it. No page links to it yet, so no reader can reach it.
  3. It writes the old page with its remaining keys, its last remaining key as its high key, and a right link to the new page.
  4. It inserts a separator for the new page into the parent it recorded on the way down, splitting the parent the same way if the parent is full.

Until step 3 is written, readers find every key in the old page. After it, they find the upper half by moving right. No reader can miss a key that exists.

Bulk loads sort each batch's keys for a tree and keep one leaf latched as a cursor, adding keys in order until it splits, so a run of neighbouring keys costs one descent instead of one each.

One page cache for every tree

Trees do not own their page memory. The engine keeps one page cache for the whole process, and every tree is given a numeric id when it is opened. Pages are stored in a sharded concurrent hash map keyed by tree id and page offset, so the pages of different trees never collide, and a per-tree list of cached offsets lets the engine drop one tree's pages without scanning the rest. A hit takes a shared lock on one shard and records the access with an atomic store, and no lock covers the whole cache.

The cache holds the trees of each type's active segment, which are the trees that writes change. Those trees are pinned, so their pages are not evicted while the segment takes writes. A write changes a page in the cache and marks it dirty, and the checkpoint, which the server runs every five seconds, writes the dirty pages to disk. Holding the live segment's pages in memory is what keeps bulk loading fast, because without it every index change would be a separate disk write.

When a segment seals, its trees are unpinned and their pages become eligible for eviction. When the unpinned pages exceed the budget, the cache sorts the clean ones by last access and evicts the oldest, a tenth of the budget at a time. A dirty page is never evicted; it has to be written first. The budget is a share of the machine: an eighth of its memory, and no more than a third of the memory that is available when the server starts.

A reader of an active tree that misses the cache reads the page from disk, and a 4 KiB positioned read is not atomic against a writer rewriting the same page. If the bytes fail their checksum, the reader takes the writer's copy of the page from the cache, where the writer placed it before writing to disk. That copy may be one change behind, and right links and high keys exist to make a slightly older page safe to read.

Sealed segments are read from their files

The cache deliberately stops at the active segment. An instance's data is usually larger than its memory, so a cache of all of it would hold a small, shifting fraction while costing work on every access to keep it consistent. Sealed segments never change, which makes their files cheap to read directly.

The engine reads sealed files through read-only memory maps. A map is created on the first read of a file and dropped for good on the first write to it, so a file that is still being appended to never keeps one. Reading a page then needs no system call: the bytes come straight from the operating system's page cache, and the engine checks and decrypts them. The operating system decides which bytes stay in memory and gives them back under pressure. Decoded column arrays and whole documents have caches of their own, described in Columnar aggregates inside a JSON document database.

What it costs

  • Reads take no latches but can take extra steps. A burst of concurrent splits can add right-link hops, so a read is lock-free without being wait-free. In practice the number of extra hops is close to zero.
  • Sealed pages are decrypted on every read. No decrypted copy of a sealed page is kept, so each read pays for a checksum and an AES-256-GCM decryption. We spend that CPU to leave memory for documents, columns and the operating system's page cache, which matters because on our 10-million-record benchmark database 84% of the bytes on disk are index trees.
  • Pinned pages ignore the budget. The active segment has three trees per field, and its pages stay in memory while it takes writes, so a type with many fields holds more memory during a large import.
  • Neighbouring writes take turns. Inserts of neighbouring keys, such as increasing timestamps, meet on the rightmost leaf of a value tree and wait for each other's latch. Inserts spread across the key space proceed in parallel.

How we test it

Concurrency faults in an index rarely show up in a single-threaded test, so the tree's own suite runs the protocol under contention:

  • 20 reader threads each read 500 keys and must get the exact stored value every time.
  • 10 writer threads each insert 100 distinct keys, after which the tree must hold exactly 1,000 entries, each readable.
  • 10 threads update the same 100 keys at once, after which the tree must still hold exactly 100 keys.
  • A mixed workload of 80% reads and 20% inserts runs on 10 threads with 500 operations each, and a stress run of 20 threads inserts 250 keys each, after which every key is read back.
  • Separate tests cover encrypted pages and a tree that is closed and reopened.

For debugging, the engine also has a structural check mode, off in normal operation, that reads the neighbouring page on every leaf write and records where an invariant broke. Above the trees, each run of our 3,511-query SQL battery compares the rows every query returns.

Where you meet it

There is no setting for any of this. Every id and field index on every type in InventDB Serverless and InventDB SOAR is a B-link tree, and every instance has one page cache sized from its own memory. What you notice is the behaviour: an import and the queries over the same type run at the same time, and the queries do not wait for the import's page splits.

# an import running
POST /api/shop/orders/bulk
[ { "order_id": "A-10231", "status": "open", "total": 129.50 }, ... ]

# a query over the same type, at the same time
POST /sql
{ "sql": "SELECT COUNT(*) AS n FROM shop.orders WHERE status = 'open'" }

Each document in the import updates the value, counts and id-to-value trees of every field it has, latching one leaf at a time in each tree. The count reads the same trees without latching any of them, and it reflects the records the import has written so far.