In Flink, what is the difference between keyed state and operator state?
answer
- Two families, not one
- One of them needs keyBy first
- The other belongs to a subtask
- CheckpointedFunction and the OperatorStateStore
- Key groups versus even-split or union
basics
~10 sKeyed state is scoped to the key of the record being processed and is only reachable on a KeyedStream after keyBy. Operator state is scoped to one parallel operator instance, with no key involved.
solid answer
~50 sFlink has two families of managed state. **Keyed state** lives in a per-key key/value store: it is only reachable on a `KeyedStream` (that is, after `keyBy`), and every access — `ValueState.value()`, `MapState.get()` — is implicitly scoped to the key of the record currently being processed. Flink organises it into *key groups* so it can redistribute it automatically when you rescale. **Operator state** (non-keyed state) is bound to one parallel subtask instead of a key; you obtain it by implementing `CheckpointedFunction` and asking the `OperatorStateStore` for a `ListState`. On rescale Flink redistributes it either by even-split (`getListState`) or by union, where every subtask gets the whole list (`getUnionListState`). Application logic almost always uses keyed state; operator state is mostly for sources and sinks that must remember something with no natural key, such as buffered records or a partition-offset map.
code
java · 21 linespublic class Dedupe extends RichFlatMapFunction<Event, Event> {
private transient ValueState<Boolean> seen;
@Override
public void open(OpenContext ctx) {
seen = getRuntimeContext().getState(
new ValueStateDescriptor<>("seen", Boolean.class));
}
@Override
public void flatMap(Event e, Collector<Event> out) throws Exception {
// scoped to e's key automatically -- no key is passed in
if (seen.value() == null) {
seen.update(true);
out.collect(e);
}
}
}
stream.keyBy(Event::getId).flatMap(new Dedupe());go deeper
Be ready to say which of the two families needs keyBy and name the primitives: ValueState, ListState, MapState, ReducingState, AggregatingState. Knowing that a state handle is implicitly scoped to the current record's key is the thing being tested.
Explain the mechanics: StateDescriptor plus RuntimeContext for keyed state, CheckpointedFunction plus OperatorStateStore for operator state, and why keyed access needs no key argument. Be able to describe both redistribution schemes for operator list state.
Show judgment about which family fits a given component. Justify why application logic should almost always be keyed, and explain the operational risk of union list state as its cardinality grows with input volume.
Own the modelling decision: choosing the key defines the unit of parallelism, the unit of state isolation, and the unit of rescaling for the whole pipeline. Be ready to argue when a design should introduce a key rather than lean on non-keyed state or broadcast state.
## Two families of managed state Flink calls state *managed* when the runtime knows its structure, serialises it itself, and can snapshot and redistribute it. Managed state comes in exactly two families, and which one you get is decided by whether the stream is keyed. ## Keyed state Keyed state can only be used on a `KeyedStream` — the stream you get back from `keyBy(KeySelector)`. `keyBy` performs a partitioned data exchange so that every record with the same key lands on the same parallel instance of the downstream operator. Because of that alignment, all state updates are local; Flink needs no transactions to keep them consistent. The crucial property is that keyed state access is *implicitly scoped*. You never pass a key to the state object. When a record with key `"user-42"` is being processed, `ValueState.value()` returns the value stored for `"user-42"`, and `update(...)` writes to it. Process the next record with a different key and the same `ValueState` handle now sees a completely different value. Think of the handle as a cursor into an embedded key/value store whose current position is set by the record in flight. The available primitives are `ValueState<T>` (one value, `value()` / `update(T)`), `ListState<T>` (`add(T)`, `addAll(List<T>)`, `update(List<T>)`, `get()` returning an `Iterable`), `MapState<UK,UV>` (`put`, `putAll`, `get`, `entries`, `keys`, `values`, `isEmpty`), `ReducingState<T>` (values added with `add(T)` are folded together by a `ReduceFunction`, so the aggregate type equals the element type) and `AggregatingState<IN,OUT>` (same idea but an `AggregateFunction` lets the accumulator and output types differ). All of them have `clear()`, which clears the state *for the current key only*. You get a handle by building a `StateDescriptor` — `ValueStateDescriptor`, `ListStateDescriptor`, `MapStateDescriptor`, `ReducingStateDescriptor` or `AggregatingStateDescriptor` — and passing it to the `RuntimeContext`, for example `getRuntimeContext().getState(descriptor)`. That means state is only available in *rich functions* (`RichFlatMapFunction`, `KeyedProcessFunction`, and so on), because only those expose a `RuntimeContext`. The usual place to fetch the handle is `open(...)`, not the per-record method. Each state has a name that must be unique within the operator; the name is what Flink uses to find the bytes again on restore. ## Operator state Operator state, also called non-keyed state, is bound to one parallel operator instance. There is no key, so there is no implicit scoping — the whole state belongs to that subtask. You reach it by implementing `CheckpointedFunction`, which requires two methods: `snapshotState(FunctionSnapshotContext)`, called whenever a checkpoint is taken, and `initializeState(FunctionInitializationContext)`, called on first start *and* on restore. Inside `initializeState` you call `context.getOperatorStateStore().getListState(descriptor)` and, if `context.isRestored()` is true, read the recovered elements back into your own fields. Operator state is list-shaped on purpose. The list elements must be independent, serialisable objects, because the element is the finest granularity at which Flink can redistribute non-keyed state when parallelism changes. ## Redistribution on rescale This is where the two families differ most. Keyed state is partitioned into key groups, and Flink reassigns whole key groups to subtasks when parallelism changes — you do nothing. Operator state offers two schemes, chosen by which accessor you call: - **Even-split redistribution** (`getListState`): all subtasks' lists are logically concatenated, then split evenly across the new set of subtasks. A subtask may end up with zero, one, or several elements. - **Union redistribution** (`getUnionListState`): every subtask receives the *complete* concatenated list. Useful when each instance needs the global picture, but dangerous at high cardinality — checkpoint metadata stores an offset per list entry, which can blow up RPC frame sizes or memory. The naming convention is deliberate: the method name carries the redistribution pattern, and a bare `getListState` means even-split. ## Broadcast state Broadcast state is a special, map-shaped form of operator state. It exists so a low-throughput stream — a rule set, a feature flag table — can be broadcast to every downstream subtask and kept identical there, then consulted while processing records of a second, non-broadcast stream. Unlike ordinary operator state it is map-shaped, it is only available to operators that have one broadcast and one non-broadcast input, and one operator may hold several named broadcast states. ## Choosing in practice In a typical stateful application you do not need operator state at all. If the thing you are remembering belongs to an entity — a user, a device, a session — key by that entity and use keyed state; you get automatic rescaling and per-key isolation for free. Reach for operator state only when there is genuinely no key: buffering records inside a sink before a flush, or a source subtask remembering which input splits and offsets it owns.
- Why can't you call getRuntimeContext().getState(...) from a plain MapFunction?State is reached through the `RuntimeContext`, which only *rich* functions expose — `RichMapFunction`, `RichFlatMapFunction`, `KeyedProcessFunction` and friends. A plain `MapFunction` has no `RuntimeContext`, so there is nowhere to register a `StateDescriptor`. On top of that, keyed state additionally requires the stream to already be a `KeyedStream`; calling it on a non-keyed stream fails at runtime even from a rich function.
- When would union list state be the wrong choice for operator state?Whenever the list can grow large. With union redistribution every subtask receives the full concatenated list on restore, and checkpoint metadata stores an offset for each entry, so a high-cardinality list can hit RPC frame-size limits or exhaust memory on the JobManager. Use it only for small, genuinely global lists; prefer even-split redistribution, or key the data and use keyed state, when cardinality grows with the input.
- What makes broadcast state different from ordinary operator state?Broadcast state is map-shaped rather than list-shaped, it is only available on operators that combine one broadcast input with one non-broadcast input, and a single operator can hold several differently named broadcast states. Its purpose is to keep an identical copy of a small, slow-changing dataset — rules, thresholds, lookup tables — on every subtask so the main stream can consult it locally.
Keyed state is a filing cabinet where the drawer is chosen for you by the record currently on your desk; operator state is the private notepad on that one desk, shared with nobody and split among desks if the office is rearranged.
saying these in an interview costs you the question
- Says any function can use ValueState without keyBy
- Thinks operator state is also scoped per record key
- Claims you pass the key into ValueState.value()
- Believes clear() wipes state for all keys
- Says operator state cannot be redistributed on rescale