In Flink's DataStream API, what does keyBy() produce and why is it required before keyed state?
answer
- the return type changes, not just the routing
- one key, always the same instance
- state has to live somewhere predictable
- arrays and hashCode-less POJOs are rejected
- it is also a network shuffle
basics
~10 skeyBy() turns a DataStream into a KeyedStream by hash-partitioning records so every record sharing a key reaches the same parallel subtask. Keyed state, keyed timers and keyed windows exist only on a KeyedStream.
solid answer
~40 s`keyBy(KeySelector)` converts a `DataStream` into a `KeyedStream`. It is a logical partitioning: Flink hashes the extracted key and routes every record carrying that key to the same parallel subtask, so one key's records are always handled by one instance. That co-location is what makes keyed state safe — a `ValueState` inside a `KeyedProcessFunction` is automatically scoped to the current key, and event-time or processing-time timers are registered per key. This is why `reduce()`, `sum()`, `window()` and stateful `process()` only appear after `keyBy()`. Two API constraints bite in practice: the key type may not be an array, and a POJO used as a key must override `hashCode()` rather than inherit `Object.hashCode()`. `keyBy()` also forces a network shuffle, so operator chaining breaks at that point.
code
java · 8 linesDataStream<Click> clicks = env.fromSource(source, WatermarkStrategy.noWatermarks(), "clicks");
// unkeyed: no keyed state, no timers
clicks.process(new ProcessFunction<Click, String>() { /* ... */ });
// keyed: ValueState and TimerService are scoped to the current key
clicks.keyBy(Click::getUserId)
.process(new KeyedProcessFunction<String, Click, String>() { /* ... */ });go deeper
Be ready to say that keyBy() returns a KeyedStream and that all records with the same key go to the same parallel instance. Know that windows, reduce and keyed state only work after it.
Explain the mechanics: hash of the selector's key into key groups, key groups assigned to subtasks, state scoped automatically per key on each call. Mention that keyBy forces a shuffle and breaks chaining.
Show production judgment about key choice: cardinality below your parallelism wastes slots, one dominant key creates a permanently hot subtask, and both show up as one lagging instance in the web UI rather than as an error.
Own the long-lived consequence: maximum parallelism fixes the key-group count for the life of the job's state, so the key design and the scaling ceiling are one decision made before the first savepoint exists.
## What keyBy() actually does In Flink's DataStream API a `DataStream<T>` is an immutable, possibly unbounded collection of records spread over N parallel subtasks. Nothing about a plain `DataStream` guarantees where a given record lands — a source may emit round-robin, a `rebalance()` deliberately scatters records, and a `map()` simply forwards whatever arrives. `keyBy()` changes that. You give it a `KeySelector` — usually a lambda such as `value -> value.getUserId()` — and Flink computes a key for every record, hashes it, and uses that hash to choose the downstream subtask. The result type changes: `DataStream<T>.keyBy(...)` returns `KeyedStream<T, K>`. That type change is the whole point. It is the API's way of proving, at compile time, that every record with a given key will be seen by exactly one parallel instance of the next operator. ```java DataStream<Click> clicks = ...; KeyedStream<Click, String> byUser = clicks.keyBy(Click::getUserId); ``` ## Why keyed state depends on it Flink's *keyed state* — `ValueState`, `ListState`, `MapState` and friends — is not a variable in your function object. It is a per-key entry that the runtime swaps in before each call. When `processElement()` runs for a record whose key is `"u-42"`, every state handle you read inside that call already points at `"u-42"`'s slot; you never pass the key yourself. That trick is only sound if all of a key's records arrive at the same subtask, which is exactly what `keyBy()` guarantees. Consequently the API refuses to let you register keyed state or a timer on an unkeyed stream: `ProcessFunction` on a plain `DataStream` gives you events but no state and no `TimerService` timers. Apply the same function after `keyBy()` and both appear. The same reasoning explains which operators are only available on `KeyedStream`: rolling `reduce()`, `sum()`/`min()`/`max()`, `window()` (as opposed to the parallelism-1 `windowAll()`), and any stateful `process()`. ## Key groups and rescaling Flink does not map a key directly to a subtask. It maps the key's hash into one of a fixed number of *key groups*, and assigns ranges of key groups to subtasks. The number of key groups equals the operator's **maximum parallelism** (`setMaxParallelism()`), which defaults to roughly `parallelism + parallelism/2`, clamped to a lower bound of 128 and an upper bound of 32768. This indirection is what lets you restart a job from a savepoint at a different parallelism: whole key groups move between subtasks, and each subtask's state can be reassembled from key-group ranges. It also means max parallelism is a hard ceiling you cannot raise later without breaking state compatibility, so it is worth setting deliberately on long-lived jobs. ## What makes a bad key The documentation names two hard rules. A key may not be an **array** of any type, because array `hashCode()` is identity-based and two equal arrays would hash differently. And a **POJO** key must override `hashCode()` — relying on `Object.hashCode()` produces per-instance hashes and scatters logically identical keys across subtasks. Flink rejects both cases rather than silently misrouting. Beyond correctness, key *cardinality and distribution* decide whether the job works at scale. A key with very few distinct values (say, `country` with three values) cannot use more subtasks than it has key groups populated, so most of your parallelism sits idle. A key with one dominant value — a bot account, a null-ish sentinel, a single hot product — sends a disproportionate share of traffic to one subtask, and that subtask becomes the bottleneck for the whole pipeline while its peers are idle. ## The cost `keyBy()` is not free. Operations like `map()`, `flatMap()` and `filter()` have a one-to-one connection pattern, so Flink chains them into a single task running in one thread with no serialization between steps. `keyBy()` requires records to move between parallel instances, which means serialization, a network hop, and a break in the operator chain. In a job graph you can read the chain boundaries directly: each `keyBy()` or `rebalance()` starts a new task. That is a reason to key once and keep working on the `KeyedStream` where possible, rather than re-keying on the same field repeatedly, and a reason why an unnecessary `keyBy()` before a stateless transformation is pure overhead.
- Why can't you register a keyed timer on a stream that has not been keyed?Timers are stored per key alongside keyed state, and the runtime fires `onTimer()` with that key's state already scoped in. On an unkeyed stream there is no key to scope to and no guarantee the same subtask would see the key's later records, so the `TimerService` offers no keyed timers there. Apply the function after `keyBy()` and both keyed state and timers become available.
- What happens to a KeyedStream job's state when you restart it at a higher parallelism?Flink redistributes key groups, not individual keys. Each operator has a fixed number of key groups equal to its maximum parallelism, and each subtask owns a contiguous range. Restarting at a new parallelism reassigns those ranges and restores each subtask's state from the savepoint. You can only scale up to the max parallelism, and changing max parallelism itself makes the savepoint incompatible.
- Why would an array field make a poor key even if Flink allowed it?Array `hashCode()` in Java is identity-based, so two arrays holding identical contents hash differently. Records that are logically the same key would be routed to different subtasks, splitting their state and silently producing wrong aggregates. Flink rejects array keys outright for this reason; use a value-typed wrapper or a derived string key instead.
A KeyedStream is like assigning every customer a fixed teller: whichever branch they walk into, the routing sends them to the one desk that already holds their file, so the file never has to be merged from two desks.
saying these in an interview costs you the question
- Says keyBy physically moves records to a single machine per key
- Claims keyBy sorts the stream by key
- Thinks keyed state is a field on the function object
- Believes keyBy is free and does not break operator chaining
- Says any object works as a key regardless of hashCode