Windowing¶
turbine.windowing ships six window primitives. They all consume the same aggregate functions (turbine.aggregates), offer the same pair of close callbacks (on_close_each= per pane, on_close= per batch), and all support opt-in early-fire triggers — they differ only in how state is stored and how their boundaries are defined.
This page covers what is common to every window: the aggregate you compute, the close callbacks, early-fire, and the time model. The two companion pages cover the rest:
- Window types — Tumbling, Sliding and Session, each in a persistent (durable, RocksDB-backed) and an in-memory variant, with a code example for each.
- Windowing — advanced — failover, the in-memory trade-offs in depth, referencing one window's result from another (
Ref/join_with), dynamic (deferred-construction) windows + the recovery hook, and sessions at high cardinality.
| Primitive | Boundaries | Storage | Crash recovery |
|---|---|---|---|
PersistentTumbling |
Fixed size_ms |
RocksDB + checkpoints | ✅ survives restart, full failover |
PersistentSliding |
Fixed size_ms / slide_ms |
RocksDB + checkpoints | ✅ survives restart, full failover |
PersistentSession |
Data-driven (gap_ms) |
RocksDB + checkpoints | ✅ survives restart, full failover |
InMemoryTumbling |
Fixed size_ms |
Per-partition memory | ❌ in-flight state lost on restart |
InMemorySliding |
Fixed size_ms / slide_ms |
Per-partition memory | ❌ in-flight state lost on restart |
InMemorySession |
Data-driven (gap_ms) |
Per-partition memory | ❌ in-flight state lost on restart |
Rules of thumb:
- Choose the shape by your boundaries: fixed buckets → Tumbling; overlapping/rolling → Sliding; activity-driven with no fixed boundary (user sessions, device bursts) → Session.
- Choose storage by the throughput/durability trade-off:
- Persistent (
PersistentTumbling/PersistentSliding/PersistentSession) is the safe default — survives crashes, full failover, sub-second close timing, browsable in the console. - In-memory (
InMemoryTumbling/InMemorySliding/InMemorySession) is the faster option: it skips the per-event durable write, so it's worth considering whenever a large number of windows are open at once (high key × window cardinality) or the event rate is high — exactly the regimes where the state-store write path becomes the bottleneck. The cost is durability: in-flight window state is lost on a crash (the source replay recreates it, so make close side-effects idempotent). The full trade-off is in In-memory windows.
- Persistent (
When @app.subscribe(..., partition_key="...", parallelism=N) is set, every window inside that subscription must be keyed (at least) by the same column. Otherwise rows for a single window key would be split across shards and each shard would aggregate over a partial view. The SDK enforces this at construction — you get a clear error pointing at the offending key=. The partition_key itself is auto-injected: you never write it again on the window. See Partitioning for the full picture.
Aggregate functions¶
Pass any aggregate function from turbine.aggregates to a window's value= parameter. The common shapes — Count(), Sum("col"), Mean("col"), Max/Min, percentile sketches — all share the same merge-across-batches contract; the dedicated Aggregate Functions page lists every one.
Column arguments accept either a string or a function expression: agg.Sum(fn.col("bytes_in") + fn.col("bytes_out")) is equivalent to materialising the derived column upstream and passing the name. See Functions for the full catalog.
Multi-Metric Aggregation¶
To compute something like SUM(errors) / COUNT(*) per window — the stream-processing equivalent of a SQL SUM(a)/COUNT(*) — use agg.Multi(...). Each keyword argument becomes an attribute on the value delivered to on_close_each:
value=agg.Multi(
err_sum=agg.Sum("is_error"),
lat_sum=agg.Sum("latency_ms"),
n=agg.Count(),
),
on_close_each=lambda state, window, v: pa.RecordBatch.from_pylist([{
"region": window.group["region"],
"error_rate": v.err_sum / v.n,
"mean_latency_ms": v.lat_sum / v.n,
"window_start_ms": window.start_ms,
}]),
Bundling sub-aggregates inside Multi is the cheap way to compute several metrics in lockstep: cost stays roughly linear in the batch size regardless of how many you add. A batch that's missing any column referenced by any sub-aggregate skips the whole update for that window, not a single metric.
Referencing another window's result¶
A window's aggregate expression can also pull in another window's last-closed result for the same key — for example, falling back to a rolling 15-minute baseline when the current minute is sparse. The referenced window names its output (with .alias("name") or Multi), and the reader declares join_with=[that_window] and reads win.Ref("name").field inside its value=. References are persistent-windows-only and resolve to the last closed value (null until the first close, so pair them with fn.coalesce). The full mechanics — key/partition rules, ordering, recovery and cost — are in Windowing — advanced.
Window close callback¶
You pick one of two callback shapes when constructing a window — they're mutually exclusive, and the SDK errors at construction if you set both or neither. The choice is the same for every window kind.
on_close_each — one call per pane (the common shape)¶
Called once per (group, window) pane when the window closes:
state— the sameTurbineStateyourprocessmethod receives. Use it to read/write auxiliary state (e.g. cooldown machines).window— aWindow(start_ms, end_ms, group)dataclass.groupis a dict mapping eachkeyfield to its value. For session windowsend_msis data-driven (last event + gap), notstart + size.value— the aggregate's finalized scalar (aCountint, aMeanfloat, aMultinamespace, …).
Return a pyarrow.RecordBatch to publish, or None to skip emission for this pane.
Use this when each closed window produces an independent output row and the cost is dominated by the per-pane logic (typical for alerts, dashboards, per-key emissions).
on_close — one call per batch of panes (the vectorised shape)¶
Called once per tick with all panes that closed at the same time packed into a single Arrow Table. The table carries one row per pane, with the key columns, window_start_ms, window_end_ms, and the finalized aggregate columns.
Use this when:
- Many panes close together and the per-pane computation is amenable to vectorised PyArrow compute (
pc.divide,pc.if_else, …). Computing a score over 10 000 panes in one Arrow pass is dramatically faster than 10 000 Python calls. - The downstream emission is a single
RecordBatchregardless of how many panes contributed.
Return a pyarrow.RecordBatch (the merged emissions for this tick) or None to skip the tick entirely.
examples/dev_app.py (score_alert) is the canonical example of the batch shape.
Early-fire triggers (provisional emissions)¶
By default a window emits once, at close. For a long-lived window — a 5-minute revenue rollup, an hour-long session still receiving events — you often want to see the current aggregate before the window closes (a live dashboard, a "latest" gauge). Set early_fire_interval_ms= together with a separate early callback and the window will also emit its current-so-far aggregate every interval while it is still open. The close emission is unchanged — it still fires once, at close, and remains authoritative.
It is opt-in and available on all six window kinds. When unset, behaviour is exactly the single emission at close.
Two callbacks, not a flag¶
Early emission gets its own callback, distinct from the close one:
@app.subscribe(kafka.topic("orders"), output=kafka.topic("revenue_5m"))
class revenue:
def __init__(self) -> None:
self._win = win.PersistentTumbling(
self,
name="revenue_5m",
size_ms=5 * 60_000,
key=["region"],
value=agg.Sum("amount"),
time_column="event_ts_ms",
allowed_lateness_ms=5_000,
on_close_each=self._commit, # 1×, authoritative, at close
early_fire_interval_ms=10_000, # every 10s while the window is open
on_early_fire_each=self._live_gauge, # n×, provisional, current-so-far
)
def _commit(self, state, window, value): # authoritative total → ledger
return pa.RecordBatch.from_pylist([{
**window.group, "revenue": float(value), "final": True,
"window_start_ms": window.start_ms, "window_end_ms": window.end_ms,
}])
def _live_gauge(self, state, window, value): # running total → live dashboard
return pa.RecordBatch.from_pylist([{
**window.group, "revenue": float(value), "final": False,
"window_start_ms": window.start_ms, "window_end_ms": window.end_ms,
}])
def process(self, batch: RecordBatch, state: TurbineState) -> RecordBatch | None:
self._win.update(batch, state)
return None
The early fire typically pushes a provisional value somewhere cheap (a gauge, a "latest" topic), while the close triggers the authoritative effect (commit, raise an alert). Two callbacks put that intent in the signature instead of an if is_final: branch. If you genuinely want identical handling, pass the same function to both.
Form mirroring + the rules¶
- The early callback mirrors the close form: pair
on_close=(batch:state, panes) withon_early_fire=, andon_close_each=(per pane:state, window, value) withon_early_fire_each=. Same data shape on both doors. Mixing forms is rejected at construction. early_fire_interval_msand the early callback are co-required (one without the other is an error), and the interval must be> 0. All of this is validated when the window is built, never at fire time.- The early callback returns a
RecordBatchto publish, orNoneto skip — exactly like the close callback. Returned batches are published tooutputlike any other output.
What's guaranteed (and what isn't)¶
- Provisional vs authoritative. Anything from the early callback is best-effort by construction; the close emission is the authoritative one. Tag your output (a
finalcolumn as above) or route the two callbacks to different topics. - Ordering. For a given window, every early fire strictly precedes the single close, and nothing is emitted after the close. An early fire never duplicates or races the close.
- Sessions can supersede an early value. Session windows merge: a bridging event can fuse two open sessions into one. An early fire emitted for a session that later merges describes an aggregate that, in hindsight, no longer exists on its own — the
on_closetotal supersedes it. There are no retractions; if a downstream cannot tolerate provisional double-counting across a merge, either don't wire the early callback for sessions or upsert on a stable key. Tumbling and sliding windows have fixed bounds and never merge, so their early fires are monotone refinements of the same window. - A closed window never re-opens. Once a window has closed, an event belonging to it is dropped — or routed to
late_data=, if you asked for that; either way it stays out of the window (see lateness). A given(group, window)produces exactly one close, ever. A downstream keyed on the group alone can still seeclosethen an early fire for the same key, but that is a different window (a later one), not the closed one coming back.
Timing¶
The interval rides the same clock the window closes on: event-time (the watermark) when time_column= is set and data is flowing, processing-time (wall-clock) otherwise. On a quiet partition the two families differ: the persistent windows keep early-firing on a wall-clock cadence (sub-second) even with no traffic, while the in-memory windows only advance on the periodic idle tick (so early fires are quantised to ~10 s there). If you need sub-10 s early fires on a partition that may go silent, use a persistent window.
Under exactly-once, early fires are outputs like any other: they ride the same per-batch transaction and are never observable twice. They do multiply output volume, so size the interval to your downstream's tolerance rather than making it as small as possible.
Time model: processing time, event time, lateness¶
Every window runs in one of two time modes, selected the same way regardless of kind.
Processing time (default). With no time_column, every row in an incoming batch is treated as having happened now (now_ms()). Use it when your data carries no timestamp or when clock skew doesn't matter. For session windows this means the gap is measured on wall-clock inactivity.
Event time. Each row is placed by its own timestamp — a single batch can fan out across windows, and out-of-order / late events are handled correctly (for sessions, a late event can even bridge two open sessions). The close timer then fires in wall-clock time at window_end + allowed_lateness_ms. There are two ways to select the timestamp.
The recommended way is to declare it once on the input topic with event_time="<column path>". Turbine reads that field from each record, converts it to an internal timestamp, and every window under that subscription automatically uses it — you don't repeat time_column= on each window:
@app.subscribe(kafka.topic("events", event_time="event_ts", event_time_unit="ms"), output=kafka.topic("counts"))
class count_per_region:
def __init__(self) -> None:
self._win = win.PersistentTumbling(
self, name="events_per_region", size_ms=60_000, key=["region"],
value=agg.Count(), on_close_each=self._emit,
allowed_lateness_ms=5_000, # event-time inferred from the subscription
)
...
event_time accepts a dotted path to reach a nested field (event_time="meta.event_ts"), and the source field may be any of:
| Source field type | Handling |
|---|---|
| Integer (epoch) | Scaled to milliseconds per event_time_unit ("ms" default, or "s", "us", "ns"). |
| Float (epoch seconds) | Same, scaled per event_time_unit. |
| Timestamp | Converted directly; event_time_unit is ignored. |
| RFC 3339 / ISO-8601 string | Parsed (e.g. "2021-01-01T00:00:00Z"); event_time_unit is ignored. |
A record whose field is missing or whose string can't be parsed is dropped from windowing (counted, never fatal) — one malformed record won't stall the stream.
The explicit way is to pass time_column="<column>" on an individual window. This bypasses the subscription default and requires a column that is already an Int64 epoch-millisecond value (no conversion is applied):
win.PersistentTumbling(
self, name="events_per_region", size_ms=60_000, key=["region"],
value=agg.Count(), on_close_each=self._emit,
time_column="event_ts_ms", # an Int64 epoch-ms column you produced
allowed_lateness_ms=5_000, # wait 5 s after window end before closing
)
Prefer the input topic's event_time= for new code; use time_column= when you've already materialised an epoch-ms column or need different windows in one subscription to key off different columns. (Setting both — an input-topic event_time= and a window time_column= that points elsewhere — warns, and the window-level value wins.)
allowed_lateness_ms is how long after a window's end Turbine keeps accepting events that belong to it. It delays the close: the window stays open until window_end + allowed_lateness_ms in event time, and an out-of-order event arriving before that is an ordinary in-window event — it is merged normally and the window still emits exactly one authoritative result.
Once that deadline passes, the window closes and its state is released. An event that belongs to it and arrives after that point is dropped by default, and counted (see below). Turbine does not re-open a closed window: doing so would publish a second result for the same window carrying only the straggler, contradicting the one already sent — and it would put no bound at all on how far back a single old record could reach. So allowed_lateness_ms is the whole of your tolerance budget: size it for the worst upstream stall you want to absorb.
A window closed because its partition went quiet rejects stragglers just like one closed by data — Turbine tracks how far closing has actually got, so a straggler cannot re-open a window that already published. Rows rejected only for that reason are counted separately on turbine_late_rows_idle_closed_total; a non-zero value means part of what you are losing depends on your traffic pattern rather than on the data's own timestamps.
Late drops are never silent — they increment turbine_late_rows_dropped_total, labelled by worker and window name. A non-zero, growing value means your allowed_lateness_ms is smaller than your real out-of-orderness, and you are losing data you asked to keep. Alert on it.
Keeping late rows instead of dropping them¶
Dropping is the default, not the only option. Set on_late="route" on the subscription and name a destination with late_data= — every row past the deadline is published there instead of being discarded:
@app.subscribe(
kafka.topic("events", event_time="ts"),
output=kafka.topic("counts"),
on_late="route",
late_data=kafka.topic("events-late", message_key="tenant_id"),
)
What you get and what you give up:
- Nothing is lost, and no result is contradicted. The straggler leaves on its own topic; the window it belonged to keeps the single authoritative result it already published. This is the reason routing exists rather than re-opening: both avoid losing the row, but only one of them avoids emitting two conflicting answers for the same window.
- The routed record is the input row as Turbine decoded it, encoded like any other output, plus three columns describing the rejection:
_late_window(the window's declared name),_late_window_end_ms(the end of the last window that could still have owned it) and_late_watermark_ms(how far the data's clock had already advanced). With several windows on one subscription,_late_windowis what tells them apart. Turbine's internal columns are stripped. late_datamust be its own topic. Pointing it at the subscription's input loops forever; pointing it atoutput=merges rejected rows into your result stream. Both are refused at startup.- Delivery follows the subscription's guarantee. The routed rows are produced with the batch they were rejected from — under
exactly_oncethey commit inside the same transaction as that batch's offsets, so a crash cannot leave a row routed-but-not-committed or committed-but-not-routed. - Routed rows are counted separately, on
turbine_late_rows_routed_total{worker, window}—turbine_late_rows_dropped_totalstays for the rows a"drop"subscription discards, so switching policy reads as one series going quiet and another waking up. - Session windows use a different, more permissive rule — see below.
on_lateapplies to them too.
late_data is one of the extra destinations a subscription can declare; the others (side_outputs= for imperative ctx.emit, branch= for a declared split) work the same way and are covered in Side outputs.
Lateness for session windows¶
A session has no fixed end: activity keeps extending it. So "is this row late?" is decided on the session the row would land in after merging, not on the row's own timestamp plus the gap.
The practical consequence is worth stating plainly: a session window accepts rows far older than allowed_lateness_ms, as long as they fall inside a session that is still open. If a session has been kept alive for an hour by continuous activity, an event from forty minutes ago still belongs to it and is merged normally. Only once that session has closed do rows belonging to it become late.
This is deliberate, and it is what gap_ms means: a session is bounded by inactivity, not by the clock. Judging a straggler on timestamp + gap_ms would drop exactly the out-of-order rows that bridging exists to absorb.
What you get is the same guarantee as the fixed windows: a (group, session) produces exactly one result, and a row arriving after that result is dropped or routed — never merged into a second, overlapping session for the same key.
allowed_lateness_ms also doubles as the failover-completeness budget for persistent windows (see Failover-resilient windows).
How an event-time window decides to close¶
In event-time mode a window closes when the data's own clock reaches its end — not when your wall clock does. Turbine tracks, per partition, the largest event timestamp it has seen (its watermark); a window [start, end) closes once that watermark passes end + allowed_lateness_ms. Two consequences worth knowing:
- Out-of-order and historical data just work. Replaying a day-old backfill closes its windows as the watermark sweeps through the old timestamps — at full speed, not pinned to real time. A burst of future-dated events closes every earlier window up to their timestamp in one step.
- A window still closes when the stream goes quiet. If events stop arriving, the watermark can't advance from data — so Turbine lets it drift forward at real time from the last event seen. In practice a window then closes after roughly the part of it not yet covered by data, plus
allowed_lateness_ms, of real silence. Example: a 60-second window whose last event landed 20 seconds in waits about 40 more seconds of silence before closing. A window that's already fully covered by data closes immediately. This guarantees no window stays open forever on an idle partition, while never closing one early just because the wall clock moved on.
Driving windows off the Kafka timestamp¶
If your producer doesn't embed a timestamp in the payload, use the Kafka broker timestamp. Opt in on the input topic with with_kafka_timestamp=True; Turbine injects an Int64 _kafka_ts_ms column on every batch, which you then pass as time_column="_kafka_ts_ms".
@app.subscribe(kafka.topic("events", with_kafka_timestamp=True), output=kafka.topic("counts"))
class count_per_region:
def __init__(self) -> None:
self._win = win.PersistentTumbling(
self, name="events_per_region", size_ms=60_000, key=["region"],
value=agg.Count(), on_close_each=self._emit,
time_column="_kafka_ts_ms", allowed_lateness_ms=5_000,
)
...
Two broker-side configs shape what _kafka_ts_ms means:
| Topic config | Source of timestamp |
|---|---|
CreateTime (default) |
Producer clock at send. Reflects event time if the producer is close to the source, but vulnerable to producer clock skew. |
LogAppendTime |
Broker clock at append. Monotone per-partition, but closer to processing time than true event time. |
Use LogAppendTime when you don't trust producer clocks; use CreateTime (or a payload column) when the producer is authoritative.
Naming a window: name= vs label=¶
Every window takes a name=. It is an identity, not a title: it prefixes
every state and timer key the window owns, so changing it abandons the
partial aggregates, panes and timers already on disk — the window starts
from nothing, and whatever was mid-flight is never emitted.
That is fine when you write the name yourself. It is a problem when windows
come from configuration — one per rule, per tenant, per customer — because
the only stable thing to name them after is an id, and an id is exactly what
nobody can read. A console listing alert-rule:17761784b6 a hundred times
tells you nothing about what any of them watches.
label= separates the two. It is display metadata: it never appears in a
key, it can be anything, and it can change at any time.
self._win = win.PersistentSliding(
self,
name=f"alert-rule:{rule.id}", # identity — stable for the life of the state
label=rule.title, # what a human reads: "Net latency per runner"
size_ms=600_000, slide_ms=60_000,
key=["tenant_id", "runner_id"],
value=agg.Mean("latency_ms"),
on_close_each=self._emit,
)
The console shows the label as the window's title and keeps the name
visible underneath — the name is what your logs, alert payloads and API
calls carry, so it must stay findable. GET /windows returns both.
If your configuration is reloaded while the app runs, call
window.set_label(rule.title) when a title changes. It republishes the
metadata — no key is touched, nothing is lost, and the console catches up
on its next poll instead of at the next restart.
Reading a window while it is open¶
A window's value is normally something you receive: the close callback fires and the result goes downstream. But while a window is open — and especially when you are asking why one fired and its neighbour didn't — you often want the numbers behind it.
GET /windows/{name}/series answers exactly that for one group: the value
of each pane a sliding window is aggregating (or of each window a group has
open), plus the aggregate they currently roll up to. Values come back
finalized — a mean is a number — so nothing on the reading side needs to
know how the aggregate function stores itself. The
REST API reference has the full shape, and the console
plots the same data as a curve.
The series covers what is still open: a sliding window's held panes (the
last size_ms), or the windows not yet closed. Closed windows leave state
as they emit; their history lives in the output stream you sent them to.