A streaming algorithm reads its input once, in arrival order, and keeps a summary far smaller than the input. That restriction is not academic. A router sees tens of millions of packets a second and can keep a few megabytes of state per counter family. A metrics pipeline sees billions of events a day and must answer "which keys are hot" and "what is the p99" without storing the events. In both, the data is too large or too fast to store, and you accept an approximate answer with a known error bound in exchange for memory that does not grow with the stream.

This article covers the model itself and the toolkit around it. It explains the three stream models, why some questions cannot be answered exactly or deterministically in small space, and three algorithms worth knowing by heart: Misra-Gries for heavy hitters, the AMS sketch for the second frequency moment, and DGIM for counting over a sliding window. The best-known sketches have their own pages here: Count-Min Sketch, HyperLogLog and reservoir sampling.

The streaming model

Formally, a stream is a sequence of updates to a frequency vector f over a universe of n possible items. After m updates you want some function of f: the number of distinct items, the most frequent items, a quantile, or a sum of squares. A streaming algorithm is judged on three resources: space, measured in bits as a function of n and m; time per update; and passes over the data, which is almost always one.

The target is polylogarithmic space, O(log n + log m) bits per counter times a number of counters that depends only on the error you accept. If the space grows with n or m you have not built a streaming algorithm, you have built a cache.

Three models differ in which updates are allowed. In the insert-only (cash register) model every update adds one occurrence of an item: page views, packets, log lines. In the turnstile model updates may be negative, so f[i] can go up and down: inventory, open connections, a difference between two streams. Sketches that are linear functions of f, such as Count-Min and AMS, handle the turnstile model naturally, because a deletion is just an update with weight -1. Counter-based algorithms such as Misra-Gries do not. In the sliding-window model only the last W items or the last W seconds matter, and expiring old items without storing them is the whole difficulty.

What small space cannot do

Two impossibility results shape everything else. First, exact answers need linear space for most useful questions. To report the exact number of distinct items, an algorithm must be able to tell, at every point, whether the next item is new, which means its state must distinguish every possible set of items seen so far. A counting argument shows that takes on the order of min(n, m) bits.

Second, determinism is expensive. Alon, Matias and Szegedy showed in 1996 that a deterministic algorithm that even approximates the number of distinct items needs linear space, while a randomized one can do it in logarithmic space. So the standard contract for a streaming algorithm is an (epsilon, delta) guarantee: the estimate is within a factor (1 plus or minus epsilon), or within an additive epsilon times m, of the truth, with probability at least 1 minus delta. Space grows with 1/epsilon or 1/epsilon squared and only with log(1/delta), which is why driving delta to one in a million is cheap and driving epsilon to 0.1% is not.

Heavy hitters are the exception that proves the rule: Misra-Gries is deterministic, because it only promises an additive error of m/(k+1), which is large for rare items and irrelevant for frequent ones.

Heavy hitters with Misra-Gries

The question: which items occur more than a fraction phi of the time? Misra-Gries answers it with k counters. When an item arrives, increment its counter if it has one; otherwise claim a free counter; if none is free, decrement every counter by one and drop the ones that reach zero.

def misra_gries(stream, k):
    """Keep at most k counters. Each estimate is low by at most m/(k+1)."""
    counters = {}
    for x in stream:
        if x in counters:
            counters[x] += 1
        elif len(counters) < k:
            counters[x] = 1
        else:
            for key in list(counters):          # the decrement step
                counters[key] -= 1
                if counters[key] == 0:
                    del counters[key]
    return counters

def mg_merge(a, b, k):
    """Mergeable: add, then subtract the (k+1)-th largest count and keep positives."""
    total = dict(a)
    for key, v in b.items():
        total[key] = total.get(key, 0) + v
    if len(total) <= k:
        return total
    cut = sorted(total.values(), reverse=True)[k]
    return {key: v - cut for key, v in total.items() if v > cut}

Why it works: each decrement step removes one from k counters and also discards the arriving item, so it destroys k+1 occurrences at once. There can be at most m/(k+1) such steps, and an item's counter is low by at most the number of decrement steps. So every estimate satisfies f[x] - m/(k+1) at most estimate, estimate at most f[x], and any item with true frequency above m/(k+1) is guaranteed to still hold a counter. To find every item above 1% of traffic, set k = 100.

Worked example with k = 2 on the stream a b a c a b d a. After a, b, a the counters are {a: 2, b: 1}. c arrives with no free counter, so everything drops by one: {a: 1}. a makes it {a: 2}, b claims the free slot {a: 2, b: 1}, d triggers a decrement: {a: 1}, and the final a gives {a: 2}. The true count of a is 4, m = 8, and the bound says the estimate is low by at most 8/3, so at most 2 here; it is low by exactly 2. a is the only item above m/3, and it survived.

The output is a candidate set, not an answer. Items near the threshold may be false positives, so a second pass or an exact counter on the candidates confirms them when that matters. Space-Saving, the common production variant, replaces the minimum counter instead and overestimates rather than underestimates.

The AMS sketch and median of means

The second frequency moment F2 is the sum of f[i] squared. It measures skew, it is the self-join size a query planner needs, and it is the squared length of the frequency vector. The AMS or "tug-of-war" sketch estimates it with one counter. Hash each item to a random sign s(i) of +1 or -1 and keep Z = sum of s(i) f[i]. Expanding Z squared, the cross terms s(i) s(j) f[i] f[j] average to zero over the random signs, leaving E[Z squared] = F2. With four-wise independent signs the variance is at most 2 F2 squared.

import hashlib, statistics

def sign(item, seed):
    h = hashlib.blake2b(f"{seed}:{item}".encode(), digest_size=8).digest()
    return 1 if h[0] & 1 else -1                # stand-in for a 4-wise independent hash

class AMS:
    def __init__(self, groups=7, per_group=64):
        self.g, self.r = groups, per_group
        self.z = [0] * (groups * per_group)

    def update(self, item, weight=1):           # weight may be negative: turnstile-safe
        for j in range(len(self.z)):
            self.z[j] += weight * sign(item, j)

    def f2(self):
        sq = [v * v for v in self.z]
        means = [statistics.fmean(sq[i*self.r:(i+1)*self.r]) for i in range(self.g)]
        return statistics.median(means)          # median of means

One counter has a standard deviation comparable to the answer, so you average r copies to cut the variance by r: with r about 8/epsilon squared, Chebyshev puts each mean within (1 plus or minus epsilon) of F2 with probability at least 3/4. Then take the median of g such means; the median is wrong only if half of the means are, which a Chernoff bound makes exponentially unlikely in g. This median-of-means pattern, average for accuracy and median for confidence, is how almost every randomized sketch turns a weak guarantee into an (epsilon, delta) one.

Because Z is linear in f, two AMS sketches built with the same seeds add, and subtracting them gives the sketch of the difference of two streams. That is the turnstile model in action.

Sliding windows with DGIM

Counting the 1s among the last W bits of a stream looks like it needs W bits. DGIM, from Datar, Gionis, Indyk and Motwani, does it approximately in O(log squared W) bits. It keeps buckets, each recording the timestamp of its most recent 1 and a size that is a power of two. New 1s start a size-1 bucket. Whenever three buckets share a size, the two oldest merge into one of twice the size. Buckets whose timestamp has left the window are dropped.

To estimate the count, add the sizes of all buckets except the oldest and add half the oldest. Only the oldest bucket straddles the window boundary, so the error is at most half of it, and because there is at least one bucket of every smaller size the true count is at least as big as that bucket. The result is a relative error of at most 50%, and allowing r buckets per size instead of two shrinks it to about 1/r at proportionally more space.

In practice many systems take a cheaper route to windows: keep one mergeable sketch per time slice, for example per minute, and merge the last sixty to answer "last hour". That costs sixty sketches instead of one but works with any mergeable summary, and it is the pattern behind most stream-processor windowing.

Mergeable summaries in a distributed pipeline

A sketch is mergeable if the summary of the union of two streams can be computed from the two summaries alone, with the same error guarantee as if one sketch had seen everything. Linear sketches (Count-Min, Count Sketch, AMS) merge by adding arrays. HyperLogLog merges by taking the register-wise maximum. Misra-Gries merges with the subtract-the-(k+1)-th rule above. Quantile summaries such as KLL and t-digest are designed to merge.

Mergeability is what makes sketches fit distributed systems. Each partition or worker keeps its own sketch, sketches are serialized per time window, and queries merge whatever range they need. Two rules make it safe: every sketch that will ever be merged must be built with the same parameters and hash seeds, and the serialized form must carry those parameters so a mismatched merge is rejected instead of silently producing garbage.

A streaming pipeline: sketch per partition, merge, queryEvent sourcesclicks, logs, flowsPartition 1sketch S1Partition 2sketch S2Partition psketch SpMergeS = S1 + S2 + ... + SpQueryestimate plus error boundSerialized sketch storeper window: bytes, not eventspersistMemory per partition is fixed by the error target, not by traffic volume.The merge must be exact for the error bound to survive the fan-in.
Each partition updates a fixed-size sketch; window sketches are stored as bytes and merged at query time. Memory depends on the error target, not on traffic.

Choosing a sketch

QuestionAlgorithmSpace for error epsilonMergeTurnstile
Distinct countHyperLogLogabout (1.04/epsilon)^2 small registersregister maxno
Frequency of one keyCount-Min(e/epsilon) x ln(1/delta) countersadd arraysyes
Heavy hittersMisra-Gries / Space-Saving1/epsilon countersadd, then cutno
F2, self-join sizeAMSO(1/epsilon^2 x log(1/delta)) countersadd arraysyes
QuantilesKLL, t-digest, GKroughly 1/epsilon items, up to log factorsyes (GK is awkward)no
Uniform sampleReservoirk itemsweighted mergeno
Count in last WDGIMO(log^2 W) bitsper-slice insteadno

Read the error column carefully. Count-Min and Misra-Gries promise additive error epsilon times m, which is excellent for heavy keys and useless for rare ones: with m = 10^9 and epsilon = 0.001 every key's estimate may be off by a million. HyperLogLog and AMS promise relative error. Quantile sketches promise rank error, so a p99.9 from a sketch with epsilon = 0.01 could really be anywhere from p98.9 up; that is why tail-focused designs such as t-digest and relative-error KLL variants exist.

Failure modes

  • Adversarial or correlated input. The guarantees assume the hash is independent of the data. Keys crafted to collide, or a weak hash on structured keys such as sequential IDs, break them. Use a seeded, well-mixed hash and keep the seed private if inputs are user-controlled.
  • Mismatched parameters at merge time. Merging sketches with different widths, precisions or seeds produces numbers that look plausible and mean nothing. Version the serialized format and check it.
  • Additive error read as relative. A Count-Min estimate of 3,000 for a key that truly had 0 is within spec at large m. Report the bound alongside the estimate, or confirm candidates exactly.
  • Deletions in an insert-only sketch. HyperLogLog and Misra-Gries cannot unsee an item. If the domain has retractions, use a linear sketch or recompute per window.
  • Window boundaries. Per-slice merging answers "last 60 full minutes", not "last 3,600 seconds". Say which one a dashboard shows.

Testing a sketch

Test a sketch statistically, not with one example. Generate streams with known answers, including a uniform stream, a Zipf stream with exponent near 1.1 and an adversarial one, run many independent seeds, and check that the observed error stays inside the claimed bound at the claimed rate. A Misra-Gries test can assert the bound on every run, because it is deterministic:

import random
from collections import Counter

def test_misra_gries(trials=200, m=20_000, k=50):
    for t in range(trials):
        rng = random.Random(t)
        stream = [int(rng.paretovariate(1.1)) for _ in range(m)]
        truth, est = Counter(stream), misra_gries(stream, k)
        for x, f in truth.items():
            e = est.get(x, 0)
            assert f - m / (k + 1) <= e <= f, (t, x, f, e)

test_misra_gries()

For randomized sketches assert a rate instead: across 1,000 seeds, the fraction of runs outside (1 plus or minus epsilon) should be below delta. Also test merge: the merged sketch of two halves must meet the same bound as one sketch over the whole stream.

What to do next

  1. Write down the question, the stream model (insert-only, turnstile or windowed) and whether you need additive, relative or rank error before choosing a sketch.
  2. Size the sketch from its error formula and your m, and compute the absolute error you will actually see.
  3. Implement Misra-Gries and the test above to build intuition, then use a maintained library such as Apache DataSketches in production.
  4. Fix hash, seed and parameters per sketch family and store them in every serialized sketch.
  5. Keep sketches per partition and per time slice, and merge at query time.
  6. Run the statistical test harness across seeds before trusting a dashboard built on it.
Key takeaway: A streaming algorithm buys memory that is independent of stream length with a stated error and failure probability. Choose by question, stream model and error type; size from the error formula; build mergeable sketches with fixed seeds and parameters; and verify the bound statistically before you rely on it.