Feature and Data Pipelines in Production

Streaming feature aggregation

Streaming aggregation keeps a running total updated as each new event arrives, instead of recomputing it from scratch by rescanning everything.

On this page 7
  1. The short answer
  2. The analogy you have already lived
  3. Why it exists
  4. How it works
  5. A real example you have seen
  6. Remember this
  7. What to learn next

One lesson, three depths. Pick the one that fits you today — you can switch any time.

Beginner — No maths. Plain English.

The short answer

Streaming aggregation keeps a running total updated as each event arrives, instead of recomputing it from scratch every time.

The analogy you have already lived

Watch a live cricket scoreboard. The run rate updates ball by ball, instantly, without anyone re-adding up the whole match from the first over.

Compare that to a scorer who waits until the innings ends, then goes back through every single ball to compute the total. Both get the same number. One is ready the moment you need it.

Why it exists

Many useful features are running counts: "transactions in the last 5 minutes", "page views in the last hour", "failed logins today". A freshness-conscious pipeline wants these updated the moment a new event arrives.

The slow way rescans every raw event each time the number is needed. As the history grows, that rescan gets slower and slower — exactly backwards from what a live feature needs.

How it works

naive:      every request  ->  scan ALL past events  -> sum
                                (gets slower as history grows)

streaming:  each new event ->  add it to a running total
            old event falls out of window -> subtract it
                                (same tiny amount of work, forever)

The trick is a sliding window. Keep only the events currently inside the time window, add new ones as they enter, remove old ones as they age out. The total is always ready, never recomputed from zero.

A real example you have seen

A card network's "transactions in the last 10 minutes" fraud counter. It cannot afford to re-scan a card's entire history on every swipe — it keeps a small, constantly updated window instead.

Remember this

  • Streaming aggregation updates a total incrementally, as events arrive.
  • The naive alternative — rescanning everything — gets slower as history grows, which is the wrong direction for a live feature.
  • A sliding window keeps only what is currently relevant, adding and removing events as time moves.

What to learn next

  • Low-latency feature lookup — where a streaming aggregate's running total needs to live so a request can read it in milliseconds.
  • Feature freshness — the freshness guarantee a well-run streaming pipeline provides almost for free.
  • Streaming data — the broader discipline of building pipelines that process events as they arrive.

Developer — Code and libraries.

Setup

No installation needed — this uses only Python's standard library collections.deque.

Naive rescan versus a sliding window

streaming_agg.py
import time
from collections import deque
import random

random.seed(0)

WINDOW_SECONDS = 300  # 5-minute window
events = [(t, round(random.uniform(10, 200), 2)) for t in range(20000)]  # one per second


def naive_rescan(all_events, now):
    """Recomputes the 5-minute sum by scanning every event seen so far."""
    return sum(amt for t, amt in all_events if now - WINDOW_SECONDS <= t <= now)


def streaming_window(all_events):
    """Keeps a running sum and a deque of events currently inside the window.
    Each event is added once and evicted once, however long the stream runs."""
    win = deque()
    running_sum = 0.0
    totals = []
    for t, amt in all_events:
        win.append((t, amt))
        running_sum += amt
        while win and t - win[0][0] > WINDOW_SECONDS:
            old_t, old_amt = win.popleft()
            running_sum -= old_amt
        totals.append(running_sum)
    return totals


# naive: recompute from scratch at every step (this is what "rescan the raw
# table on every request" looks like) -- expensive, so only time a slice.
sample_size = 2000
t0 = time.perf_counter()
naive_totals = [naive_rescan(events[:i+1], events[i][0]) for i in range(sample_size)]
naive_time = time.perf_counter() - t0

t0 = time.perf_counter()
stream_totals_full = streaming_window(events)
stream_time_full = time.perf_counter() - t0
stream_totals = stream_totals_full[:sample_size]

agree = all(abs(a - b) < 1e-6 for a, b in zip(naive_totals, stream_totals))

print(f"events processed (naive sample)   : {sample_size}")
print(f"events processed (streaming, full): {len(events)}")
print(f"naive rescan time  (first {sample_size} events)  : {naive_time:.3f}s")
print(f"streaming window time (all {len(events)} events): {stream_time_full:.3f}s")
print(f"both approaches agree on every value: {agree}")
Output
events processed (naive sample)   : 2000
events processed (streaming, full): 20000
naive rescan time  (first 2000 events)  : 0.066s
streaming window time (all 20000 events): 0.003s
both approaches agree on every value: True

These are wall-clock times from one run on one machine — the exact numbers will differ on yours. The comparison that matters: streaming processed ten times more events in twenty times less time. That gap only grows as history grows further.

Line-by-line walkthrough

naive_rescan is given a growing list each time — it is deliberately doing the "rescan everything" approach so the cost is visible. In a real system this would be a database query, not a Python loop, but the shape of the cost is the same.

streaming_window uses a deque as the sliding window. Each event enters exactly once (append) and leaves exactly once (popleft), no matter how long the stream runs — that is what keeps the per-event cost constant instead of growing.

The agree check matters as much as the timings. A faster wrong answer is worse than a slow right one — always verify both approaches produce the same numbers before trusting the faster one.

Common mistakes

Rebuilding the window from a raw event table on every single request. This is the naive path above, running in production. It works fine in a demo with a hundred rows and falls over at real volume.

Forgetting to evict. A window that only grows, and never removes events that aged out, is not a sliding window — it is a slow leak that eventually reports the total for "all time" instead of "the last 5 minutes".

Building this yourself when a stream processor already does it. Kafka Streams, Flink and Spark Structured Streaming all provide windowed aggregation as a built-in primitive, tested against real out-of-order and late data. Reach for one before hand-rolling this at production scale.

Try it yourself

Change WINDOW_SECONDS to 60 and re-run. The streaming approach's total work barely changes, because it was already only ever touching events inside the current window — confirm that by timing it.

What to learn next

  • Low-latency feature lookup — where a streaming aggregate's running total needs to live so a request can read it in milliseconds.
  • Feature freshness — the freshness guarantee a well-run streaming pipeline provides almost for free.
  • Streaming data — the broader discipline of building pipelines that process events as they arrive.

Researcher — Mathematics and papers.

Amortised cost

For a stream of $n$ events with a bounded window of $w$ time units and roughly $k$ events per window, the naive rescan costs $O(n \times k)$ in the worst case — quadratic when $k$ grows with $n$. The sliding-window accumulator costs $O(n)$ total: each event is inserted once and evicted at most once, giving $O(1)$ amortised cost per event regardless of window size.

This is the same amortised-analysis argument used for any monotonic-pointer or two-pointer technique: total work is bounded by the number of insertions plus removals, not by the number of times the window is queried.

Beyond sum and count: algebraic and holdable aggregates

Sum, count and mean are invertible — removing an element is a direct subtraction, which is what makes the deque-based eviction above trivial. Not every useful aggregate is invertible:

  • max/min over a window needs a monotonic deque (keep only candidates that could still be the max as the window slides), not a plain running value, because the maximum cannot be "un-computed" by subtraction when it leaves the window.
  • distinct count (how many unique users in the window) needs a sketch structure — HyperLogLog is standard — because exact distinct counting over a sliding window is memory-expensive at scale.
  • percentiles need a mergeable summary (t-digest, KLL sketch), since no simple running statistic captures a quantile exactly.

Choosing the wrong data structure for a non-invertible aggregate is a common production bug: a naive running-max that never handles eviction correctly will silently report a stale maximum long after the event that set it has left the window.

Watermarks and late data

Real streams are not perfectly ordered. An event with timestamp $t$ can arrive after events with timestamp $t' > t$, due to network delay or retries. Correctly windowed aggregation needs a watermark: a declared bound on how late an event can arrive before it is dropped or triggers a correction to an already-emitted window result.

Flink and Spark Structured Streaming both implement the Dataflow model's watermark-and-trigger semantics directly. Choosing a watermark too tight drops legitimately late data; too loose delays every window's result and inflates state size, since more history must be kept "open" for possible correction.

Papers and systems

  • Akidau et al., The Dataflow Model, VLDB 2015 — the formal treatment of windows, watermarks and triggers this lesson's window is a simplified instance of.
  • Heule, Nunkesser and Hall, HyperLogLog in Practice, EDBT 2013 — the sketch behind streaming distinct-count.
  • Kreps, I heart logs, O'Reilly 2014 — why treating the event stream itself as the source of truth simplifies exactly this class of problem.

What to learn next

  • Low-latency feature lookup — where a streaming aggregate's running total needs to live so a request can read it in milliseconds.
  • Feature freshness — the freshness guarantee a well-run streaming pipeline provides almost for free.
  • Streaming data — the broader discipline of building pipelines that process events as they arrive.