skip to content

In Flink windowing, what is the difference between a Trigger returning FIRE and FIRE_AND_PURGE?

level: middleimportance: should knowfreq 36%

answer

  1. one of them forgets what it just saw
  2. matters only when a window fires twice
  3. cumulative totals versus disjoint deltas
  4. the window itself survives either way
  5. a wrapper turns any trigger into the purging kind

basics

~20 s

Both emit a result for the window. FIRE leaves the window's buffered elements in place, so a later firing recomputes over them again; FIRE_AND_PURGE clears the contents afterwards, so the next firing starts from empty. Neither removes the window itself.

solid answer

~50 s

A `Trigger` decides *when* a window is ready to be computed, and its `onElement`, `onEventTime` and `onProcessingTime` callbacks each return a `TriggerResult`: `CONTINUE`, `FIRE`, `PURGE` or `FIRE_AND_PURGE`. `FIRE` runs the window function and emits a result while **keeping** the window's contents, so if the window fires again the result covers everything accumulated so far — the right behaviour for speculative early results that get superseded by a final one. `FIRE_AND_PURGE` emits and then **discards** the contents, so the next firing covers only records that arrive afterwards — the right behaviour when firings should partition the data rather than accumulate. Crucially, purging removes only the elements: the window's metadata and the trigger's own state survive, and new records can still be assigned to that window. The basic built-in triggers (`EventTimeTrigger`, `ProcessingTimeTrigger`, `CountTrigger`) return plain `FIRE`; `PurgingTrigger.of(...)` wraps another trigger to convert its firings into purging ones.

code

java · 11 lines
java
// replaces EventTimeTrigger entirely -- window no longer fires on watermark
input.keyBy(Event::getUserId)
     .window(TumblingEventTimeWindows.of(Duration.ofMinutes(10)))
     .trigger(CountTrigger.of(100))
     .aggregate(new CountAgg());

// same count trigger, but each firing clears the window contents
input.keyBy(Event::getUserId)
     .window(TumblingEventTimeWindows.of(Duration.ofMinutes(10)))
     .trigger(PurgingTrigger.of(CountTrigger.of(100)))
     .aggregate(new CountAgg());

go deeper

for a junior

Know that a trigger is what decides when a window emits, and that FIRE_AND_PURGE additionally throws away what the window had buffered.

for a middle

List the four TriggerResult values and explain the downstream consequence: FIRE produces cumulative results that must be upserted, FIRE_AND_PURGE produces deltas that must be summed.

for a senior

Show you know a custom trigger replaces the default rather than composing with it, and that you would check what the sink does with repeated firings before choosing.

for a principal

Frame early firing as a latency-versus-correctness contract with consumers: how often speculative results are emitted, how they are identified as provisional, and who is responsible for reconciling them.

## Where the trigger sits A windowed pipeline has three moving parts. The **assigner** decides which window a record belongs to. The **trigger** decides when that window is ready for computation. The **window function** does the computation. Every `WindowAssigner` ships with a sensible default trigger, and you only supply your own via `.trigger(...)` when the default is wrong for your use case. ## The Trigger interface `Trigger` has five callbacks: - `onElement()` — called for every record added to the window. - `onEventTime()` — called when a registered event-time timer fires. - `onProcessingTime()` — called when a registered processing-time timer fires. - `onMerge()` — called when two windows merge, so stateful triggers can combine their state (this is what makes session windows work). - `clear()` — called when the window is removed, so the trigger can delete its own state and timers. The first three return a `TriggerResult`. Any of them may also register future timers through the context they receive. ## The four TriggerResult values - **`CONTINUE`** — do nothing; the window is not ready. - **`FIRE`** — compute and emit the window result, keeping the buffered contents. - **`PURGE`** — discard the window's contents without emitting anything. - **`FIRE_AND_PURGE`** — emit, then discard the contents. ## FIRE versus FIRE_AND_PURGE in practice The distinction only becomes visible when a window fires **more than once**, which happens for two reasons: a custom trigger that fires speculatively (say, every 30 seconds inside a 10-minute window so a dashboard updates early), or a late firing caused by allowed lateness. With `FIRE`, each firing recomputes over everything the window has accumulated, so successive outputs are **cumulative**: `count=40`, then `count=95`, then `count=130`. A downstream consumer must treat each output as an *update* that supersedes the previous one for that window — an upsert keyed on the window, not an append. With `FIRE_AND_PURGE`, each firing empties the window afterwards, so successive outputs are **disjoint deltas**: `count=40`, then `count=55`, then `count=35`. A downstream consumer sums them. Which is correct depends entirely on what the sink does with the records, and getting it backwards produces a metric that is either wildly overstated or quietly wrong. Note what purging does *not* do: it removes the elements only. The window's metadata and the trigger's state stay intact, so records can continue to be assigned to that same window afterwards. Window *removal* is a separate lifecycle event, driven by time passing the window's end plus allowed lateness. A subtlety for incremental functions: a window using `ReduceFunction` or `AggregateFunction` holds an accumulator rather than a list of elements, and purging discards that accumulator, so the next firing starts from a fresh one. ## Default triggers Every event-time window assigner defaults to `EventTimeTrigger`, which fires once when the watermark passes the end of the window — a plain `FIRE`, not a purge. Processing-time assigners default to `ProcessingTimeTrigger`. `GlobalWindows.create()` installs `NeverTrigger`, which never fires at all, which is why such a global window emits nothing until you supply a trigger of your own. In Flink 2.3 the `GlobalWindows.createWithEndOfStreamTrigger()` variant instead fires once, when a bounded input ends. Built-in triggers you can use directly: - `EventTimeTrigger` — fires when the watermark passes the window end. - `ProcessingTimeTrigger` — fires on processing-time progress. - `CountTrigger` — fires when the window's element count reaches the given number, then resets its count, so it fires every N elements. - `PurgingTrigger.of(otherTrigger)` — wraps another trigger and converts its `FIRE` results into `FIRE_AND_PURGE`. ## The overriding gotcha Calling `.trigger(...)` **replaces** the assigner's default trigger; it does not add to it. If you set `CountTrigger.of(100)` on `TumblingEventTimeWindows`, you lose `EventTimeTrigger` entirely — the window no longer fires when the watermark passes its end, only when it reaches 100 elements. A window that never reaches 100 elements never emits. If you want both watermark-based and count-based firing you have to write a custom trigger that combines them; no built-in trigger does this. The one stock wrapper that adds a second condition, `ProcessingTimeoutTrigger`, adds a processing-time timeout, not watermark firing. ## What to say Define the trigger's role, list the four `TriggerResult` values, then explain the practical consequence: `FIRE` gives cumulative outputs the sink must upsert, `FIRE_AND_PURGE` gives disjoint deltas the sink must sum. Adding that `.trigger(...)` overrides rather than augments the default is the detail that shows you have actually shipped a custom trigger.

  • What actually survives a PURGE, and what does not?
    Purging removes only the window's buffered elements or its accumulator. The window's metadata and the trigger's own state remain, so records can still be assigned to that same window afterwards and the trigger keeps its counters and timers. Removing the window entirely is a separate lifecycle step, driven by time passing the window's end plus any allowed lateness.
  • What happens if you set CountTrigger on a TumblingEventTimeWindows assigner?
    The count trigger **replaces** the default `EventTimeTrigger` rather than supplementing it, so the window no longer fires when the watermark passes its end — only each time the element count reaches the threshold. A window that never reaches the threshold never emits. Getting both behaviours requires writing a custom trigger; no built-in trigger combines watermark and count firing.
  • Which trigger does GlobalWindows.create() install, and why does that matter?
    `NeverTrigger`, which never returns FIRE. A global window has no natural end for a time-based trigger to key off, so without a custom trigger the job accumulates state forever and emits nothing. In practice you pair `GlobalWindows.create()` with `CountTrigger` — often wrapped in `PurgingTrigger` — to get count-based batches; for a bounded input, `createWithEndOfStreamTrigger()` fires once at end of input instead.

saying these in an interview costs you the question

  • Thinks PURGE deletes the window itself, not just its contents
  • Believes a custom trigger is added alongside the assigner's default
  • Cannot say why multiple firings would ever happen
  • Assumes downstream can append every firing regardless of purge choice
  • Expects GlobalWindows.create() to fire without a configured trigger

context