Windows are how Flink turns an unbounded stream into finite groups that can be aggregated, and they are where most event-time bugs come from: results that never appear, results that appear twice, state that grows until a TaskManager dies, and numbers that quietly leave out late data.

The general idea of windowing is covered in Stream windowing, and Flink's runtime as a whole in Apache Flink Architecture in Depth. This article is about Flink's window operator specifically: the path from element to result, the knobs that change it, and how to run it. Code is Java against the Flink 2.x DataStream API, where durations are java.time.Duration (the old Time class was removed in 2.0 under FLIP-335), and the SQL section follows the window TVF documentation for the current stable release.

Advertisement

The window operator: five parts

A windowed DataStream program has one shape: keyBy, window with an assigner, optional trigger, evictor and lateness, then a function.

  • Key. keyBy partitions the stream so that every window is per key. Each subtask owns a range of key groups and keeps window state only for those keys. Non-keyed windowAll runs with parallelism 1.
  • WindowAssigner. Maps each element to zero or more windows. A tumbling assigner returns one window, a sliding assigner returns several, and a session assigner returns a provisional window that may later be merged.
  • Trigger. Decides when a window is evaluated. It is called on every element and on every timer, and returns CONTINUE, FIRE, PURGE or FIRE_AND_PURGE. Event-time assigners default to EventTimeTrigger, which fires once the watermark passes the end of the window.
  • Evictor. Optional. Removes elements from the window buffer before or after the function runs.
  • Window function. Computes the result: ReduceFunction or AggregateFunction incrementally, or ProcessWindowFunction over all buffered elements, or a combination.
Inside one keyed window operator (per parallel subtask)keyBy(userId)hash to key groupWindowAssignerelement to window(s)Window stateper key and windowaddTriggeronElement / timersTimer serviceend-1 and cleanupregister timersEvictor (optional)before/after functionFIREWindow functionreduce / aggregate / processDownstreamone result per firingWatermark from upstreammin over input channelsadvancesLate beyond allowed latenessdropped, or side outputwindow already cleanedState, timers and the trigger are all scoped to one (key, window) pair and live in the state backend.
Each element is assigned to windows, added to per-key, per-window state, and offered to the trigger. Timers at the window end and at the cleanup time drive firing and state removal. Elements for windows that have already been cleaned up are dropped or sent to a side output.

A window appears when its first element arrives, as state plus timers keyed by (key, window), and disappears when the cleanup timer fires. A key with no data has no windows, so a quiet minute produces no empty result.

Assigners and what each one costs

AssignerWindows per elementState shapeTypical use
TumblingEventTimeWindows.of(size)1one accumulator per key per windowper-minute counts, hourly billing
SlidingEventTimeWindows.of(size, slide)size divided by slidethat many copies of every element's contributionmoving averages, "last hour, updated every minute"
EventTimeSessionWindows.withGap(gap)1, then mergedone per active session, merged as gaps closeuser sessions, device bursts
GlobalWindows.create()1 for all timegrows forever unless a trigger purgescount-based or custom-triggered windows

Tumbling and sliding windows are aligned to the epoch, so a one-day window runs midnight to midnight UTC. All of them accept an offset: a tumbling one-day window with an offset of minus eight hours, Duration.ofHours(-8), aligns days to UTC+8. Processing-time variants exist, but a replay of them produces different numbers; prefer event time.

The sliding row surprises people: a one-hour window sliding every minute puts each element into sixty windows, and with a ProcessWindowFunction that is sixty buffered copies. For long windows with short slides, pre-aggregate in small tumbling windows and combine downstream.

Advertisement

Event time and watermarks decide when windows fire

An event-time window fires when the operator's watermark reaches the window's maximum timestamp, which is the end minus one millisecond. The watermark is a claim from upstream: no more elements with a timestamp at or below this value are expected. A window operator with several input channels takes the minimum of their watermarks, so the slowest partition sets the pace for everyone.

WatermarkStrategy<Click> strategy = WatermarkStrategy
    .<Click>forBoundedOutOfOrderness(Duration.ofSeconds(5))   // wm = max seen ts - 5 s - 1 ms
    .withTimestampAssigner((click, recordTs) -> click.eventTimeMillis())
    .withIdleness(Duration.ofMinutes(1));                       // idle splits stop holding wm back

DataStream<Click> clicks = env.fromSource(kafkaSource, strategy, "clicks");

Two settings dominate correctness. The out-of-orderness bound is a trade between latency and completeness: every window waits that long after its end before firing, and anything later than that becomes late data. Idleness matters whenever a Kafka partition or a source split can go quiet. Without it, one silent partition pins the minimum watermark and no window anywhere in the job fires. That is the most common "my Flink job outputs nothing" ticket.

Incremental versus buffering window functions

Where the computation happens decides how much state each window holds. A ReduceFunction or AggregateFunction folds each element into an accumulator as it arrives, so the state per window is one accumulator. A ProcessWindowFunction receives an Iterable of every element when the window fires, so Flink must keep every element in list state until then. ProcessWindowFunction still earns its place because it sees the window bounds and watermark, so the standard pattern combines them: aggregate incrementally and let the process function decorate the single result.

class CountAgg implements AggregateFunction<Click, Long, Long> {
    public Long createAccumulator()          { return 0L; }
    public Long add(Click c, Long acc)       { return acc + 1; }
    public Long getResult(Long acc)          { return acc; }
    public Long merge(Long a, Long b)        { return a + b; }   // needed for session windows
}

class Stamp extends ProcessWindowFunction<Long, UserCount, String, TimeWindow> {
    public void process(String user, Context ctx, Iterable<Long> counts, Collector<UserCount> out) {
        long n = counts.iterator().next();                       // exactly one pre-aggregated value
        out.collect(new UserCount(user, ctx.window().getStart(), ctx.window().getEnd(), n));
    }
}

OutputTag<Click> lateTag = new OutputTag<Click>("late-clicks") {};

SingleOutputStreamOperator<UserCount> perMinute = clicks
    .keyBy(Click::userId)
    .window(TumblingEventTimeWindows.of(Duration.ofMinutes(1)))
    .allowedLateness(Duration.ofSeconds(30))
    .sideOutputLateData(lateTag)
    .aggregate(new CountAgg(), new Stamp());

DataStream<Click> late = perMinute.getSideOutput(lateTag);   // audit, or reprocess in batch

Allowed lateness and late firings

By default allowed lateness is zero: once the watermark passes the end of a window, the window fires, its state is cleared, and any later element for it is dropped. Setting allowedLateness keeps the window's state until the watermark passes the end plus the lateness. During that period a late element is added to the window and EventTimeTrigger fires again immediately. That late firing emits a new, complete result for the same window.

This has a direct consequence downstream: the same window can produce several results, and the later ones replace the earlier ones. The sink must treat window results as upserts keyed by (key, window start), for example a primary-key table, an upsert Kafka topic or an idempotent write. An append-only sink will double-count. Elements that arrive after the lateness has expired go to the side output if you configured one, and otherwise they are counted in the operator's numLateRecordsDropped metric and discarded. Alert on that metric. It is the only place silent data loss shows up.

Worked example: one key through an event-time timeline

Take the job above: one-minute tumbling windows, five seconds of out-of-orderness, thirty seconds of allowed lateness. Follow user A and the window from 10:00:00 to 10:01:00, whose maximum timestamp is 10:00:59.999. For readability, assume a watermark is emitted after every element. In reality watermarks are emitted periodically (every 200 ms by default), so firing can lag slightly behind the table.

ArrivesEvent timeWatermark afterWhat happens
110:00:2010:00:14.999window created, count 1, timers at 10:00:59.999 and 10:01:29.999
210:00:5010:00:44.999count 2
310:01:0310:00:57.999goes to the next window; this one is still open because the watermark is below its end
410:01:0610:01:00.999watermark passes 10:00:59.999: first firing, emits (A, 10:00, 2)
510:00:40unchangedlate but within lateness: count 3, fires again, emits (A, 10:00, 3)
610:01:4010:01:34.999cleanup timer at 10:01:29.999 fires; window state removed. The 10:01 window keeps its own state
710:00:45unchangedwindow gone: element goes to the late side output (or is dropped and counted)

Three things to notice. Element 3 did not fire anything, even though its timestamp is after the window end, because firing is driven by the watermark, not by element timestamps. The downstream sink saw two results for the same window, so it has to be an upsert. And element 7 was only a few seconds later in processing order than element 5, yet it was lost, because the watermark had moved on in between. Lateness is measured against the watermark, not against the wall clock.

Custom triggers and evictors

Override the trigger when "fire once at the end" is wrong, for example a dashboard that wants early results from a one-hour window. Flink ships ContinuousEventTimeTrigger and CountTrigger, and PurgingTrigger turns FIRE into FIRE_AND_PURGE. A custom trigger replaces the default entirely, so it must register the end-of-window timer itself and delete its timers and state in clear(), or they leak for every window.

Evictors are the most expensive option in the API. Because an evictor needs to see all elements, using one prevents any pre-aggregation: even an AggregateFunction falls back to buffering the full window. CountEvictor, TimeEvictor and DeltaEvictor exist; prefer an aggregate over a smaller window when you can.

Session windows and merging

A session window has no fixed bounds. Each element starts a provisional window from its timestamp to timestamp plus gap, and whenever two windows for the same key overlap, Flink merges them into one, merging their state and triggers too. EventTimeSessionWindows.withDynamicGap lets the gap depend on the element, for example a shorter gap for API clients than for humans. Merging is why the AggregateFunction above implements merge, and why a custom trigger used with sessions must support merging.

Two operational facts follow. A late element can bridge two already-emitted sessions, producing a merged result that replaces both. And a session never closes while events keep coming: a bot that calls every few seconds holds one window open forever, so cap sessions with a trigger that fires and purges after a maximum duration, or filter such keys first.

The SQL side: window TVFs

In Flink SQL and the Table API, windows are table-valued functions: TUMBLE, HOP, CUMULATE and SESSION. Each adds window_start, window_end and window_time columns to every row, and you then group by those columns. The documentation describes them as the replacement for the older grouped window functions, and they also feed window joins, Top-N and deduplication.

-- per-user clicks per minute; final result once per window
SELECT user_id, window_start, window_end, COUNT(*) AS clicks
FROM TABLE(
  TUMBLE(TABLE clicks, DESCRIPTOR(event_time), INTERVAL '1' MINUTES))
GROUP BY user_id, window_start, window_end;

-- "today so far", refreshed every 10 minutes: CUMULATE(table, time, step, max size)
SELECT window_start, window_end, SUM(amount) AS revenue
FROM TABLE(
  CUMULATE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '10' MINUTES, INTERVAL '1' DAY))
GROUP BY window_start, window_end;

-- sessions per user with a 30-minute gap
SELECT user_id, window_start, window_end, COUNT(*) AS events
FROM TABLE(
  SESSION(TABLE clicks PARTITION BY user_id, DESCRIPTOR(event_time), INTERVAL '30' MINUTES))
GROUP BY user_id, window_start, window_end;

Window TVF aggregation emits one final result per window when the watermark passes its end, rather than a stream of updates, and it clears its state afterwards. That makes the output append-only, which is easy for sinks. The price is that there is no allowed-lateness equivalent here: data later than the watermark is not reflected. Check the restrictions for your version before you commit. The current documentation says the SESSION TVF is not supported in batch mode, and that session window aggregation gets none of the performance optimisations the other window types get.

State sizing and operations

Window state lives in the configured state backend, covered in Flink State Backends Compared. A back-of-envelope estimate is enough to choose between heap and RocksDB: active keys times live windows per key times bytes per window. For example, two million active users on a one-hour window sliding each minute, with a 48-byte accumulator, is 2,000,000 times 60 times 48 bytes, about 5.8 GB before backend overhead. That is arithmetic, not a measurement, but it is enough to tell you this job belongs on RocksDB and that a buffering ProcessWindowFunction would be unaffordable. Remember that live windows include those held open by allowed lateness.

Each live window also holds timers, which are state and appear in checkpoints; Flink State Architecture explains where they live. Checkpoint size that grows steadily on a windowed job usually means a stalled watermark.

Failure modes

  • No output at all. The watermark is not advancing: an idle partition, timestamps in seconds read as milliseconds, or a source without a WatermarkStrategy.
  • Output stops after a while. One partition went idle or one key's upstream stalled. Add withIdleness, and consider watermark alignment if sources drift far apart.
  • Double counting. Late firings written to an append-only sink. Upsert by key and window start.
  • Silent data loss. numLateRecordsDropped rising with no side output.
  • State explosion. Sliding windows with a small slide plus ProcessWindowFunction, an evictor, or never-ending sessions.
  • Skewed subtask. One hot key makes one subtask far slower, which also holds back the watermark downstream. Pre-aggregate with a salted key, then combine. The same tactic applies to joins, see Stateful stream joins.
  • Wrong day boundaries. Daily windows aligned to UTC when the business day is local. Use the assigner offset or the SQL session time zone.

What to do next

  1. Write down, for each windowed operator, the window size, out-of-orderness bound, allowed lateness and what happens to late data, and agree on those numbers with whoever consumes the output.
  2. Replace any ProcessWindowFunction that only aggregates with an AggregateFunction plus a thin ProcessWindowFunction.
  3. Add withIdleness to every source that can go quiet, and an alert on a window operator's input watermark lagging wall-clock time.
  4. Configure sideOutputLateData and alert on numLateRecordsDropped; make window sinks upserts if allowed lateness is above zero.
  5. Estimate window state with keys times windows times bytes, choose the state backend from that, and recheck when slide or lateness changes.
  6. Replay one hour of production data through a test job twice and confirm identical results; if they differ, find the processing-time dependency.
Key takeaway: A Flink window is per-key state plus timers, created by an assigner, fired by a trigger driven by the watermark, and evaluated by a function. Most production issues come from the watermark (not advancing, or advancing past data you still wanted) and from state that is larger than it needs to be. Use event time with an explicit out-of-orderness bound and idleness, aggregate incrementally, decide deliberately what happens to late data and make sinks idempotent if windows can fire more than once, size state before you deploy, and prefer SQL window TVFs when you only need final aggregates.