An Embedded Sketch Database: Billions of Events, Megabytes of Storage, Error Bars on Every Answer
Counting unique visitors exactly means remembering every visitor. Say 6 million distinct IPv4 addresses hit a checkout service today. Packed at 4 bytes each, that set is 24 MB before any hash table overhead, and you need one per service, per region, per day you want to compare. Redis answers the same question with a HyperLogLog in about 12 KB per key, at a standard error of about 0.81%.
The thin part is the database around it: a SQLite-shaped file of sketches with time buckets, rollups and a SQL front door (check before you start; this space moves). The version worth building stores nothing but mergeable sketches, grouped by dimension and time bucket, and every aggregate returns its error bound next to the value. You never keep the events. It’s the extreme case of reducing data near the source: keep the answer, drop the rows.
Apache DataSketches is the serious prior art: HyperLogLog, KLL quantiles, Theta sketches for set operations, frequent-items summaries, all with documented serialization formats. ClickHouse, Druid and Pinot ship approximate aggregates, postgresql-hll adds HyperLogLog columns to Postgres, and DuckDB’s approximate functions run over raw rows you’ve already stored, which is the cost this design avoids. What’s left open is the small, self-contained file.
Each Sketch Is Wrong in Its Own Way
HyperLogLog hashes each item, uses the first bits to pick one of m registers, and keeps in each register the longest run of leading zeros seen in the rest of the hash. The standard error is about 1.04/√m, so Redis’s 16,384 registers (6 bits each, 12 KB) give 0.81%, while 1,024 registers give 3.25% in 768 bytes. Size depends on m, never on the stream: a billion events cost the same 12 KB as ten thousand. HyperLogLog++ (Google, 2013) improved small-range accuracy and added a sparse form, so a cell that saw 40 distinct IPs doesn’t pay for the full array. Merging is a register-wise max, which loses nothing.
Count-Min Sketch is a grid of counters, d rows by w columns. A query takes the minimum across rows, so it never undercounts. With w about e/ε and d about ln(1/δ), the overcount stays under ε times the total event count N, with probability at least 1 minus δ. Take ε = 0.001 and δ = 0.01: 2,719 columns by 5 rows, about 54 KB at 4 bytes a counter. That suits “is this URL a big one” and fails “how often did this rare URL appear”, because the error is a slice of all traffic. At a billion events the overcount can reach a million. Count-Min also forgets its keys, so for a top 20 you want Space-Saving: k triples of key, count and error, where any key above N/k events is guaranteed to be present and every count states the most it could be too high.
KLL promises an error in rank: ask for p99 and the value you get has a true rank within ε of 99%, which at an ε of 1% could be anything from p98 to p100. For tail latency that’s loose, so size the sketch for the tail you care about. A t-digest packs its resolution toward the tails and is the practical favorite for p99.9, but it has no clean worst-case bound. A database promising a bound on every answer should default to KLL, offer t-digest as an explicit choice, and mark those answers unbounded. To report the error in milliseconds, read the sketch at q minus ε and q plus ε; that pair is the interval.
Bloom filters (1970) give no false negatives and a tunable false-positive rate: about 9.6 bits per key buys 1%, so a million keys cost 1.2 MB. Size grows with the keys you expect, which breaks the fixed-size promise, but the error bar is free, since the fraction of bits set gives the current false-positive rate. Name the function might_contain.
One Cell per Dimension and Time Bucket
Storage is a table of cells keyed by dimension values, tier and bucket start, each holding the sketches you declared, as blobs. Ingest folds events into the open bucket in memory and writes the cell when it closes. Compaction merges 60 one-minute cells into an hourly one and 24 of those into a daily one, then deletes the children. A daily answer isn’t an average of hourly answers: merged HyperLogLogs are exactly the sketch you’d have built from the whole day, and the other sketches keep their guarantees through a merge.
Cost is cells times cell size. A 4,096-register HyperLogLog (3 KB), a quantile sketch and a short top-K list come to about 5 KB, so ten thousand cells is 50 MB, against 50 GB of raw rows (a billion events at 50 bytes each). But “megabytes” means thousands of cells, not millions, so dimension values times buckets kept is the budget you manage.
That gives the schema one hard rule. Dimensions you group by go in the key and stay low in cardinality: service, region, status class. Things you count go inside the sketches: IP, URL, user. Put user ID in the key and every cell holds one user, which makes a sketch database an expensive table. The cookieless analytics idea in the nine-apps post is the natural customer, since daily unique visitors is a distinct-count problem and a HyperLogLog fed a salted hash never stores an address.
These commands are a proposal.
$ sketchdb create traffic.sk \
--dims service,region \
--tiers 1m:2d,1h:90d,1d:forever \
--sketch ip:distinct --sketch latency_ms:quantile --sketch url:topk:1000
$ sketchdb query traffic.sk "
SELECT approx_count_distinct(ip), percentile(latency_ms, .99), top(url, 3)
FROM requests
WHERE service = 'checkout' AND ts >= '2026-06-01' AND ts < '2026-06-02'"
approx_count_distinct(ip) 1.20 M 95% interval 1.17 M to 1.24 M
percentile(latency_ms, .99) 418 ms 95% interval 396 ms to 455 ms
top(url, 3) /cart 311,402 at most 8,120 too high
/pay 190,877 at most 8,120 too high
/sku/771 88,415 at most 8,120 too high
read 3 day cells (one per region), 15 KB, 0 events
Where It Breaks
Declare sketches before the first event. The events are gone by the time you change your mind, so a quantile sketch can’t give you distinct counts later, and a sketch added to a running file covers only buckets from that moment on. Older cells must answer NULL with a reason, never zero, because zero reads as “nobody visited”.
Intersections are the sharp edge. HyperLogLog does union for free and overlap badly: you get it as the size of A plus the size of B minus the size of their union, and each estimate carries its own error. Take two days of 10 million visitors with 10,000 in common: a 1% error on each is 100,000, ten times the overlap, so the answer comes out as zero or garbage. “Who came back tomorrow” is exactly that query. Theta sketches do union, intersection and difference directly, at a larger size, so offer a set sketch type and make intersections on HyperLogLog columns fail with a message that says why.
Merging works only between sketches built the same way: same hash function and seed, same width and depth. The header should pin the hash by name and version, record every sketch’s parameters, and refuse a mismatched merge by naming the field. Encoding is the quieter trap: the string “10.0.0.1” and the same address as four packed bytes hash differently, so two writers that disagree inflate a distinct count with no error raised. Fix one canonical encoding per column type in the format, and where DataSketches has a serialization format, store that. Interop beats novelty.
Then there’s the reader, who expects count(DISTINCT ip) to be exact. Reject that spelling instead of quietly approximating it, and round to the digits the error allows: “1,204,318” with a 3.3% error is false precision.
What Version 0.1 Refuses to Do
One writer, one file. Three sketch types: HyperLogLog, KLL and Space-Saving. SQL is a single SELECT with the sketch aggregates, a WHERE over dimensions and time, and a GROUP BY on a time unit; events arrive as newline-delimited JSON on stdin, the territory of a small SQL engine for JSON streams. No joins, no exact counts, no deletes, no aggregate that can’t state a bound.
The test that makes it credible ships with it. Replay a day of real traffic into the engine and into a scratch database that keeps every row, then check that the exact answer lands inside the stated interval about 95 times in 100. If the bound lies, nothing else matters. Which tiers to keep and what survives each is the retention question the telemetry rollup post works through.
The error bar is the product.