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.

Advertisement

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.

Advertisement

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.
Reservoir sampling as a system: per-shard bottom-k, mergedpartition 0events, unknown nkey(id)bottom-k sketchk smallest keys + countpartition 1events, unknown nkey(id)bottom-k sketchk smallest keys + countpartition Nevents, unknown nkey(id)bottom-k sketchk smallest keys + countmergek smallest of unionsampleuniform orweighted, k itemskey = hash(id, seed) for uniform; -ln(hash) / weight for weightedper-minute sketches merged on read give windows without evictionsame seed everywhere: duplicates collide, samples are replayable
Each partition keeps a bottom-k sketch of hashed keys; merging keeps the k smallest of the union. Weighted keys and time-bucketed sketches reuse the same structure.
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 out

The 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 weight

Use 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.

PitfallSymptomFix
Same RNG seed on every shardCorrelated samples; shards keep the same positionsSeed per shard, or use hashed keys of item ids
Naive merge of reservoirsSmall shards over-representedHypergeometric split, or bottom-k keys
Treating exponential keys as inclusion proportional to weightBiased totals from a weighted samplePriority sampling or VarOpt for estimates
Hash seed never rotatedAdversary can craft ids that are always or never sampledKeyed hash with a secret seed; rotate on a schedule
Low-quality or 32-bit keysTies and visible bias at large n64-bit hash or 53-bit floats; break ties by id

What to do next

  1. Write down which properties your sampler needs from the list at the top; most teams need mergeability and reproducibility more than raw speed.
  2. 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.
  3. 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.
  4. Decide whether a weighted sample is for inspection or for estimates, and use exponential keys or priority sampling accordingly.
  5. Build windows from per-minute sketches merged on read rather than resetting a single reservoir.
  6. Add a chi-square test over repeated runs, including unequal-shard merges, to your test suite.
Key takeaway: Algorithm R is correct but only covers one stream. For speed on one stream, Algorithm L skips most random draws. For systems, keep the k smallest keys: a keyed hash of a stable id gives a uniform sample that merges by union, deduplicates, replays identically and estimates cardinality; exponential keys add weighting without breaking any of that. Merge existing reservoirs with a hypergeometric split, never a naive pool. Remember that weighted keys do not give inclusion proportional to weight, so use priority sampling for totals. Build windows from bucketed sketches, and test the distribution rather than the output size.