Data Engineering for AI

Streaming data

Streaming means handling data that never stops arriving, where events show up late and out of order, so any number you report is only true so far.

On this page 9
  1. Batch and streaming, side by side
  2. The two times every event has
  3. Windows
  4. How long do you wait?
  5. What this does to a model
  6. Somewhere you have seen this
  7. What is honestly hard here
  8. Remember this
  9. 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.

Streaming means working with data that keeps arriving, one event at a time. There is no finished file to open.

Think about standing at a railway platform counting people getting off a train. You cannot wait until everybody in India has finished travelling. You count as they come.

Now somebody who got off at the far end walks over five minutes later. Do you add them to your count? Your count was already announced.

That is the whole problem of streaming, and it never fully goes away.

Batch and streaming, side by side

Batch is a file that has stopped growing. You read the whole thing and compute an answer. The answer is final.

Streaming never stops. You compute answers over slices of time as events arrive, and a later arrival can change a slice you already reported.

The two times every event has

This is the idea to hold on to, and it takes most people two readings.

Event time is when the thing happened. You tapped to pay at 10:04.

Arrival time is when your server found out. Your phone was in a lift with no signal, so the payment reached the server at 10:12.

Those are eight minutes apart, and they lead to different answers.

If you count payments by arrival time, your 10:00 to 10:05 report is about your network, not about your customers. If you count by event time, the report is about reality — but you have to wait for stragglers.

Windows

Since a stream has no end, you cut it into pieces. A five-minute window groups every event whose event time falls in those five minutes.

   10:00 |-----| 10:05 |-----| 10:10 |-----| 10:15
         window 1       window 2       window 3

   a payment made at 10:04 belongs to window 1
   even if it reaches you at 10:12

How long do you wait?

Wait a long time, and your dashboard is always behind. Wait a short time, and you announce numbers that later turn out wrong.

The chosen waiting period is called a watermark: a promise that you will not expect events older than this any more. Picking it is a business decision about how much lateness you tolerate for how much delay.

There is no correct answer. There is only the trade you chose.

What this does to a model

A model trained on stream data must use only what was known at decision time. Say a label came from a window that later received three more events. Your training data now holds a history production will never see.

This connects directly to feature stores, where the same rule appears from a different angle.

Somewhere you have seen this

Live cricket scores that correct themselves after a review. Or a payment app showing a balance that updates a few seconds later.

Neither is a bug. Both are systems admitting that a number is true so far.

What is honestly hard here

You will never know whether a late event is still coming. You can only decide how long to wait.

Testing is hard too. Bugs appear only in specific orderings, which are hard to reproduce. A pipeline that ran perfectly for six months fails when a mobile network is slow.

Anyone who tells you streaming is batch with smaller files has not run one at 2am.

Remember this

  • Every event has an event time and an arrival time. They differ, and it matters.
  • Group by event time. Grouping by arrival time measures your network, not your users.
  • A watermark is how long you wait for stragglers. It is a trade, not a solution.

What to learn next

Developer — Code and libraries.

Nine payments, three windows, two of them late. Pure pandas, no broker, no cluster. This is the smallest thing that shows the real problem.

Setup

bash
pip install pandas

Event time, arrival time, and a watermark

windows.py
import pandas as pd

def t(s): return pd.Timestamp("2026-08-20 " + s)

# Nine payments. event_time is when the tap happened; arrival_time is when our server saw it.
events = pd.DataFrame([
    ("p1", t("10:01:00"), t("10:01:02"), 250.0),
    ("p2", t("10:03:30"), t("10:03:33"), 120.0),
    ("p3", t("10:04:50"), t("10:12:40"), 900.0),   # phone was offline in a lift
    ("p4", t("10:06:10"), t("10:06:12"),  75.0),
    ("p5", t("10:07:00"), t("10:07:01"), 310.0),
    ("p6", t("10:09:59"), t("10:10:05"),  60.0),   # crossed the boundary in transit
    ("p7", t("10:11:20"), t("10:11:22"), 480.0),
    ("p8", t("10:13:00"), t("10:13:01"), 150.0),
    ("p9", t("10:02:10"), t("10:31:00"), 700.0),   # 29 minutes late
], columns=["id", "event_time", "arrival_time", "amount"])

def windows(df, time_col):
    w = df.groupby(pd.Grouper(key=time_col, freq="5min"))["amount"].agg(["count", "sum"])
    return w[w["count"] > 0]

print("windowed by EVENT time (what actually happened):")
print(windows(events, "event_time"))
print()
print("windowed by ARRIVAL time (what our server felt like):")
print(windows(events, "arrival_time"))
print()

WATERMARK = pd.Timedelta(minutes=5)

def accumulating(now):
    """Re-report every open window from everything seen so far. Late data revises old answers."""
    seen = events[events["arrival_time"] <= now]
    out = []
    for start, grp in seen.groupby(pd.Grouper(key="event_time", freq="5min")):
        if len(grp) and now >= start + pd.Timedelta(minutes=5) + WATERMARK:
            out.append((str(start.time()), len(grp), float(grp["amount"].sum())))
    return out

def closed(now):
    """Freeze each window when the watermark passes. Anything later is late data."""
    out, late = [], []
    for start, grp in events.groupby(pd.Grouper(key="event_time", freq="5min")):
        if not len(grp):
            continue
        shuts = start + pd.Timedelta(minutes=5) + WATERMARK
        if now < shuts:
            continue
        in_time = grp[grp["arrival_time"] <= shuts]
        out.append((str(start.time()), len(in_time), float(in_time["amount"].sum())))
        late += list(grp[(grp["arrival_time"] > shuts) & (grp["arrival_time"] <= now)]["id"])
    return out, late

print("ACCUMULATING mode - the window is re-reported as late events land:")
for now in ["10:10:00", "10:15:00", "10:20:00", "10:40:00"]:
    print(f"  at {now}: {accumulating(t(now))}")
print()
print("CLOSED mode - the window is frozen at the watermark, late events are dropped:")
for now in ["10:10:00", "10:15:00", "10:20:00", "10:40:00"]:
    rows, late = closed(t(now))
    print(f"  at {now}: {rows}  dropped as late: {late}")
Output
windowed by EVENT time (what actually happened):
                     count     sum
event_time                        
2026-08-20 10:00:00      4  1970.0
2026-08-20 10:05:00      3   445.0
2026-08-20 10:10:00      2   630.0

windowed by ARRIVAL time (what our server felt like):
                     count     sum
arrival_time                      
2026-08-20 10:00:00      2   370.0
2026-08-20 10:05:00      2   385.0
2026-08-20 10:10:00      4  1590.0
2026-08-20 10:30:00      1   700.0

ACCUMULATING mode - the window is re-reported as late events land:
  at 10:10:00: [('10:00:00', 2, 370.0)]
  at 10:15:00: [('10:00:00', 3, 1270.0), ('10:05:00', 3, 445.0)]
  at 10:20:00: [('10:00:00', 3, 1270.0), ('10:05:00', 3, 445.0), ('10:10:00', 2, 630.0)]
  at 10:40:00: [('10:00:00', 4, 1970.0), ('10:05:00', 3, 445.0), ('10:10:00', 2, 630.0)]

CLOSED mode - the window is frozen at the watermark, late events are dropped:
  at 10:10:00: [('10:00:00', 2, 370.0)]  dropped as late: []
  at 10:15:00: [('10:00:00', 2, 370.0), ('10:05:00', 3, 445.0)]  dropped as late: ['p3']
  at 10:20:00: [('10:00:00', 2, 370.0), ('10:05:00', 3, 445.0), ('10:10:00', 2, 630.0)]  dropped as late: ['p3']
  at 10:40:00: [('10:00:00', 2, 370.0), ('10:05:00', 3, 445.0), ('10:10:00', 2, 630.0)]  dropped as late: ['p9', 'p3']

Compare the first two tables

Event time says the 10:00 window held 4 payments totalling 1970.0. Arrival time says 2 payments totalling 370.0.

Worse, arrival time invents a window at 10:30 that corresponds to nothing that happened. Payment p9 occurred at 10:02 and is filed under half past ten because of a network delay.

Grouping by arrival time is the default in most hand-written code, because arrival time is the timestamp you have lying around. It measures your infrastructure and reports it as customer behaviour.

Now watch the same window change three times

Follow the 10:00:00 window down the accumulating lines:

reported atcountsum
10:102370.0
10:1531270.0
10:2031270.0
10:4041970.0

Same window, four reports, three different answers. Every one was correct given what had arrived.

This is the fact to take away: a streaming number is not a fact, it is a fact so far. Any consumer of that number — a dashboard, an alert, a training set — has to tolerate revision, or it will be wrong and confident.

The other mode, and what it costs

Accumulating means old answers change. Some sinks cannot accept that: an append-only ledger, a bill you have already sent, a training label already written. So the alternative is to freeze the window when the watermark passes and discard whatever comes later.

Look at what closed mode reports for the 10:00 window: 2, 370.0, permanently. The truth was 4, 1970.0. Eighty-one percent of the money in that window was thrown away, and p3 and p9 appear in the dropped list.

Neither mode is correct. Accumulating gives right answers that keep changing. Closed gives stable answers that are wrong. You are choosing which failure your downstream systems can survive, and there is no third option.

Raising WATERMARK to 30 minutes would rescue both events in closed mode, and delay every report by half an hour. That is the trade, stated in one sentence.

The professional middle ground is to close the window, emit the stable number, and route late events to a side output for reconciliation. You get a stable answer now and a correction path later. It is more code, and it is what production systems do.

Where this bites a model

Suppose your label is "did this customer make another payment within five minutes". Compute that label at 10:10 for the 10:00 window and you get one answer; compute it at 10:40 and you get another.

Training on labels computed after the fact, with all stragglers included, teaches the model a version of the world that the live system will never observe. The model will underperform in production by roughly the amount of information the late events carried, and nothing in your offline evaluation will show it.

The fix is to compute training labels with the same watermark the live system uses. Simulate the delay rather than wishing it away.

Common mistakes

Using arrival time because it is convenient. Demonstrated above. Ask your source for an event timestamp; if it will not give you one, record the gap and quantify what you are losing.

Assuming events arrive in order. They do not. Retries, partitions and mobile networks all reorder. Any code assuming monotonic timestamps has a bug awaiting the right day.

Setting a watermark from a guess. Measure it. Plot the distribution of arrival_time - event_time over a real week and pick a percentile. p99 is a common choice, and you should know what fraction it drops.

No late-arrival policy. Decide explicitly: drop them, correct the window and re-emit, or route them to a side output for later reconciliation. All three are valid. Having no policy means you picked "drop silently".

Testing only on ordered data. Your test fixture should include an out-of-order event and one arriving after the watermark. If it does not, your handling code has never run.

Try it yourself

Change WATERMARK to 30 minutes and re-run. Closed mode now reports the 10:00 window as 4, 1970.0, matching the truth — and reports nothing at all until 10:35. Correctness bought with latency, priced exactly.

Then add a payment with event time 10:03 and arrival time 11:15. At a five-minute watermark it is dropped for good; at thirty minutes it is still dropped. Some event is always too late, whatever you choose. Measuring how much money that is, per week, is part of the job.

What to learn next

Researcher — Mathematics and papers.

The Dataflow model

Akidau et al. (2015), The Dataflow Model, VLDB, is the reference formalisation and reframes the batch/stream distinction as unnecessary. Four questions define any computation:

  1. What results are computed — the aggregation.
  2. Where in event time — the windowing function.
  3. When in processing time results are emitted — triggers.
  4. How refinements relate — accumulation mode: discarding, accumulating, or accumulating-and-retracting.

Batch is the special case where the trigger fires once, at the end. This decomposition is why the developer example can express both event-time and arrival-time windowing with the same aggregation.

Windowing functions

A windowing function maps each element to a set of windows:

$$ \text{AssignWindows} : (k, v, t) \;\mapsto\; {\,w_1, \dots, w_m\,} $$

  • Tumbling — fixed, non-overlapping, of width $T$. Each element lands in one window; storage is $O(1)$ per key.
  • Sliding — width $T$, period $P < T$. Each element lands in $\lceil T/P \rceil$ windows; storage grows proportionally.
  • Session — dynamic, defined by a gap $g$. Windows merge when a new element bridges them, so session windows are data-dependent and require merge support in the runner.

Sliding windows are the usual cause of unexpected state growth. A 24-hour window sliding every minute puts every element into 1,440 windows.

Watermarks

A watermark is a monotonically non-decreasing function $W(t_p)$ from processing time to event time, asserting that no element with event time below $W(t_p)$ will be observed after processing time $t_p$.

Two regimes:

  • A perfect watermark requires complete knowledge of the source and is available only for bounded or strictly ordered inputs.
  • A heuristic watermark estimates the bound from observed lateness, source partition progress, or a fixed allowance. It can be wrong in both directions, producing late data when too aggressive and latency when too conservative.

The fundamental trade-off is unavoidable: completeness, latency and cost — pick two. Triggers decouple the three by allowing early speculative emissions before the watermark, an on-time emission at the watermark, and late refinements after it.

Kafka Streams takes a different route with a grace period per window and no global watermark, since partition-local ordering plus a bounded grace covers most cases at lower coordination cost.

Accumulation and correctness

When a window emits more than once, downstream consumers must know how to combine emissions:

  • Discarding — each emission covers only elements since the last. Downstream must sum. Suits an idempotent additive sink.
  • Accumulating — each emission is the current total. Downstream must overwrite. Requires an upsert-capable sink, which is the idempotency property from ETL for machine learning.
  • Accumulating and retracting — emits a retraction of the previous value alongside the new one. Necessary when the downstream aggregation is not invertible, for example when feeding a top-k that must remove the stale entry.

Choosing accumulating mode with an append-only sink is a common and quiet corruption: every refinement adds a row and the totals inflate.

Exactly-once, restated

As in the batch case, exactly-once delivery is unachievable; exactly-once effect is achievable through idempotent writes or transactional sinks with a two-phase commit. Flink implements the latter via distributed snapshots (Carbone et al., 2015, Lightweight Asynchronous Snapshots for Distributed Dataflows), an adaptation of the Chandy-Lamport algorithm: barriers flow with the data, and operator state is checkpointed consistently without halting the stream.

Recovery restores the last complete checkpoint and replays from the recorded source offsets. Correctness therefore requires a replayable source — a log with retention, not a queue that discards on acknowledgement. This is the architectural reason Kafka's log-with-offsets model displaced traditional message queues for this workload.

Stream-table duality

A stream of changes and a table of current state are two views of the same information (Kreps, 2013). A table is the fold of a stream; a stream is the changelog of a table:

$$ \text{table}{t} \;=\; \text{fold}\big(\oplus,\; \text{table}{0},\; \text{stream}_{[0,t]}\big) $$

This duality is what makes change data capture — reading a database's write-ahead log as a stream — and materialised views over streams the same mechanism. It is also the cleanest justification for event sourcing as a storage design: the log is the truth, and every table is a derived, recomputable projection.

Consequences for machine learning

  1. Label maturity. A label defined over a window is not final until the watermark passes. Training on fully matured labels while serving on immature ones is a form of leakage that offline evaluation cannot detect.
  2. Feature freshness against correctness. The freshest feature is computed before late data arrives and is therefore wrong; the correct feature is stale. Quantify the trade for your specific lateness distribution rather than assuming one end.
  3. Backfill asymmetry. Reprocessing history from a log has no lateness, since everything is already there. A model trained on backfilled features sees a cleaner world than production, and will overperform offline. Simulate the watermark during backfill.
  4. Concept drift. Streams shift under you. Continuous evaluation on a delayed-label holdout is the only reliable detector. See monitoring and model drift.

Reading

  • Akidau et al., The Dataflow Model, VLDB 2015 — research.google/pubs/pub43864/
  • Akidau, Chernyak and Lax, Streaming Systems, O'Reilly 2018 — the book-length version, and the best treatment available.
  • Carbone et al., Lightweight Asynchronous Snapshots for Distributed Dataflows, 2015 — arxiv.org/abs/1506.08603
  • Kreps, Narkhede and Rao, Kafka: a Distributed Messaging System for Log Processing, NetDB 2011.
  • Kleppmann, Designing Data-Intensive Applications, O'Reilly 2017 — chapter 11.

What to learn next