InventDB
All articles Query engine

Columnar aggregates inside a JSON document database

InventDB stores records as JSON documents, but it answers COUNT, SUM, AVG and GROUP BY without reading them. Counts and sums come from totals that the indexes keep on every write, and larger aggregates walk dense column files that the engine writes when a segment seals.

The straightforward way for a JSON database to answer SUM(amount) GROUP BY region is to read every document, decrypt it, parse it and pick out two fields. On a table of 10 million records, that is 10 million decryptions and JSON parses to produce a few dozen numbers.

InventDB keeps records as JSON, because that is what applications send and read. It answers aggregates from structures built for them: totals that the field indexes keep up to date on every write and, for each sealed segment, column files that hold one dense array per field. Each table is stored as segments of 100,000 records, one active and the rest sealed, as described in How a single instance holds 250 million records. This article walks through the aggregate paths from the cheapest to the most general, and then covers what they cost.

Counting without reading records

Three common counts never touch a record:

  • COUNT(*) with no WHERE reads the type's stored record count. If writes are waiting in memory for the next checkpoint, the engine checks each pending insert and delete against the id index and adjusts the count, so the cost follows the number of pending writes rather than the size of the table.
  • A single-column COUNT with GROUP BY reads the field's counts tree, which holds each distinct value with the number of rows that have it. Every write updates those tallies, so the query reads one entry per distinct value: a status field with six values costs six entries per segment, whether the table holds a thousand rows or a hundred million.
  • A two-column COUNT with GROUP BY and no WHERE reads, for each sealed segment, a stored tree of value pairs and their counts. The active segment is still changing, so it computes its pairs fresh from its two field indexes and keeps the result against a write counter. Every insert, update and delete advances the counter, which discards the stored result the moment the segment changes.

A filtered COUNT(*) costs what its filter costs to produce the matching ids, plus a length. WHERE col IS NOT NULL and WHERE col IS NULL are answered from the field index's entry counts without collecting any ids.

SELECT status, COUNT(*) AS n FROM shop.orders GROUP BY status

Sums, averages and ranges from running totals

For an unfiltered SUM, AVG or COUNT of a numeric field, every segment keeps a running count and sum for the field that each write adjusts. The engine adds the segments' totals and subtracts the contribution of any row that was updated after its segment sealed and still has a stale copy there. MIN and MAX come from the smallest and largest live entries in each segment's value index.

A range on the field being aggregated never collects the matching ids. For SUM(amount) WHERE amount BETWEEN 100 AND 500, sealed segments scan their column for amount, a flat array of numbers, and the active segment walks its counts tree between the two bounds. If that scan cannot be used, every segment walks its counts tree, adding the count and sum stored with each distinct value in the range, and a segment whose minimum and maximum lie outside the range is skipped without a walk.

A one-sided range that asks only for COUNT, SUM or AVG can be cheaper still. The engine checks which side of the bound covers less of the field's values, and when the excluded side is smaller it computes the totals minus that side:

SELECT SUM(amount) AS s, AVG(amount) AS a
FROM shop.orders WHERE amount > 20

-- if 20 is near the bottom of the field's values, this is answered as
-- the field's totals minus COUNT and SUM where amount <= 20

MIN and MAX of a range cannot be derived by subtraction, so they always take the direct route.

Column files for sealed segments

When a segment seals, the engine writes one column file per indexed field, on a background thread so that writers do not wait. It builds each file from the field's id-to-value tree. Numbers and dates become two dense arrays of the same length: 64-bit hashes of the record ids, and 64-bit floating-point values. Strings become a dictionary of the distinct values, an array of id hashes and an array of 32-bit positions in the dictionary, so a string that repeats a million times is stored once.

Each column file is encrypted with AES-256-GCM and carries a CRC32 of its contents, which the engine checks when it loads the file. A column file is written once and never updated. Before using one, the engine compares its row count with the number of entries the field's index holds for that segment. If the two differ, because cleanup has since removed stale copies from the segment's indexes, the engine answers that segment from the index instead. The active segment has no column files and is always answered from its indexes.

Two columns of a sealed segment, sorted by id hash, walked together into one accumulator slot per dictionary entry, then merged across segments ONE SEALED SEGMENT, COLUMNS LOADED region string i0#0a712 i1#1c3e0 i2#2f902 i3#4b121 dictionary 0 EU, 1 US, 2 APAC amount number i0#0a71120.00 i1#1c3e75.50 i2#2f9042.00 i3#4b12310.00 sorted by id hash, in the same order as region position i is the same record in every column One slot per dictionary entry 0 EU75.501 row 1 US310.001 row 2 APAC162.002 rows Merge across segments sums and counts add, min and max keep extremes Each sealed segment fills its own slots; the active segment is answered from its indexes.
Sorting each column by id hash on load lines up columns that cover the same records, so a GROUP BY reads position i from both arrays and updates one slot, with no hashing per row.

GROUP BY as two arrays walked in step

When a column is loaded, the engine sorts its arrays by id hash. Two columns that cover the same records then line up exactly, and position i in each describes the same record. That turns a GROUP BY into a walk over two arrays with one index:

SELECT region, SUM(amount) AS revenue, COUNT(*) AS orders
FROM shop.orders GROUP BY region

For each sealed segment the engine loads the region and amount columns and walks them together. The accumulator is a plain array with one slot per dictionary entry, so adding a row means reading the dictionary position at i and the amount at i and updating that slot, with no hashing and no string comparison per row. Each slot keeps a sum, a minimum, a maximum, a count of rows that have a value and a count of rows; AVG is the sum divided by the count of values. If some records lack the aggregated field, its column is shorter and the positions no longer line up, so the engine finds each record's value through a map from id hash to position, built once per column the first time it is needed.

Each segment produces a small partial result keyed by group, and the engine folds the partials together: sums add, counts add, and minimum and maximum keep the extremes. Memory grows with the number of distinct groups rather than the number of rows, so a GROUP BY over 250 segments never holds the rows themselves. The active segment produces its partial from its field indexes.

With a WHERE clause, the engine first produces the matching ids from the indexes and hashes them once into a set, and the walk skips rows whose id hash is not in it. A record that was updated after its segment sealed is counted once, because the walk skips its stale copy in every segment except the one holding the newest version. This path serves one grouping column with at most one aggregated field. GROUP BY over two or more columns is assembled from the field indexes instead, with a key built from every grouping column.

The column cache

Loading a column costs a read, a decryption, a checksum and a sort, so the engine keeps decoded columns in one cache shared by every type in the process. A hit takes a shared lock on one shard and stamps the entry with an atomic counter, and the arrays are shared between concurrent queries without copying. A numeric column costs 16 bytes per row in memory, so one column of a 100,000-record segment is about 1.6 MB; a string column costs 12 bytes per row plus its dictionary.

The cache is a cap rather than a reservation, because entries are added only when a query reads them. Its size is a share of the machine: an eighth of its memory, and never more than a twelfth of the memory available when the cache is created. When it is full it evicts whole columns, least recently used first.

What it costs

  • More files on disk. Every sealed segment holds one extra file per indexed field.
  • A first load. The first query to touch a column pays to decrypt the whole column, check it and sort it. Later queries reuse the decoded arrays until the cache evicts them.
  • Recent data takes the slower path. The active segment has no column files, so up to 100,000 records per type are aggregated from index pages.
  • Updates to sealed records cost the fast paths. While a type has stale copies waiting for cleanup, the columnar range scan declines and the engine uses paths that subtract or skip the stale copies.
  • Floating-point sums. Numeric values are summed as 64-bit floating-point numbers.
  • Not every aggregate has a fast path. COUNT(DISTINCT col) is not one of the paths above, and a GROUP BY that aggregates several fields at once takes a more general path.

Using it

Nothing needs to be declared. Every field is indexed on write and every sealed segment gets its column files, in InventDB Serverless and InventDB SOAR alike. Two habits keep aggregates on the fast paths:

  • Store numbers as JSON numbers and dates as ISO 8601 strings. The engine types a field from its values: a JSON number is numeric, a string that starts with a YYYY-MM-DD date is a date, and anything else is a string. An amount sent as "129.50" in quotes is a string, and SUM over it falls back to reading documents.
  • Aggregate the field you filter on when you can. A range on the aggregated field is answered from totals and columns without collecting ids.
POST /sql?metrics=1
Authorization: Bearer <token>
Content-Type: application/json

{
  "sql": "SELECT region, SUM(amount) AS revenue, COUNT(*) AS orders FROM shop.orders GROUP BY region ORDER BY revenue DESC"
}

With metrics=1 the response wraps the rows with the elapsed milliseconds and the row count, which shows the difference between the first run of a query, which loads its columns, and the runs after it. For a comparison across whole workloads, our benchmarks page reports InventDB at 6.3x PostgreSQL, warm, on the published 143-query set over 10 million records, with InventDB encrypted and PostgreSQL 18 tuned and unencrypted, and it publishes the categories PostgreSQL wins alongside.