In a Flink window, why is a ReduceFunction cheaper than a ProcessWindowFunction?
answer
- how much does the operator keep per window
- one value versus every element
- the Iterable is the expensive part
- you can pass both functions at once
- combined form yields a one-element Iterable
basics
~20 sA ReduceFunction aggregates each record into the window as it arrives, so Flink stores one value per window. A ProcessWindowFunction receives an Iterable of every element, so Flink must buffer the whole window's contents in state until it fires.
solid answer
~40 sFlink can compute a window in two ways. With a `ReduceFunction` or an `AggregateFunction`, the operator folds each arriving record into the window's running value immediately, so the state held per window is a **single element or accumulator** regardless of how many records land in it. With a `ProcessWindowFunction`, the operator has to hand the function an `Iterable` of *all* the elements when the window fires, so it must buffer every record until then — state grows linearly with the window's cardinality, and long or high-volume windows become the job's dominant state cost. You do not have to choose, though: passing both, as in `.reduce(reducer, processWindowFunction)`, gives you incremental aggregation *and* the window metadata and per-window state that only `ProcessWindowFunction` exposes. The `ProcessWindowFunction` then receives an `Iterable` containing exactly one pre-aggregated value.
code
java · 11 lines// buffers every element until the window fires
input.keyBy(e -> e.userId)
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
.process(new ProcessWindowFunction<Event, Long, String, TimeWindow>() {
public void process(String key, Context ctx,
Iterable<Event> elements, Collector<Long> out) {
long n = 0;
for (Event e : elements) n++; // whole window is in state
out.collect(n);
}
});go deeper
Be ready to say that a ReduceFunction folds records as they arrive while a ProcessWindowFunction gets all of them at once, and that the second one holds far more state.
Explain the type signatures — ReduceFunction is same-type, AggregateFunction has a separate accumulator — and demonstrate the combined reduce/aggregate plus ProcessWindowFunction form and what the resulting Iterable contains.
Show that you reach for the combined form by default, that you can point to window state as the thing that decides checkpoint size, and that you know per-window state must be cleared explicitly.
Frame this as a state-budget decision across a platform: which teams need element-level access at all, what state per key per window the cluster can afford, and where the answer is a KeyedProcessFunction rather than any window function.
## The two computation models After you have chosen a window assigner, you must attach a **window function** — the code that turns the contents of a window into a result. Flink offers `ReduceFunction`, `AggregateFunction` and `ProcessWindowFunction`, and the difference between the first two and the third is not a matter of API taste. It determines how much state the window operator holds. `ReduceFunction` and `AggregateFunction` are **incremental**. The window operator applies them as each record arrives and keeps only the result so far. `ProcessWindowFunction` is **buffering**: its signature receives an `Iterable<IN>` of the whole window, so nothing can be discarded before the window fires. ## ReduceFunction `ReduceFunction<T>` combines two elements of type `T` into one element of the same type `T`. That same-type constraint is its whole character — it can express a sum, a max, or a "keep the latest" but not an average, because an average needs a running sum *and* a count, which is not the input type. ```java input.keyBy(e -> e.userId) .window(TumblingEventTimeWindows.of(Duration.ofMinutes(5))) .reduce((a, b) -> new Tuple2<>(a.f0, a.f1 + b.f1)); ``` The state for each open window is one `Tuple2`, whether the window received three records or three million. ## AggregateFunction `AggregateFunction<IN, ACC, OUT>` generalises this with three separate types: the input type, an **accumulator** type, and an output type. It has four methods — `createAccumulator()`, `add(value, acc)`, `getResult(acc)` and `merge(a, b)`. The accumulator is free to be a different shape from the input, which is what lets it compute an average (accumulator = sum and count, output = a double), a distinct-count sketch, or a percentile estimate. `merge()` matters specifically for merging assigners: session windows combine two windows by merging their accumulators. Like `ReduceFunction`, the state is one accumulator per open window per key. ## ProcessWindowFunction `ProcessWindowFunction<IN, OUT, KEY, W extends Window>` has a `process(key, context, elements, out)` method. It receives: - the key the window is scoped to, - an `Iterable` of every element in the window, - a `Context` exposing `window()` (with the window's start and end timestamps), `currentProcessingTime()`, `currentWatermark()`, `windowState()`, `globalState()` and `output(...)` for side outputs, - a `Collector` to emit zero or more results. That access is genuinely useful — you cannot label a result with its window's start and end timestamp from a `ReduceFunction`, and you cannot compute a true median without seeing every element. But it costs you: the operator buffers all elements in the window's state, and for a one-hour window over a high-volume key that can be millions of records per key. ## Combining them — the answer interviewers want The framing "incremental *or* full access" is a false choice. Both `reduce` and `aggregate` have an overload that also takes a `ProcessWindowFunction`: ```java input.keyBy(e -> e.userId) .window(TumblingEventTimeWindows.of(Duration.ofMinutes(5))) .aggregate(new AverageAggregate(), new EmitWithWindowBounds()); ``` Flink aggregates incrementally as records arrive, and when the window fires it invokes the `ProcessWindowFunction` with an `Iterable` holding **exactly one** element: the aggregated result. You get the window metadata and per-window state, at one-value-per-window state cost. This combined form should be your default whenever you need the window's timestamps in the output — which is most of the time, because a windowed result stream carries no window information otherwise. ## Per-window state One capability only `ProcessWindowFunction` has is state scoped to *this key and this window instance*, reached through `context.windowState()`, alongside `context.globalState()` for per-key state that outlives the window. It is how you record that a window has already fired once — relevant when a custom trigger fires speculatively or when allowed lateness causes a late firing. Per-window state is not cleaned up for you: you must implement `clear(Context)` to delete it, or it leaks for the lifetime of the job. ## What to say in an interview State the mechanism (incremental fold versus buffered `Iterable`), state the consequence (constant versus linear state per window), name the `merge()` method as the reason `AggregateFunction` works with session windows, and finish by pointing out that the combined `reduce(reducer, processFn)` form removes the tradeoff for the common case. Mentioning `clear()` for per-window state is the detail that separates someone who has read the docs from someone who has operated the job.
- Why can a ReduceFunction not compute an average, while an AggregateFunction can?`ReduceFunction` combines two elements into one of the **same** type, so the intermediate value must have the shape of the input. An average needs a running sum and a count together. `AggregateFunction` has a separate accumulator type, so it can carry `(sum, count)` while inputs are events and the output is a double, with `getResult()` doing the division at the end.
- What does AggregateFunction.merge() do, and when is it actually called?It combines two accumulators into one. It is required by merging window assigners — session windows, above all, which open a window per record and merge windows that fall within the gap. When two session windows merge, their accumulators are merged with this method. For tumbling and sliding windows it is generally never invoked.
- What is per-window state in a ProcessWindowFunction, and why must you clean it up?`context.windowState()` gives keyed state scoped to this key *and* this window instance, useful for recording that the window has already fired when a custom trigger or a late firing causes multiple firings. Flink clears window contents and trigger state automatically but not your per-window state, so you must delete it in `clear(Context)` or it leaks.
saying these in an interview costs you the question
- Claims ProcessWindowFunction also aggregates incrementally
- Thinks you must choose between incremental aggregation and window metadata
- Says ReduceFunction can output a different type than its input
- Believes per-window state is cleaned up automatically
- Assumes buffering cost depends on window duration, not record count