skip to content

In a Flink windowed aggregation, why does adding an Evictor disable pre-aggregation?

level: middleimportance: nice to knowfreq 22%

answer

  1. it needs to look at the records themselves
  2. you cannot inspect what was already folded away
  3. state goes from one value to all of them
  4. runs after the trigger, around the function
  5. upstream filter is usually the better tool

basics

~20 s

An evictor inspects and removes individual elements after the trigger fires, so every element must still exist when the window is evaluated. That forces Flink to buffer the whole window instead of folding records into a single accumulator on arrival.

solid answer

~50 s

An `Evictor` is an optional third component alongside the assigner and the trigger. Its two methods, `evictBefore()` and `evictAfter()`, receive an `Iterable` of the window's elements and can drop individual ones — before the window function runs and/or after. That contract is only satisfiable if the individual elements still exist at firing time, so attaching an evictor **prevents any incremental aggregation**: even if you call `.reduce(...)` or `.aggregate(...)`, the operator must buffer every record in the window rather than folding it into one accumulator. Flink's documentation is explicit that this creates significantly more state. Flink ships three: `CountEvictor` (keep at most N), `DeltaEvictor` (drop elements whose delta from the last one exceeds a threshold) and `TimeEvictor` (keep only elements within an interval of the window's largest timestamp). All three evict before the function by default. Note also that Flink gives no ordering guarantee within a window, so "the first N" is not a meaningful eviction.

code

java · 6 lines
java
// a true sliding-count window: last 100 records, recomputed every 10
input.keyBy(Event::getSensorId)
     .window(GlobalWindows.create())
     .trigger(CountTrigger.of(10))
     .evictor(CountEvictor.of(100))
     .process(new SensorStats());

go deeper

for a junior

Know that an evictor is optional, that it removes individual records from a window around the window function, and that it is rarely needed in everyday jobs.

for a middle

Explain why the Iterable-of-elements signature makes incremental aggregation impossible, and name CountEvictor, DeltaEvictor and TimeEvictor with what each keeps.

for a senior

Be able to say what an evictor does to checkpoint size on a busy window, and to argue for an upstream filter or a different assigner instead in almost every case you are shown.

for a principal

Treat it as a review signal: an evictor in a production topology usually means the window shape was chosen wrongly, and is worth challenging before it becomes the job's dominant state cost.

## The optional third component A Flink windowed transformation always has a **window assigner** and a **trigger**, and it optionally has an **evictor**, attached with `.evictor(...)`. The evictor sits between the trigger's decision to fire and the window function's computation, and its job is to remove individual elements from the pane. ## The interface ```java void evictBefore(Iterable<TimestampedValue<T>> elements, int size, W window, EvictorContext ctx); void evictAfter(Iterable<TimestampedValue<T>> elements, int size, W window, EvictorContext ctx); ``` `evictBefore()` runs after the trigger fires but before the window function is applied — anything it removes is never seen by the computation. `evictAfter()` runs after the function has produced its result, trimming what remains in the window for future firings. By default all three built-in evictors do their work in `evictBefore()`; each also has an `of(..., doEvictAfter)` overload that moves it to `evictAfter()`. ## Why pre-aggregation becomes impossible The signature is the whole answer. Both methods take an `Iterable` of the window's individual elements. For Flink to be able to hand that over, the individual elements must still be in state at firing time — which means the operator cannot have folded them into a single value as they arrived. So the moment you attach an evictor, the incremental path is switched off. Calling `.reduce(...)` or `.aggregate(...)` on an evicting window still produces the right answer, but the operator now buffers every record and applies the reduce or aggregate at firing time over the surviving elements. The state profile changes from *one accumulator per open window* to *every element in every open window* — Flink's documentation warns in as many words that windows with evictors create significantly more state. On a long window over a busy key, that is the difference between a job that checkpoints in seconds and one that does not checkpoint at all. The same is true of an evicting window that also happens to be a sliding window: the element-copies-per-window multiplier and the no-pre-aggregation penalty compound. ## The built-in evictors - **`CountEvictor`** — keeps up to a user-specified number of elements and discards the rest from the beginning of the window buffer. - **`DeltaEvictor`** — takes a `DeltaFunction` and a threshold, computes the delta between the last element in the buffer and each of the others, and removes those whose delta meets or exceeds the threshold. Useful for "only keep readings within X of the most recent one". - **`TimeEvictor`** — in Flink 2.3, `TimeEvictor.of(Duration)` takes the interval as a `java.time.Duration`; it finds the maximum timestamp `max_ts` among the window's elements and removes every element whose timestamp is at or before `max_ts - interval`. This is how you get a trailing sub-window inside a longer or global window. ## The ordering caveat Flink provides no guarantee about the order of elements within a window. `CountEvictor` removes from the beginning of the *window buffer*, but that buffer's order is not necessarily arrival order or timestamp order. So "keep the most recent 100" is not something `CountEvictor` reliably expresses — if you need recency you want `TimeEvictor`, or your own logic in a `ProcessWindowFunction`. The Python DataStream API does not support evictors at all; in Flink 2.3 they are a Java DataStream feature, since the Scala DataStream API was removed in 2.0. ## When to use one — and when not to Evictors are a genuinely narrow tool. The two honest uses are: a `GlobalWindows` assigner with a `CountTrigger` and a `CountEvictor`, giving a true sliding-count window ("the last 100 records per key, recomputed every 10"); and a `TimeEvictor` inside a global or long window when you want a trailing time slice without paying for a sliding assigner's many open windows. For almost everything else, an evictor is the wrong answer to a filtering question. If you want to exclude records, filter them upstream with `filter(...)` before the window — that costs nothing and preserves pre-aggregation. If you want a trailing window, use a sliding assigner or a `KeyedProcessFunction` with your own timers. Reaching for an evictor to drop unwanted records is the classic misuse: it does the filtering at the most expensive possible point in the pipeline, after they have already been stored. ## Answering well Lead with the mechanism — the evictor needs the individual elements, so they cannot be pre-folded — then quantify the consequence in state terms, then name the three built-ins and the narrow cases where an evictor genuinely earns its place. Mentioning that upstream `filter()` is almost always the better tool shows judgment rather than recall.

  • What is the difference between evictBefore() and evictAfter()?
    `evictBefore()` runs after the trigger fires but before the window function is applied, so anything it removes never reaches the computation. `evictAfter()` runs once the function has produced its result and trims what remains in the window for subsequent firings. By default all three built-in evictors — CountEvictor, DeltaEvictor and TimeEvictor — do their work in `evictBefore()`; an `of(..., true)` overload moves each to `evictAfter()`.
  • Why is CountEvictor a poor way to keep the most recent N elements?
    Flink makes no guarantee about the order of elements within a window. CountEvictor discards from the beginning of the window buffer, and that buffer's order is not necessarily arrival or timestamp order, so "the beginning" is not reliably "the oldest". If recency is the requirement, use TimeEvictor or implement the logic yourself in a ProcessWindowFunction.
  • If you just want to exclude certain records from a window, is an evictor the right tool?
    Almost never. A `filter(...)` upstream of the window drops the records before they are ever assigned or stored, costs nothing, and leaves pre-aggregation intact. An evictor filters at the most expensive possible point — after every record has been buffered in window state — and forces the whole window to abandon incremental aggregation to do it.

Pre-aggregation is like tallying receipts into a running total as they come in; an evictor is a rule that says 'before totalling, throw out any receipt over a month old' — which only works if you kept every receipt instead of just the total.

saying these in an interview costs you the question

  • Thinks an evictor is just a filter with no state cost
  • Believes reduce still pre-aggregates when an evictor is attached
  • Claims CountEvictor reliably keeps the newest N records
  • Confuses evicting elements with purging or removing the window
  • Reaches for an evictor when an upstream filter would do

context