skip to content

Can an external service query a Flink job's keyed state directly, and should it?

level: middleimportance: nice to knowfreq 18%

answer

  1. The feature that did this is deprecated
  2. Off by default, awaiting removal
  3. Make it an output, not internal state
  4. Offline inspection has its own API
  5. SavepointReader for a point-in-time view

basics

~20 s

Technically yes, via Queryable State, but it should not: Flink 2.3 ships it only as a deprecated, off-by-default leftover slated for removal. Publish readable state to an external store, or read a savepoint offline with the State Processor API.

solid answer

~50 s

Flink's **Queryable State** lets an external client read keyed state straight out of the TaskManagers, and in Flink 2.3 the code is still there — `KeyedStream.asQueryableState`, `StateDescriptor.setQueryable`, `QueryableStateClient` — but all of it is `@Deprecated` and slated for removal in a future major version, it has no page in the documentation, it runs only when `queryable-state.enable` is set, and it cannot be used on state with TTL. It was never a good design either: reads are not consistent with checkpoints, clients couple to a job run and its TaskManagers, and serving traffic competes with record processing. The two supported patterns are **push** — emit what should be readable from a `KeyedProcessFunction` into a key/value store, a database or a compacted topic — and **offline inspection** with the State Processor API's `SavepointReader`, which reads a savepoint, not live state.

go deeper

for a junior

Recall that a live external lookup into a running Flink job's state exists only through the deprecated Queryable State, and that the usual answer is to write what needs to be readable out to an external store.

for a middle

Explain what Queryable State is, that it has been deprecated since 1.18 and survives in Flink 2.3 only off by default and slated for removal, and describe both replacements: pushing state out from the job, and reading a savepoint offline with the State Processor API.

for a senior

Argue the design point — consistency, isolation, availability coupling and topology discovery all break when a stream job doubles as a lookup service — and specify how you would bound and monitor the staleness of the pushed copy.

for a principal

Own the boundary: whether a serving store belongs in the streaming platform at all, who owns its availability and schema contract, and how to stop teams from reaching into pipeline internals for data that should be a published output.

## The short answer There is a mechanism, and you should not build on it. Flink's **Queryable State** feature lets an external client look up a key's value directly in a running job's keyed state. It was formally deprecated in Flink 1.18 with notice that it would be removed in a future major version. Flink 2.0 did not remove it, so it is still in Flink 2.3, but in a state that rules it out for new work: - Every public entry point — `KeyedStream.asQueryableState(...)`, `StateDescriptor.setQueryable(...)`, the `QueryableStateClient` and `QueryableStateOptions` — is marked `@Deprecated`, with removal still announced for a future major version. - It is off by default: the TaskManager-side proxy and server start only when `queryable-state.enable` is `true`, and the runtime ships as an optional `flink-queryable-state-runtime` jar in the distribution's `opt/` directory. - It no longer has a page in the Flink documentation. - It cannot be combined with state TTL: `setQueryable` rejects a descriptor that has TTL enabled. Interviewers ask this because the instinct is natural: the job already holds a per-key value in an embedded key/value store, so why not read it? The expected answer is that you still can, barely, and should not. ## Why it was a bad idea even when it was supported Four reasons, all worth being able to state: **Consistency.** A running job's state is mid-flight. It reflects everything processed so far, including work that has not been committed by a checkpoint and would be rolled back on failure. A reader gets a value that may never have existed in any consistent snapshot of the job. **Isolation.** Serving lookups from the TaskManagers puts external query traffic on the same processes — and, for the state access itself, competing with the same RocksDB instances and memory budget — that are doing record processing. A burst of queries becomes a throughput problem for the pipeline, and vice versa. **Availability coupling.** The state lives inside job execution. Redeploy the job, rescale it, or let it fail over, and the endpoints move, keys land on different subtasks, and the reader sees gaps. External consumers expect an availability model a stream job does not offer. **Topology coupling.** The client connects to a proxy on some TaskManager by host and port and names the job by its `JobID`; the proxy asks the JobManager which TaskManager holds the key's key group and caches that location. The client never picks the subtask itself, but it is bound to one job run's ID and to TaskManager endpoints — exactly the sort of coupling you do not want in a service boundary. ## Pattern one: push what should be readable The idiomatic answer is that if a value should be readable from outside, it is an output, not internal state. Emit it. A `KeyedProcessFunction` that maintains a `ValueState` can also `collect` the updated value whenever it changes — or on a timer, if you want to rate-limit the output — into a sink that writes a key/value store, a database, or a compacted topic. Readers query that store. This inverts every problem above. The store's consistency and availability are its own concern; serving load never touches the TaskManagers; the shape of the served record is a deliberate contract rather than an accident of your state layout; and you can change the internal state representation — swap `ValueState<Map<K,V>>` for `MapState<K,V>`, add a field, change the backend — without breaking a single consumer. The cost is a second system to run and an extra write path, and the output is eventually consistent with the stream by design. Bound that lag explicitly rather than pretending it is zero. ## Pattern two: read a savepoint offline When the requirement is inspection rather than serving — "what does the job think this user's state is?", "how many keys are we holding?", "why did this session never close?" — the State Processor API is the right tool. Take a savepoint, then load it in a batch job with `SavepointReader` and read the keyed state out as ordinary records that you can filter, count or export. The same API works in the other direction: `SavepointWriter` produces a savepoint from batch-computed data, which is how you bootstrap a new job with state derived from history rather than waiting for it to warm up from the live stream. It is also the escape hatch for state migrations Flink will not perform automatically — changing a field's type, re-keying, or changing maximum parallelism. What it is not: a live lookup. It operates on a savepoint file, so it is always a point-in-time view and always costs a batch run. ## How to answer it Say that Queryable State survives in Flink 2.3 only as a deprecated leftover — off by default, undocumented, incompatible with TTL and slated for removal — and that moving away from it is a design correction rather than a regression. Then give the two patterns and be explicit about which requirement each serves: push to an external store for low-latency serving, State Processor API for offline inspection and bootstrapping. If pressed on the trade-off, the honest one is that the push pattern trades an extra system and a bounded staleness window for isolation, a stable contract, and freedom to evolve the job's internals — which is nearly always the right trade. ## A neighbouring point of confusion Other stream processors made different choices here, and candidates sometimes import the vocabulary. Serving reads directly out of a stream processor's local state store is a real design in other systems; it is not the model current Flink supports, and answering the Flink question with another engine's mechanism is a common and visible slip.

  • What is the State Processor API actually for?
    Reading and writing savepoints as batch data. `SavepointReader` loads a savepoint and exposes its keyed and operator state as records you can inspect, count or export — useful for debugging state growth or auditing what a job believes. `SavepointWriter` goes the other way, producing a savepoint from computed data, which lets you bootstrap a new job with historical state or perform migrations Flink will not do automatically, such as changing a field type or maximum parallelism.
  • If you push state out to an external store, how do you bound the staleness?
    Decide when the job emits. Emitting on every state change gives the tightest lag at the cost of write volume; emitting from a timer gives a bounded, predictable refresh interval and collapses bursts. Either way, stamp each record with the event time or processing time it reflects so readers can see the lag rather than guess, and monitor the end-to-end delay between an input event and the corresponding row being readable.
  • Why is serving read traffic from the TaskManagers a bad idea beyond consistency?
    It couples two workloads with different shapes. Query bursts compete with record processing for CPU and, on `EmbeddedRocksDBStateBackend`, for the same block cache and memory budget, so serving latency degrades pipeline throughput and vice versa. Availability is coupled too: a redeploy, a rescale or a failover moves key ownership between subtasks, which readers experience as an outage rather than a routine job operation.

saying these in an interview costs you the question

  • Proposes Queryable State as a supported option for new work
  • Believes Flink 2.0 already deleted Queryable State from the codebase
  • Thinks the State Processor API reads live running state
  • Assumes reading job state is consistent with checkpoints
  • Imports another engine's local-store serving model into Flink

context