skip to content

In Flink's DataStream API, how do tumbling, sliding and session windows differ?

level: juniorimportance: must knowfreq 72%

answer

  1. three shapes, plus one that never ends
  2. count how many windows one record joins
  3. one type has no fixed end time
  4. inactivity closes the third kind
  5. overlap comes from slide smaller than size

basics

~20 s

Tumbling windows are fixed-size and never overlap, so a record lands in exactly one. Sliding windows are fixed-size but start every slide interval, so a record lands in several. Session windows have no fixed size and close after a gap of inactivity.

solid answer

~50 s

All three are `WindowAssigner` implementations you pass to `window(...)` after a `keyBy(...)`. **Tumbling** (`TumblingEventTimeWindows.of(...)`) cuts the stream into fixed-size, non-overlapping buckets — every record belongs to exactly one, which makes it the right choice for "revenue per hour". **Sliding** (`SlidingEventTimeWindows.of(size, slide)`) also has a fixed size but a new window starts every slide interval; when the slide is smaller than the size the windows overlap and each record is copied into `size / slide` windows — the right choice for "revenue over the last hour, refreshed every minute". **Session** (`EventTimeSessionWindows.withGap(...)`) has no fixed size at all: a window closes once no record for that key arrives for the configured gap, so it models user activity bursts. A fourth assigner, `GlobalWindows.create()`, puts every record for a key in one never-ending window and never fires unless you attach a custom trigger; its sibling `GlobalWindows.createWithEndOfStreamTrigger()` fires once, when a bounded input ends.

code

java · 16 lines
java
DataStream<Event> input = ...;

// exactly one window per record
input.keyBy(Event::getUserId)
     .window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
     .aggregate(new CountAgg());

// 12 open windows -> each record copied 12 times
input.keyBy(Event::getUserId)
     .window(SlidingEventTimeWindows.of(Duration.ofHours(1), Duration.ofMinutes(5)))
     .aggregate(new CountAgg());

// boundaries defined by 30 minutes of silence
input.keyBy(Event::getUserId)
     .window(EventTimeSessionWindows.withGap(Duration.ofMinutes(30)))
     .aggregate(new CountAgg());

go deeper

for a junior

Be ready to name the four built-in assigners and say, for each, whether a single record can land in more than one window. Interviewers usually follow up by asking which one fits a metric they describe out loud.

for a middle

Explain the mechanics: epoch alignment and the offset argument, the size-over-slide copy factor for sliding windows, and the fact that a session assigner opens a window per record and merges them afterwards.

for a senior

Show that you check the state cost of an assigner before choosing it, and that you know when to abandon built-in windows entirely for a KeyedProcessFunction with your own timers.

for a principal

Own the question of whether the metric should be windowed at all — whether the business definition is a calendar period, a trailing interval or a behavioural session, and what each choice commits the platform to in state size and reprocessing cost.

## Why windows exist A stream is unbounded, so an aggregation like `SUM` has no natural point at which to emit a result. A **window** slices the infinite stream into finite buckets over which a computation can be run and a result emitted. In Flink's DataStream API the shape of those buckets is decided by a **window assigner**, the object you pass to `window(...)` on a keyed stream (or `windowAll(...)` on a non-keyed one). A `WindowAssigner` has exactly one job: given a record, decide which window or windows it belongs to. A keyed window operator runs one independent set of windows *per key*, and different keys can be processed by different parallel subtasks. A non-keyed `windowAll(...)` has no key to partition on, so all the windowing runs in a single task at parallelism 1 — fine for a small final roll-up, a bottleneck for anything else. In Flink 2.3 every time-based assigner takes `java.time.Duration` arguments; the old `Time` helper class they used to take was removed in 2.0. ## Tumbling windows `TumblingEventTimeWindows.of(Duration.ofMinutes(5))` produces fixed-size, non-overlapping, gap-free buckets: `[12:00, 12:05)`, `[12:05, 12:10)`, and so on. The start timestamp is inclusive and the end timestamp exclusive, so a record at exactly `12:05` belongs to the second window. Every record belongs to exactly **one** window, which makes tumbling windows the cheapest option in state: one entry per key per open window. By default the windows are aligned to the Unix epoch, so hourly windows begin on the hour in UTC. The assigner takes an optional second argument, an **offset**, precisely so you can realign them: `TumblingEventTimeWindows.of(Duration.ofDays(1), Duration.ofHours(-8))` gives you a daily window that begins at local midnight in a UTC+8 timezone rather than at UTC midnight. Tumbling is what you want for periodic reports: orders per hour, errors per minute, revenue per day. ## Sliding windows `SlidingEventTimeWindows.of(Duration.ofHours(1), Duration.ofMinutes(5))` also produces fixed-size windows, but a new one starts every *slide* interval rather than every *size* interval. With size 1 hour and slide 5 minutes, twelve windows are open at any moment and a single record is assigned to all twelve. That is the defining property: **sliding windows overlap, and overlap means the record is copied into every window it belongs to**, so the state cost is roughly `size / slide` times that of a tumbling window of the same size. A sliding window of size one day and slide one second is a well-known way to destroy a cluster. Sliding windows answer "the last N minutes, refreshed every M minutes" — trailing averages, moving thresholds, rolling alert conditions. Like tumbling, they take an optional offset for alignment. ## Session windows `EventTimeSessionWindows.withGap(Duration.ofMinutes(30))` has neither a fixed size nor a fixed alignment. A session window for a key extends as long as records keep arriving; it ends when a **gap of inactivity** longer than the configured gap occurs, and the next record opens a fresh session. Two sessions never overlap, and different keys will have completely different session boundaries at completely different times. Internally, the operator creates a *new* window for every arriving record and then **merges** windows that are closer to each other than the gap. That merging behaviour is why a session window assigner requires a merging trigger and a merging window function (`ReduceFunction`, `AggregateFunction` or `ProcessWindowFunction`). A variant, `EventTimeSessionWindows.withDynamicGap(...)`, lets you compute the gap per record with a `SessionWindowTimeGapExtractor` — useful when a premium user should get a longer inactivity timeout than a free one. Sessions are the right model for user behaviour: a browsing session, a support conversation, a device that reports in bursts. ## Global windows `GlobalWindows.create()` assigns every record for a key to one single window whose end timestamp is effectively infinite. Created with `GlobalWindows.create()`, its default trigger is `NeverTrigger`, so on its own it emits nothing at all; it becomes useful with a custom trigger, most commonly `CountTrigger`, to build count-based rather than time-based windows. Flink 2.3 also offers `GlobalWindows.createWithEndOfStreamTrigger()`, whose trigger fires exactly once, when the input ends: useful for a whole-input aggregate over a bounded source, and it never fires on an unbounded one. Because that window never ends, no record is ever considered late in it. ## Choosing between them Start from the metric the business is asking for. "Per calendar period" is tumbling. "Trailing N, refreshed every M" is sliding — and check `size / slide` before you ship it. "Per burst of activity, boundaries defined by the data" is session. "Every N records regardless of time" is a global window with a count trigger. If the metric maps onto none of these cleanly, a `KeyedProcessFunction` with your own state and timers is often simpler and far cheaper than bending a built-in assigner into shape.

  • What does the optional second Duration argument on TumblingEventTimeWindows.of do?
    It is an offset that shifts window alignment. Without it, windows align to the Unix epoch, so hourly windows start on the hour in UTC. Passing `Duration.ofHours(-8)` shifts daily windows to begin at local midnight in a UTC+8 timezone. It is the standard way to make calendar-day aggregations match a business's timezone rather than UTC.
  • Why does GlobalWindows.create() emit nothing unless you configure it further?
    A global window has no natural end — its end timestamp is effectively infinite — so there is nothing for a time-based trigger to fire on. `GlobalWindows.create()` installs `NeverTrigger`, which never returns FIRE. You must supply a custom trigger via `trigger(...)`, typically `CountTrigger`, to get count-based firing. The one built-in exception is `createWithEndOfStreamTrigger()`, which fires once when a bounded input ends.
  • What is the difference between window(...) and windowAll(...)?
    `window(...)` is called on a `KeyedStream` and gives one independent set of windows per key, so the operator can run in parallel across subtasks. `windowAll(...)` is called on an unkeyed stream and there is nothing to partition on, so the whole windowing computation runs in a single task at parallelism 1. Use it only for small final roll-ups.

Tumbling windows are like consecutive pages of a calendar; sliding windows are like a magnifying glass dragged across it, showing overlapping stretches; session windows are like paragraphs, ending wherever the writer stopped typing long enough.

saying these in an interview costs you the question

  • Says sliding windows never overlap or are just tumbling with an offset
  • Claims a session window has a fixed configurable duration
  • Thinks a record in a sliding window is stored only once
  • Believes windowAll still runs in parallel across subtasks
  • Expects GlobalWindows.create() to emit results without a custom trigger

context