The textbook reservoir sampler fits in six lines: keep the first k items, then replace a random slot with item i with probability k/i. The companion article on reservoir sampling proves why that gives every item the same chance of ending up in the sample. This article starts where that one stops: what happens when the stream is a billion events a day, spread over 64 Kafka partitions, when some events matter more than others, and when you want the sample for the last hour rather than for all time.
Those requirements change the design more than the algorithm. You will see that the fastest single-stream method is not the one you want across machines, that the obvious way to merge two reservoirs is wrong, that the popular weighted method does not do what most people assume, and that one formulation, keeping the k smallest random keys, solves merging, weighting, deduplication and windows at once.
The requirements a production sampler must meet
A production sampler needs more than uniformity, and each variant trades one property for another:
- Bounded memory: exactly k items, however long the stream runs.
- One pass: each event is seen once, in arrival order, and the total count n is unknown until the end.
- Low per-item cost: at millions of events a second, even one random number per event is measurable.
- Mergeability: samples taken independently on shards must combine into a correct sample of the union, without rereading data.
- Weighting: errors or expensive requests should be more likely to be kept than routine ones.
- Windows: a sample of the last hour, not of everything since the process started.
- Reproducibility: the same input should give the same sample, so a bug report can be replayed.
Algorithm R (attributed to Waterman in Knuth's volume 2) meets the first two and fails the rest. The sections below add the others one by one.
Making it fast: skip instead of flip
Late in a long stream, almost every item is rejected. With k = 1,000 and n = 109, item i is kept with probability 1000/i, so Algorithm R spends a billion random draws to make about k(1 + ln(n/k)) = 1000 × (1 + 13.8) ≈ 14,800 replacements. The idea behind Vitter's 1985 algorithms (X, Y and Z) and Li's simpler Algorithm L (1994) is to draw the gap to the next accepted item directly, so the cost becomes proportional to the number of replacements rather than to n: expected O(k(1 + log(n/k))).
Algorithm L tracks w, the current largest of k uniform random keys. The gap to the next item whose key beats w is geometric, so one logarithm draws it:
import math
import random
def _u(rng):
u = rng.random() # [0, 1); zero would break log()
while u == 0.0:
u = rng.random()
return u
def algorithm_l(stream, k, rng=None):
rng = rng or random.Random()
it = iter(stream)
reservoir = []
for item in it: # fill phase: the first k items are all kept
reservoir.append(item)
if len(reservoir) == k:
break
if len(reservoir) < k:
return reservoir # the stream was shorter than k
w = math.exp(math.log(_u(rng)) / k)
while True:
skip = math.floor(math.log(_u(rng)) / math.log(1.0 - w))
try:
for _ in range(skip):
next(it) # skipped items cost no random draw
item = next(it)
except StopIteration:
return reservoir
reservoir[rng.randrange(k)] = item
w *= math.exp(math.log(_u(rng)) / k)The draw must never be exactly zero, hence the helper. Skipped items are still consumed from the iterator but cost no random draw; when the source supports random access, such as an offset-addressed log, you can seek past the skip instead, which is where the speed-up becomes dramatic. Algorithm L is the right choice for one fast stream, but it is not mergeable.
Merging reservoirs: the tempting mistake
Suppose shard A saw 9 million events and shard B saw 1 million, and each kept a uniform reservoir of 1,000. The tempting merge takes 1,000 items at random from the 2,000 pooled. That gives each shard half the sample, when A should contribute about 90 percent. The result is biased towards small shards, and the bias is invisible because the sample still has the right size.
The correct merge first decides how many items come from each shard, by simulating a draw of k items without replacement from a population of n1 + n2. That count follows a hypergeometric distribution. Then it takes that many items uniformly from each reservoir, which is valid because each reservoir is itself a uniform sample of its shard:
def merge_algorithm_r(res1, n1, res2, n2, k, rng):
# res_i is a uniform sample of min(k, n_i) items from shard i, which saw n_i items.
want = min(k, n1 + n2)
a, b, take1 = n1, n2, 0
for _ in range(want): # draw 'want' of the n1 + n2 items without replacement
if rng.random() * (a + b) < a:
take1 += 1
a -= 1
else:
b -= 1
# take1 is hypergeometric; pick that many uniformly from each reservoir
return rng.sample(res1, take1) + rng.sample(res2, want - take1)This works, but every shard must report its count, merges need fresh randomness, and it does not extend cleanly to weights.
The architecture: keep the k smallest random keys
Give every item an independent uniform random key and keep the k items with the smallest keys. Because the keys are independent and identically distributed, every set of k items is equally likely to hold the k smallest, so this is a uniform sample. This is called a bottom-k sample, and Algorithm R and L are clever ways of simulating it without storing keys. Storing them buys four properties:
- Merge is a union. The k smallest keys of A∪B are among the k smallest of A plus the k smallest of B. No counts, no extra randomness, any merge order, any tree shape.
- Reproducible and deduplicating. If the key is a keyed hash of the item's identifier rather than a fresh random draw, the same event always gets the same key. Replaying a partition gives the same sample, and an event delivered twice collides with itself at merge instead of doubling its chance.
- Consistent across datasets. With a shared seed, two services sampling the same trace ids keep the same traces, so the sampled request in the gateway log is also the sampled request in the database log.
- Free cardinality estimate. With uniform keys, (k − 1) divided by the k-th smallest key estimates the number of distinct ids seen, the KMV sketch.
import hashlib
import heapq
import math
def key_uniform(item_id: bytes, seed: bytes) -> float:
h = hashlib.blake2b(item_id, key=seed, digest_size=8).digest()
return (int.from_bytes(h, "big") + 0.5) / 2**64 # strictly inside (0, 1)
def key_weighted(item_id: bytes, weight: float, seed: bytes) -> float:
return -math.log(key_uniform(item_id, seed)) / weight # exponential with rate = weight
class BottomK:
# Keeps the k items with the smallest keys. Mergeable and order-independent.
def __init__(self, k):
self.k, self.heap, self.seen, self.ids = k, [], 0, set() # max-heap via negated keys
def offer(self, key, item_id):
self.seen += 1
if item_id in self.ids: # same id, same key: already held
return
if len(self.heap) < self.k:
heapq.heappush(self.heap, (-key, item_id))
self.ids.add(item_id)
elif key < -self.heap[0][0]: # beats the current worst
_, evicted = heapq.heapreplace(self.heap, (-key, item_id))
self.ids.discard(evicted)
self.ids.add(item_id)
def merge(self, other):
out = BottomK(self.k)
for neg, item_id in self.heap + other.heap: # offer() drops duplicate ids
out.offer(-neg, item_id)
out.seen = self.seen + other.seen
return outThe cost is one hash per item and a heap of k entries (see heap operations); late in the stream most items need only one comparison against the heap's maximum. That hash per item is the trade-off against Algorithm L: bottom-k spends CPU on every item to buy mergeability and determinism.
Weighted sampling, and what it actually guarantees
To prefer heavy items, change the key. Efraimidis and Spirakis showed that keeping the k largest values of u1/w is equivalent to drawing items one at a time without replacement, each draw proportional to weight among those remaining. Taking logarithms gives the form in the code above: keep the k smallest values of −ln(u)/w, an exponential random variable with rate w. It stays a bottom-k sample, so it merges and deduplicates exactly like the uniform version.
The common mistake is to assume each item's inclusion probability is k·w/W, where W is the total weight. That holds only for k = 1. For larger k, heavy items saturate towards certainty and light items' chances are not proportional to their weights, so dividing a sampled value by k·w/W to estimate a total gives a biased answer. If you need inclusion proportional to weight, use Chao's method or VarOpt sampling. If you need unbiased sums, use priority sampling (Duffield, Lund and Thorup), which is again a top-k of random keys:
# Priority sampling: estimate sums over the whole stream from k items.
# priority q_i = w_i / u_i; keep the k largest; tau = the (k+1)-th largest priority.
def estimate_total(sample, tau):
return sum(max(w, tau) for _, w in sample) # unbiased for the total weightUse exponential keys for a debugging sample that favours errors, and priority sampling when the sample feeds totals such as bytes by customer. Counting the most frequent items is a different problem again: top-k and heavy hitters.
Windows: sampling the recent past
A reservoir cannot forget. Once an item is in, it stays until a newer item replaces it, so a reservoir over a day's stream is a sample of the whole day, not of the last hour. There are three designs, in rising cost:
- Tumbling windows. Emit the reservoir at each window boundary and start a new one. Simple and exact, but the sample for 10:59 covers only a minute of data.
- Bucketed sketches. Keep one bottom-k sketch per minute and answer "last hour" by merging the latest 60. Because merging is a union, this is exact at minute granularity and costs 60k entries of memory. This is the design most systems should use, and it lines up with the stream processor's windowing model.
- Exact sliding windows. Babcock, Datar and Motwani (2002) keep only candidates that could still become one of the k smallest keys as older items expire. Worth it only when minute granularity is too coarse.
Worked example: a trace sample across 64 partitions
A service emits 109 spans a day over 64 Kafka partitions, about 15.6 million per partition. We want 1,000 traces per hour for debugging, with errors weighted ten times, and the same trace kept in every service that sees it.
Each consumer computes key_weighted(trace_id, w, seed) with w = 10 for error traces and 1 otherwise, using one seed shared across services and rotated monthly. It keeps a bottom-1,000 sketch per partition per minute. Per partition, a minute holds about 10,800 spans, and the heap is touched roughly k(1 + ln(10,800/1,000)) ≈ 3,400 times, so most spans cost a hash and one comparison. Each minute the 64 sketches merge into one of 1,000, and the hourly view merges 60 of those.
Because the key depends only on the trace id and its weight, a trace whose spans land on several partitions gets the same key everywhere, so it is sampled on all of them or on none, the property trace-id-based head sampling relies on. That holds only if the weight is a function of the whole trace; if only some spans carry the error flag, apply the weight where traces are assembled. Error traces are over-represented by design, so any rate computed from this sample must be reweighted or taken from unsampled counters instead.
Operating a sampler: seeds, tests and pitfalls
Sampling bugs pass functional tests because the output always has the right size. Test the distribution: sample 10 of 100 items many times, count each item's selections, and apply a chi-square test against 10 percent. Repeat on merges of unequal shards, which catches the naive merge immediately.
| Pitfall | Symptom | Fix |
|---|---|---|
| Same RNG seed on every shard | Correlated samples; shards keep the same positions | Seed per shard, or use hashed keys of item ids |
| Naive merge of reservoirs | Small shards over-represented | Hypergeometric split, or bottom-k keys |
| Treating exponential keys as inclusion proportional to weight | Biased totals from a weighted sample | Priority sampling or VarOpt for estimates |
| Hash seed never rotated | Adversary can craft ids that are always or never sampled | Keyed hash with a secret seed; rotate on a schedule |
| Low-quality or 32-bit keys | Ties and visible bias at large n | 64-bit hash or 53-bit floats; break ties by id |
What to do next
- Write down which properties your sampler needs from the list at the top; most teams need mergeability and reproducibility more than raw speed.
- If you merge Algorithm R reservoirs today, check whether the merge accounts for shard counts; if not, switch to the hypergeometric merge or to bottom-k keys.
- Replace per-item random draws with a keyed hash of a stable id, and share the seed across services that must agree on what is sampled.
- Decide whether a weighted sample is for inspection or for estimates, and use exponential keys or priority sampling accordingly.
- Build windows from per-minute sketches merged on read rather than resetting a single reservoir.
- Add a chi-square test over repeated runs, including unequal-shard merges, to your test suite.