Probabilistic data structures: Wrong on purpose
by Davor Lozic, CTO
You have 100 million URLs in a table, and every request asks the same question: have we seen this one before? Almost always the answer is no. Keeping the keys in memory costs 1.6 GB before you count anything, three or four times that once Python or Java has wrapped each key in an object. You pay all of it to mostly say no.
Here is the alternative with the numbers attached. 1.2 MB of bits answers the same question in seven array lookups. It is never wrong when it says no. It is wrong about one time in a hundred when it says yes, and a wrong yes costs you a single database read you were going to do anyway. You never store a URL.
That is the whole family in one example: give up one direction of correctness by a measured amount, and your memory stops depending on your data.
The trade: a Bloom filter and what it actually costs
A Bloom filter is m bits, all zero, and k hash functions. To insert a key, set the bits its k hashes point at. To query, check whether all k of those bits are set. Sixteen bits and three hashes, by hand:
| Operation | Hashes to | Result |
|---|---|---|
insert("/login") | 3, 9, 14 | sets bits 3, 9 and 14 |
insert("/logout") | 1, 9, 12 | sets bits 1 and 12; bit 9 was already on |
query("/admin") | 3, 12, 15 | bit 15 is zero, so definitely absent |
query("/signup") | 1, 3, 14 | all three on, so maybe present |
After the two inserts, five of the sixteen bits are on:
/signup was never inserted. Its three bits were set by the other two keys, and the filter cannot tell the difference. That is a false positive, and it is the only way this structure lies. The other direction is airtight: no false negatives, ever, because a member's bits were set at insert and nothing clears a bit.
The same filter also shows what k is for. One hash per key and /signup collides constantly; ten hashes per key and all sixteen bits are ones within three inserts, after which everything is maybe. The balance point is half the bits set, at k* = (m/n) · ln 2, where the false-positive rate for n keys in m bits works out to
p = (1 - e^(-kn/m))^k
Inverted, that is the sizing rule m/n = -1.44 · log2(p) bits per element. For a million keys:
| Target false-positive rate | Bits per element | Optimal k | Total size |
|---|---|---|---|
| 10% | 4.79 | 3 | 0.60 MB |
| 1% | 9.59 | 7 | 1.20 MB |
| 0.1% | 14.38 | 10 | 1.80 MB |
| 0.01% | 19.17 | 13 | 2.40 MB |
Each factor of ten in accuracy costs a flat 4.79 bits per element. But the number to carry out of the table is the one not in it: the space depends only on the error rate you accept, and not at all on the size of your elements. A filter over 100-byte URLs costs what a filter over 4-byte integers costs, and that 1% filter is 1.2 MB against 40 to 100 MB for a real hash set once you count pointers, load factor and object headers.
The economics, which is where these get misapplied
A filter is a rejection stage, never a source of truth. Every maybe still goes to the authoritative store, which confirms the hit or reveals the false positive, so correctness is never at risk - only the number of lookups changes.
So count the lookups. Take 10,000 requests a second and that 1.2 MB filter in front of the table. At a 2% hit rate, 200 requests are genuine hits and every one of them goes through; of the 9,800 genuine misses the filter waves 1% through anyway, which is another 98. Traffic to the database: 298 instead of 10,000.
Nothing about the filter changes in the rows below. Only the hit rate does.
| Hit rate | Reaching the database | Of those, false positives | Load removed |
|---|---|---|---|
| 0.5% | 150 | 100 | 98.5% |
| 2% | 298 | 98 | 97.0% |
| 10% | 1,090 | 90 | 89.1% |
| 50% | 5,050 | 50 | 49.5% |
| 90% | 9,010 | 10 | 9.9% |
Read the last column down and the deal collapses. Read the third column down and you can see why: the false positives, the thing the filter is criticised for, shrink as the hit rate climbs. They are never the problem. The real hits are, because a filter cannot reject those - they belong in the database and it has to go there. Formally the traffic still reaching the store is h + (1 - h)p, which tends to h as h grows, and by the bottom row you are paying 1.2 MB and a hash per request to remove one query in ten.
The structure only pays when no is the common answer. If you cannot state your hit rate, you are not ready to add a filter.
Every LSM-tree engine - RocksDB, Cassandra, HBase - puts a filter on each SSTable to sit in that top row deliberately: a key lives in one file out of dozens, so the per-file hit rate is a fraction of a percent by construction, whatever the hit rate of the query that triggered the read. Different route, same goal as the B-tree indexes of relational engines: do not touch the disk unless you have to.
The failure mode you will actually hit
A filter is sized for a maximum, not an average, and nothing enforces it at runtime. You size for a million keys at 1%, the table grows to three million over two quarters, and the filter keeps its 9.6 million bits and its seven hashes:
| Keys in a filter sized for 1M | Actual false-positive rate |
|---|---|
| 1,000,000 | 1.0% |
| 2,000,000 | 15.7% |
| 3,000,000 | 43.6% |
No error, no log line. Database load triples while every dashboard says the cache layer is healthy.
You also cannot delete. Clearing /logout in the example above would zero bit 9 and make /login vanish - a false negative, destroying the one guarantee you had. Use a counting or cuckoo variant.
The proof: exact is not merely expensive
The second case is harder, because there the cheap structure is not an optimisation but the only option. Count distinct users per day, then per country, per platform, per campaign, and per whatever somebody wants to slice by next quarter. Exact means keeping every ID seen in every bucket: gigabytes per dimension per day, for a number nobody reads past two significant figures. You will try to be clever about it and fail, in a way that feels like it should be solvable. It is not, and there is a theorem saying so:
Computing the number of distinct elements exactly in one pass requires Ω(n) bits of space, where n is the size of the universe.
The argument is short once you look at a streaming algorithm the right way. The algorithm's state is a message. Shard your stream across two machines: the first consumes its half and ships its s-byte state to the second, which resumes from it and finishes the count. Whatever that state does not carry, the answer cannot depend on.
Now use the pair as an oracle. They report |A ∪ B| exactly, and each machine knows its own |A| and |B|, so subtraction hands you |A ∩ B| and in particular whether it is zero. Deciding that, set disjointness, is known to require Ω(n) bits of communication. The machines exchanged s. Therefore s = Ω(n): your state is as big as the universe, and no cleverness in the algorithm changes it.
Approximation is what buys the exemption. Off by 1% on a billion-element union, the count cannot tell an intersection of size 0 from one of size 1, so it is useless as a disjointness oracle and the reduction never starts.
On the other side of that line, HyperLogLog counts distinct elements to within 0.81% standard error in 12 KiB, at any cardinality up to billions. The trick is one you can check in your head. Somebody flips a coin for a while and tells you the longest run of tails they hit was 12. You can tell them they flipped about 2^12 = 4,096 times, without having watched a single flip, because that is how long you wait for twelve tails in a row. HyperLogLog runs that backwards over your data: hash each element to a uniform bit string, which is the coin, so the chance a value begins with k zeros is 2^-k and across N distinct values the longest run of leading zeros is about log2 N. Track that maximum, report 2^max. Duplicates are free, since rehashing an element gives the same bit string and cannot raise the maximum.
One run is a terrible estimate, because a single lucky element doubles your answer. So split the stream: use the first few bits of the hash to pick one of m registers, keep a maximum in each, and average with the harmonic mean, which drags a freak register back toward the pack. Standard error is 1.04/√m, and since a register only holds "longest run seen", 6 bits covers it.
| Registers | Memory | Standard error |
|---|---|---|
| 4,096 | 3 KiB | 1.62% |
| 16,384 | 12 KiB | 0.81% |
| 65,536 | 48 KiB | 0.41% |
Four times the memory for half the error, every time. That is a theorem rather than a tuning problem, so 1% to 0.1% costs 100x. Redis picked 16,384 registers because 12 KiB at sub-1% is where both numbers stop mattering to most applications.
Count-Min does the same job for frequencies. Keep d rows of w counters, each row with its own hash; on update, increment one cell per row; to query, report the smallest of the d cells. Four rows of 2,718 counters at 4 bytes is 43 KB, and a lookup reads four numbers:
request count for "/api/v2/users"
row 0 -> cell 1,904: 806
row 1 -> cell 77: 812 <- something else hashed here too
row 2 -> cell 2,310: 1,190 <- and here, worse
row 3 -> cell 1,455: 806
report min = 806 true count is 806 or slightly under
Every cell counted the real increments plus whatever else collided there, so the estimate can only overshoot, never undershoot - one-sided error again, the same shape as the Bloom filter. A row's overshoot averages N/w and all d rows being unlucky at once has probability e^-d, which is why four rows is already enough. Set w = e/ε and d = ln(1/δ) for every frequency in a stream of any length, to within 0.1% of the total, in under 80 KB.
Both structures merge - two HyperLogLogs by per-register maximum, two Count-Mins by cellwise sum - so a 200-node cluster counts distinct users by having every node send 12 KiB to a coordinator that maxes them together. No coordination, no re-reading. That is the real reason every analytics system ships them, from Redis PFCOUNT to BigQuery's APPROX_COUNT_DISTINCT.
One cost arrives later than the benefit. HyperLogLog's raw estimator is biased at small cardinalities, so the original paper switches to linear counting below 2.5m and Google's HLL++ uses a correction table. A naive implementation reports nonsense under a few thousand elements, which is exactly the range your integration tests run in.
The arithmetic: the birthday bound sizes all of it
Both structures above want collisions and budget for them. The same arithmetic bites when you did not intend any.
Twenty-three people in a room, more likely than not that two share a birthday: the version everyone has heard and nobody believes. The engineering version is that with N equally likely values and m items drawn at random you cross 50% at m = 1.177 · √N, and below that
P(collision) = 1 - e^(-m² / 2N) ≈ m² / 2N when m² << N
It surprises people because you count items and the arithmetic counts pairs: twenty-three people are 253 pairs. Each pair collides with probability 1/N, so expected collisions are m²/2N - set that to 1, take the square root, and you get the rule that an n-bit hash gives you n/2 bits of collision resistance.
| Digest | Space | 50% collision at |
|---|---|---|
| 32-bit | 2^32 | ~77 thousand |
| 64-bit | 2^64 | ~5 billion |
| 128-bit (MD5) | 2^128 | ~2^64 |
| 160-bit (SHA-1) | 2^160 | ~2^80 |
| 256-bit (SHA-256) | 2^256 | ~2^128 |
Seventy-seven thousand. Content-address your uploads with CRC32 and you will serve somebody the wrong file before the bucket outgrows a laptop, not at the four billion the bit width suggests. Truncating a stronger digest to 32 bits buys you nothing either: the truncation sets the bound, not the algorithm.
The same theorem gets used to justify real waste, though. "A UUIDv4's 122 random bits are only 61 bits of safety, so we need a central ID service" has cost teams a database dependency they did not need: 2^61 is about 2.3 × 10^18 UUIDs, so at a billion per second you wait 73 years. The right response to a factor of two in an exponent is to check the arithmetic, not to add a coordination point. And the bound only applies to random selection in the first place - a counter, a sequence, or a Snowflake ID from a machine number and a timestamp never collides at all, by construction.
The fine print: every bound above is a claim about your hash
Three sections of probabilities. Probabilities over what, exactly? Not over your keys - those are usernames, URLs and whatever a stranger typed into a form, and somebody else gets to choose them.
So move the randomness. A family of hash functions is universal if, for every pair of distinct keys x ≠ y,
Pr[h(x) = h(y)] ≤ 1/m over the random choice of h from the family
The "over" clause is the whole idea: the probability is taken over which function you drew at startup, not over which keys arrived. Same move as randomised quicksort - stop hoping the input is not sorted, pick the pivot at random instead - and nobody can attack a coin they did not see.
No fixed function gets there, and this is not hypothetical. Your key space dwarfs your table, so pigeonhole guarantees some bucket attracts at least |U|/m keys, and with a published deterministic hash an attacker finds that set offline at leisure. In 2011 somebody did precisely that for the string hashes in PHP, Java, Python and Ruby, then put a few thousand colliding keys in one POST body: every insert walked the same bucket, the O(1) hash table became an O(n²) linked list, and a single request pinned a CPU core for minutes. Every major runtime shipped an emergency release.
The usual fix is a fast non-cryptographic hash such as xxHash plus a per-process random seed, and it is worth being honest about what that buys. The seed kills accidental clustering and untargeted flooding, but a non-cryptographic hash's seed can be recovered from observed behaviour, which is why runtimes expecting real adversaries use a keyed cryptographic function - SipHash, in both Python's and Rust's case. Either way the protection is lost the same way. A universal family used with a constant seed is a fixed function, and provides nothing. A seed hardcoded once so two services would agree on a shard, or so a test would stop flaking, reverts you to a deterministic hash, and every test still passes.
Which lands back on the first section. A Bloom filter fed keys a user can choose does not have a 1% false-positive rate - it has whatever rate the attacker decides, because finding keys that light up bits already set is the same offline search as finding colliding bucket keys.
The five questions
Before you put one of these in a system:
- Which direction can you be wrong in, and who catches it? With no authoritative store behind the structure to resolve a maybe, you do not have a filter. You have a bug with good asymptotics.
- Is no the common answer? Compute h + (1 - h)p with your real hit rate. Close to h, build it. Close to 1, do not.
- What is the maximum n, not the expected one? One million keys to three million took that filter from 1% to 43.6%, silently.
- What does the next digit of accuracy cost? A flat 4.79 bits per element for Bloom, 4x the memory per halving of error for HyperLogLog. Neither is negotiable at runtime.
- Where does the seed come from, and is it per-process? If the answer is a constant in a config file, the bounds you quoted in the design review do not apply.
These trade-offs come up most in infrastructure work, where the memory budget and the latency budget turn out to be the same budget. What makes any of it engineering rather than gambling is that the error is bounded, computable in advance, and yours to choose before you ship.