In this article
A single InventDB instance keeps a whole database on one machine. That is simple and fast, and it is what our customers run today. It also has three limits: one machine's disks decide how large the database can grow, one machine's processors decide how much work it can do, and when that machine stops, the database stops with it.
InventDB Cluster takes those limits away. It spreads one database across many machines in three availability zones, keeps three copies of every piece of data in three different zones, and moves work and data between machines without anyone placing anything by hand. It is designed so that applications see the same database they already know: the same SQL, the same REST API, the same MCP tools and the same row rules. A single instance is simply a cluster of one machine.
Why a cluster
Larger enterprises bring larger databases: billions of rows rather than millions, more people and systems writing at once, and a business that cannot stop when one server does. On one machine every one of those needs eventually meets a wall, and buying a bigger machine only moves the wall a little further away.
A cluster grows by adding machines instead. Data, and the work on it, is divided into many small pieces that can live on any machine, so more machines mean more room and more work done at the same time. And because every piece is kept three times, in three separate data centres, losing a machine, or even a whole data centre, does not lose data and does not stop the database.
How a cell is built
We call one cluster a cell. A cell lives in one cloud region and spreads over three availability zones, which are separate data centres with their own power and network. It is made of four kinds of part, and every machine in it runs the same program, with a setting at start-up that decides which parts it plays.
- Data nodes hold the data and run queries next to it, so most of the work happens where the rows already are.
- Gateways take every request, check who is asking, and send the request to the data nodes that hold its rows. They also merge the partial answers that come back. They keep nothing that must survive a restart.
- The metadata group is three members, one in each zone, that keep the cell's map: which machines belong to the cell and where every piece of data lives. They agree on every change to that map before it takes effect.
- The object tier is the cell's own object store, which we wrote ourselves. It keeps the cell's backups, each object on three machines, one in each zone. Nothing in the data path is a third-party service.
How a request flows
Every table is divided into tablets, each holding a share of its rows. A tablet keeps three copies, one in each zone, and one of the three is its leader. Leaders are spread evenly over the data nodes, so every machine leads some tablets and follows others.
A write goes from the gateway to the leader of the tablet that holds its row. The leader applies it, records it in its log and sends the entry to the other two copies. The write is committed once the leader and at least one other copy have it on disk, which means two of the three zones, and only then does the client get its answer. A transaction that touches several tablets is committed across their leaders in two phases, and each step is itself written to the replicated log, so a crash at any point leaves a state that the machines still running can finish.
A read is answered by the tablet's leader alone, with no round trip to the other copies. While the leader holds its lease, no other copy can be leading, so the answer is current.
A query is planned once at the gateway and sent to every tablet it needs at the same time. Filters, sorting and the first stage of grouping run on the data nodes next to the rows, and the gateway merges what comes back. A query over a billion rows therefore runs on thirty machines at once rather than on one.
When something fails
The cell is designed to choose correct answers over availability whenever the two conflict. Within the failures it is built to survive, it keeps answering. Beyond them, the affected data waits until it is repaired, and the cell never returns a wrong, stale or partial answer in the meantime. The table shows what happens in each case, and which of them we have measured so far.
| What fails | What the cell does | Where we stand |
|---|---|---|
| A machine stops | For each tablet it led, another copy is elected leader. Writes to those tablets are refused until then, and every other tablet carries on. | Measured: a new leader in 5.4 seconds, and none of the 53,144 writes acknowledged during the run were lost. |
| The machine comes back | It rejoins as a follower and catches up from the other copies. | Measured: afterwards all three copies were byte-identical at the same point in the log. |
| A whole zone goes dark | Every tablet still has two of its three copies, so writes can still commit, and leaders move to the two zones that are left. The cell is sized so two zones can carry the load. | By design, not yet measured. |
| The network splits | The side with a majority carries on. The other side refuses writes and reads rather than drift apart. | By design. |
| A copy is damaged on disk | Checksums catch it when it is read or during a background check, the damaged block is fetched again from a healthy copy, and a block that fails its check is never served. | By design. |
The two measurements come from our six-node cell, where we killed the leader of a tablet a third of the way through 45 seconds of continuous writes, then read back every single write the cell had acknowledged.
How a cell grows
More data nodes for more data and more work. A new data node joins the cell with no data and no configuration of its own. The metadata group then moves copies of tablets onto it, one at a time, until it carries its share, and leaders follow. For our scale tests we grew the cell from six data nodes to thirty.
More gateways for more requests. Gateways hold nothing that must survive a restart, so adding one moves no data at all. On our six-node cell, point reads rose with every gateway we added while the data nodes stayed the same:
| Gateways | Reads a second | Against one machine on its own |
|---|---|---|
| 1 | 86,536 | 0.81 times |
| 3 | 194,403 | 1.80 times |
| 7 | 365,923 | 3.43 times |
One machine on its own answered 106,688 reads a second on the same data. With a single gateway the cell lost to it, because one machine was handling every request while six data nodes sat half idle. With seven gateways, the busiest machines were data nodes.
What it has measured so far
On 4 October 2026 we grew our test cell to thirty data nodes, ten in each of three zones in the AWS Hyderabad region. Each was the same machine, with 8 vCPUs and 16 GiB of memory, and we ran two tests on the cell.
TPC-C at 1,000 warehouses
TPC-C models a wholesale business with many warehouses: new orders, payments, deliveries and stock checks, all arriving at once. We ran the full specification of its five transactions as stored procedures, with 2,500 clients working on 1,000 warehouses. The database held 499 million rows, which took 8.7 TB of the cell's disks with three copies and the object tier, and it loaded in 29 minutes.
| Measure | Result |
|---|---|
| New orders a minute (tpmC) | 235,707 |
| Transactions a second, all five kinds | 9,035 |
| Failed transactions | 0 |
| New-Order time | 0.40 s on average, 0.65 s for 90 in 100 |
| Payment time, 90 in 100 | 0.31 s |
| Delivery, queued to done, 90 in 100 | 1.2 s |
| Busiest resource | every data node's processors, at 94 to 98% |
The same benchmark on six data nodes, with 100 warehouses and 500 clients, gave 52,165 tpmC. Five times the machines gave 4.5 times the work: 7,857 tpmC per data node at thirty, against 8,694 at six, which is 90% of linear.
How to read these numbers. We measured without the keying and think times that the TPC-C rules require between a client's transactions. Our clients sent their next transaction at once, which is 236 tpmC for each warehouse, where the rules allow at most 12.86. The results compare one cell with another; they are not comparable with published TPC-C results.
SQL over a billion rows
Our SQL benchmark for a single instance runs on four tables of 2.5 million rows each. For the cell we generated the same tables one hundred times larger, 250 million rows each and a billion in total, and loaded them in 2,305 seconds: 433,839 rows a second from 16 loaders, with no failed batch. Each table was split into 60 tablets over the thirty data nodes.
Then we ran the same two sets of queries as at 10 million rows, from one client, timing each query on its first run. The cell answered 3,548 of 3,551 queries in 145.9 seconds, which is 7.1 times the time it took at 10 million rows for 100 times the data. Of a shorter set of 143 queries that we also run on PostgreSQL and SQLite, it answered 123 in 32.1 seconds. The slowest single query, counting the distinct salaries of 250 million customers, took 8.5 seconds. These are the heaviest kinds of query, timed on their first run:
| Kind of query | Queries | 1 billion rows | 10 million rows |
|---|---|---|---|
| Count of distinct values | 40 | 39.0 s | 3.3 s |
| Pages deep into a result, up to row 10,000 | 80 | 27.3 s | 1.6 s |
| Grouping with a filter | 184 | 18.2 s | 1.2 s |
| Sorting with a filter | 320 | 11.9 s | 1.6 s |
| Ranges | 253 | 10.3 s | 1.5 s |
| Sorting on two keys | 64 | 9.5 s | 1.5 s |
| Text patterns (LIKE) | 180 | 6.4 s | 1.0 s |
| Several conditions together | 1,860 | 5.7 s | 3.1 s |
| Exact matches | 240 | 1.3 s | 0.4 s |
The queries the cell did not answer, 3 of the 3,551 and 20 of the 143, each compute totals over a join of whole tables, between 389 million and a billion rows. Today a join's totals are computed on one machine, so the cell refuses those joins within milliseconds and says why, rather than run out of memory. Join queries that return a page of rows are answered, in 0.1 to 1.2 seconds.
Where it is behind
- Writes grow more slowly than reads. Every write is carried out on three machines, so the cell's capacity for writes grows with its number of data nodes divided by three, while a read is carried out once. On our six-node cell, with disks at their baseline speed, the cell made 1,796 writes a second against 1,090 for one machine on its own: 1.65 times, where reads reached 3.43 times.
- Large queries are slower than on a small database. On their first run, the full set of 3,551 queries took 7.1 times as long at a billion rows as at 10 million, and counts of distinct values and deep pages took 12 to 17 times as long.
- Totals over joins of whole tables are refused. Until a join's totals are computed where its rows live, the cell refuses joins of more than 20 million rows that ask for totals, as described above.
- A whole zone failing has not been measured yet. We have measured a machine failing and coming back. A zone failing is covered by the design, and measuring it is part of the work before release.
The work before release
We plan to release InventDB Cluster by the end of 2026. Between now and then, the work is about making it stable and ready for production at larger enterprises and for much larger databases, so that a cell needs as little looking after as a single instance does today. These are the main pieces of that work:
- Data stored once in the object tier. Today every copy of a tablet keeps its finished data on its own disk. When a tablet's finished data lives once in the cell's object store and is shared by its three copies, machines hold much less, a replacement copy can start in minutes, and data that is updated often stops growing on disk.
- Tablets that split and merge by themselves as data and traffic grow and shrink. Splitting and merging work today as operations; the cell does not yet decide when to do them.
- Many more tablets on each machine, by lowering what each tablet costs in files and memory, so a cell can hold far larger databases.
- Totals over large joins computed where the rows live, so those queries are answered instead of refused.
- Scaling up and down with nobody deciding. The cell will ask for machines when it is busy and give them back when it is quiet, within limits its owner sets, and a database nobody is using will cost only its storage.
- The failures we have not yet measured, starting with a whole zone, measured the same way we measured a machine.
We measure every step on a full cell, against the numbers on this page, before we keep it.
InventDB Cluster is one of three research programmes at InventDB, alongside InventDB Sparkle, our inference engine for open models, and InventDB Brahma, our research toward a human-like intelligence. All three are listed under Research. If your database has outgrown one machine, or soon will, we would like to hear from you at contact@inventdb.com.