A HyperLogLog sketch answers one question, how many distinct items have I seen, in a few kilobytes with an error of one or two percent. The articles on the estimator and precision and on registers, sparse encoding and merge explain the sketch itself. This one is about what happens when you put it inside a stream processor: millions of keys, events arriving out of order and more than once, windows that slide, and questions like how many distinct ports has this host touched in the last five minutes that have to be answered continuously.

Streams change the engineering more than the mathematics. The sketch is the same; what changes is where its state lives, when a bucket is closed, how a window is assembled from pieces, what a replayed event does, and what you can never undo. By the end you should be able to design a windowed distinct-count pipeline, size its state, choose between a ring of buckets and a sliding sketch, and build a fan-out detector on top of it.

Why max-merge suits streams

Adding an item hashes it, uses the first p bits to pick one of m = 2p registers, and stores the position of the first one-bit in the rest of the hash if it is larger than what the register holds. Merging two sketches takes the register-wise maximum. Maximum has three properties that make it unusually friendly to streams.

It is commutative and associative, so events can be processed in any order and partitions can be merged in any tree shape: per-core sketches merge into per-task sketches, which merge into per-minute sketches, with the same result as one sketch fed serially. It is idempotent: adding the same item twice changes nothing, so an at-least-once source that replays a batch after a crash cannot inflate the count. Exactly-once delivery, which is expensive, buys you nothing for this metric. And merged sketches have the same error as a single sketch of the same precision over the union, so assembling an hour from sixty minutes costs no accuracy.

The flip side is that maximum has no inverse. A register that reached rank 9 because of one item cannot be lowered when that item is deleted, because the register does not know which item set it. Everything in this article about windows, retractions and privacy follows from that one missing operation.

The pipeline architecture

Distinct counts in a stream: sketches are state, numbers are viewsEvent sourceat-least-oncePartition by keye.g. source IPKeyed sketch statekey x minute bucketadd = hash, take max per registerwatermarkClosed bucketsserialized sketchesSketch storekey, bucket, bytesmerge on readQuery layerunion of N bucketsAlerts, dashboardsestimate + errorSliding sketchtime, rank listssame addsany window WLate events land in their own bucket until the watermark closes it; replays are harmless because max is idempotentDeletes are not possible: a register cannot forget which item set it
Events are partitioned by key, folded into per-key, per-bucket sketches, closed by the watermark and stored as bytes; queries merge buckets on read. A sliding sketch fed by the same adds can answer any window length up to its retention.

The architecture has one rule worth stating plainly: store sketches, not numbers. A stored count of 41,000 distinct users for 10:00 and 39,000 for 10:01 cannot tell you the distinct users across both minutes; the stored sketches can, by merging. Every downstream rollup, every window and every dimension you might later slice by needs the sketch bytes. The number is a view computed at read time.

Tumbling windows: a ring of bucket sketches

The simplest window is a tumbling bucket, one sketch per key per minute. A window of the last W minutes is the merge of the last W buckets. Keeping them in a ring indexed by minute modulo W bounds memory and makes expiry free: a bucket is reset when its slot is reused by a newer minute.

class RingWindow:
    """Last-N-buckets distinct count from per-bucket HLL sketches."""
    def __init__(self, p, buckets):
        self.p, self.n = p, buckets
        self.ring = [HLL(p) for _ in range(buckets)]
        self.bucket_of = [None] * buckets      # which minute each slot holds

    def add(self, item, minute):
        slot = minute % self.n
        if self.bucket_of[slot] != minute:     # slot holds an old minute: recycle it
            self.ring[slot] = HLL(self.p)
            self.bucket_of[slot] = minute
        self.ring[slot].add(item)

    def count(self, now_minute):
        acc = HLL(self.p)
        for slot in range(self.n):
            b = self.bucket_of[slot]
            if b is not None and now_minute - self.n < b <= now_minute:
                acc.merge(self.ring[slot])
        return acc.count()

The query cost is W merges of m registers, so a 60-minute window at p = 12 touches about 245,000 registers. That is cheap for a dashboard and too slow for a per-event check on a million keys. Hierarchical buckets (minutes roll into hours) and caching the merged window until a bucket closes both help. The ring also fixes the window length at build time and its edge at bucket granularity.

Sliding windows: keeping possible future maxima

A sliding HyperLogLog, described by Chabchoub and Hebrail, answers any window length up to a retention limit from one structure. Instead of a single rank, each register keeps a short list of (timestamp, rank) pairs that could still be the register's maximum for some future window. When a new pair arrives, every older pair with a rank no larger than the new one is removed, because for any window containing the old pair the newer, larger-or-equal pair is also inside it. What remains is ordered: oldest first with the highest rank, newest last with the lowest.

from collections import deque

class SlidingHLL:
    """Per register, keep only (time, rank) pairs that could still be a window max."""
    def __init__(self, p):
        self.p = p
        self.lists = [deque() for _ in range(1 << p)]

    def add(self, item, t):
        i, r = split(item, self.p)             # register index, rank
        q = self.lists[i]
        while q and q[-1][1] <= r:             # dominated: older and not larger
            q.pop()
        q.append((t, r))

    def expire(self, oldest):
        for q in self.lists:                   # drop pairs no query will ever reach
            while q and q[0][0] < oldest:
                q.popleft()

    def count(self, now, window):
        # the oldest surviving pair inside the window has the largest rank
        regs = [next((r for t, r in q if t > now - window), 0) for q in self.lists]
        return estimate(regs)

Call expire periodically with the current time minus the longest window you will ever query; that retention limit is the only sizing decision. The expected list length per register grows roughly with the logarithm of the number of distinct items per register, so memory is a small multiple of a plain sketch rather than a multiple of the window length.

In a scratch run with p = 12 (4,096 registers, theoretical standard error 1.63%), 2,000 events per minute for 180 minutes and the population shrinking after minute 120, a 60-minute window queried every ten minutes gave estimates identical to the ring of per-minute sketches at every check, with a worst relative error of 1.40% against the exact count. The sliding structure held between 10,200 and 11,900 pairs, about 2.9 per register. So at bucket granularity the two designs agree exactly; the sliding sketch buys arbitrary window lengths and edges at a few times the memory, and the ring buys simplicity and cheap serialization.

Keyed state and what it costs

Keyed state is where streaming HLL gets expensive. A dense sketch at 6 bits per register costs 0.75 × 2p bytes, and you hold one per key per open bucket.

Precision pRegistersDense bytesStandard error10 million keys
101,0247683.25%7.7 GB
124,0963,0721.63%30.7 GB
1416,38412,2880.81%123 GB

Most keys in real streams are small: a typical source IP touches a handful of destinations, a typical user visits a few pages. So production pipelines start each key as an exact set or a sparse list of (index, rank) pairs and promote it to dense registers only when the sparse form would be larger. With a heavy-tailed key distribution this cuts state by one to two orders of magnitude. Choose p per question, not per pipeline: a fan-out detector that only needs to tell 20 from 2,000 is well served by p = 8 or 10, while billing-grade unique counts want 14.

Two more state rules. Use a fast, well-mixed 64-bit hash and the same hash and seed everywhere; sketches built with different hashes merge without error and produce garbage. And version the serialized format, because the store will outlive the code that wrote it.

Late data, replays and retractions

Event time and processing time differ, so a bucket for 10:00 must stay open until the watermark says no more 10:00 events are expected. Because adds are idempotent and order-free, a late event that arrives inside the allowed lateness is simply added to its own bucket, and the stored sketch for that bucket is overwritten with the updated bytes. Downstream merged windows must be invalidated or recomputed, which is a caching problem, not a correctness one.

Events later than the allowed lateness have two honest options: drop them and count the drops, or add them to the stored bucket sketch by read-merge-write, which the max-merge makes safe to retry; against a concurrent writer of the same bucket, use compare-and-set or a store that merges sketches natively. Never add a late event to the current bucket; that silently moves it into the wrong window.

Retractions are where the missing inverse bites. A stream that emits corrections, such as a changelog where a row is deleted, cannot apply them to a sketch. If a deletion must take effect, for example a privacy request, you need either short bucket lifetimes so the item ages out, or the raw data to rebuild affected buckets. Plan the retention of sketch bytes with the same care as the retention of raw logs, because a sketch is derived from personal identifiers even if it does not expose them directly.

Worked example: a fan-out scan detector

Distinct counts per key are a classic security and abuse signal. A host scanning a network touches many distinct destination ports or addresses; a stolen credential logs in from many distinct devices; a scraping client requests many distinct URLs. Total request volume does not separate these from heavy but normal users. Fan-out, the number of distinct targets per source, does.

SCAN_THRESHOLD = 500          # distinct (dst_ip, dst_port) in 5 minutes
P_DETECT = 10                 # 3.25% standard error is plenty for 20 versus 2,000

state = {}                    # src_ip -> RingWindow(P_DETECT, buckets=5)

def on_flow(src_ip, dst_ip, dst_port, minute):
    win = state.get(src_ip)
    if win is None:
        win = state[src_ip] = RingWindow(P_DETECT, buckets=5)
    win.add((dst_ip, dst_port), minute)

def on_bucket_close(minute):
    for src_ip, win in state.items():
        fanout = win.count(minute)
        if fanout > SCAN_THRESHOLD * 1.1:      # margin of about 3 standard errors
            emit_alert(src_ip, minute, round(fanout))

Walk one source through it. A web server's health checker hits 4 targets every second: millions of flows, fan-out 4, never alerts. A compromised laptop probes 1,024 ports on each of 3 hosts in two minutes: a few thousand flows, fan-out about 3,072, alerts at the next bucket close. Volume would have ranked the health checker first. The 10% margin above the threshold keeps estimates sitting right at 500 from flapping; with a 3.25% standard error, three standard errors is just under 10%.

Serving and reporting the counts

For serving, many stores keep HLL natively. Redis has PFADD, PFCOUNT and PFMERGE; PFCOUNT over several keys returns the cardinality of their union, its sketches use a sparse form for small sets and at most about 12 KB dense, with a documented standard error of 0.81%. BigQuery has HLL_COUNT.INIT, HLL_COUNT.MERGE_PARTIAL and HLL_COUNT.EXTRACT for storing and combining sketch columns. These formats are generally not interchangeable with each other or with your stream processor's, so pick one canonical encoding for stored buckets and convert at the edges.

When you report a number, report its uncertainty with it: an estimate of 51,850 with a 1.6% standard error is 51,850 plus or minus about 830. Differences between windows, such as new users today computed as a union minus yesterday, carry the error of both terms and can be meaningless when the difference is small; for those, other sketches or exact sampling via reservoir sampling may fit better.

Failure modes

  • Storing counts instead of sketches. Windows and rollups become impossible; teams then sum per-minute distinct counts and report numbers that can be sixty times too large.
  • Mixed hash functions or seeds. Two services sketch the same users with different hashes; the merged count roughly doubles with no error raised.
  • Dense sketches for every key. State grows to hundreds of gigabytes, checkpoints slow down, and the job falls behind; start keys sparse.
  • Late events added to the current bucket. Counts drift between windows and replays of the same input give different dashboards.
  • Expecting deletes to work. A retraction stream is silently ignored or, worse, implemented as a second sketch that is subtracted, producing negative counts.
  • Alerting on raw estimates at the threshold. Estimates within a few standard errors of the threshold flap; add margin or require two consecutive buckets.
  • Small-range bias ignored. Below about 2.5 m items the raw estimator is biased; the linear-counting correction or the sparse exact path must be in place.

Trade-offs

A ring of tumbling buckets is simple, compact and easy to serialize, but fixes the window grid. A sliding sketch answers any window at a few times the memory and with a more complex encoding. Higher precision quarters the variance for four times the bytes per key. Exact sets are free of error and cheaper for small keys, so the hybrid of exact-then-sparse-then-dense is almost always right.

What to do next

  1. List every distinct-count metric you compute in streams and confirm each stores sketch bytes, not numbers.
  2. Fix one hash function, seed and serialization version across all producers and consumers.
  3. Choose p per question from the error you can tolerate, and size keyed state with the table above.
  4. Start keys as exact or sparse sets and measure the fraction that ever promote to dense.
  5. Define bucket width, allowed lateness and the policy for events later than that.
  6. Pick a ring for fixed windows or a sliding sketch for arbitrary ones, and test it against an exact count.
  7. Add error bars to every dashboard and alert margins of about three standard errors.
  8. Document that deletions cannot be applied and set sketch retention accordingly.
Key takeaway: In a stream, a HyperLogLog sketch is state you merge, not a number you store: max-merge makes it order-free and replay-safe, rings or sliding sketches give windows, sparse-first keys keep state small, and the one thing it can never do is forget an item.