skip to content

What is the IQv2 API (KIP-796), and how do you interactively query windowed/session stores versus key-value stores?

level: principalimportance: nice to knowfreq 25%

answer

  1. IQv2 = StateQueryRequest.inStore().withQuery() → StateQueryResult
  2. Query types: KeyQuery / RangeQuery / WindowKeyQuery / custom
  3. partitionResults() + Position bound = read-your-writes
  4. windowed: fetch(key, tFrom, tTo) — key+time, not just key
  5. metadata APIs shared; IQv2 changes read, not locate

basics

~20 s

IQv2 is a newer, type-safe query API (StateQueryRequest + Query types like KeyQuery/RangeQuery) returning a StateQueryResult, designed to be extensible and partition-aware. Windowed stores use a different queryable type (windowStore/sessionStore) and you query by key plus a time range, not just a key.

solid answer

~40 s

Classic IQ (IQv1) hands back fixed store interfaces (ReadOnlyKeyValueStore, ReadOnlyWindowStore, ReadOnlySessionStore) via QueryableStoreTypes. IQv2 (KIP-796) replaces that with a pluggable request/response model: you build a StateQueryRequest.inStore(name).withQuery(query) where query is a Query implementation (KeyQuery.withKey, RangeQuery.withRange, WindowKeyQuery, WindowRangeQuery, or a custom one), optionally restrict to partitions, and get a StateQueryResult mapping partition → QueryResult with positions and execution info. It exposes partition-level results and Position bounds for read-your-writes consistency that IQv1 hid. For windowed data you must query by key and a time window: ReadOnlyWindowStore.fetch(key, timeFrom, timeTo) (IQv1) or WindowKeyQuery/WindowRangeQuery (IQv2). Session stores return sessions via fetch(key). IQv2 also makes building custom store query types feasible.

code

java · 8 lines
java
// IQv2 point lookup with read-your-writes bound
StateQueryRequest<Long> req = StateQueryRequest
    .inStore("word-counts")
    .withQuery(KeyQuery.withKey("hello"))
    .withPositionBound(PositionBound.at(lastSeenPosition));
StateQueryResult<Long> result = streams.query(req);
result.getPartitionResults()
      .forEach((p, qr) -> System.out.println(p + " -> " + qr.getResult()));

go deeper

for a junior

Aware that there is a newer query API and that windowed stores are queried by key plus a time range.

for a middle

Use ReadOnlyWindowStore.fetch(key, tFrom, tTo) and know IQv2 exists with KeyQuery/RangeQuery.

for a senior

Build StateQueryRequest with partition restriction and read per-partition QueryResults; handle window/session semantics.

for a principal

Exploit Position/PositionBound for read-your-writes, custom Query types, and design the scatter-gather + payload/pagination strategy around windowed result sizes.

**Two generations of the API.** *IQv1 (the original).* You name a store, pick a fixed interface via `QueryableStoreTypes`, and call its methods: - `keyValueStore()` → `ReadOnlyKeyValueStore<K,V>`: `get`, `range`, `all`. - `windowStore()` → `ReadOnlyWindowStore<K,V>`: query *by key and time range* — `fetch(key, timeFrom, timeTo)` returns a `WindowStoreIterator<V>` of `(windowTimestamp, value)` pairs; `fetch(keyFrom, keyTo, tFrom, tTo)` and `fetchAll(tFrom, tTo)` scan ranges. The crucial difference from KV: a windowed key is `(key, window)`, so a single application key has *many* values across time windows; you can't just `get(key)`. - `sessionStore()` → `ReadOnlySessionStore<K,V>`: `fetch(key)` returns the session windows for that key. IQv1's limits: the interfaces are closed (you can't add query shapes), results are merged so you lose per-partition visibility, and there's no first-class notion of *how current* the answer is. *IQv2 (KIP-796, evolving through later KIPs).* A pluggable model: ``` StateQueryRequest<V> req = StateQueryRequest .inStore("word-counts") .withQuery(KeyQuery.withKey("hello")); StateQueryResult<V> result = streams.query(req); ``` - **`Query` types**: `KeyQuery.withKey(k)`, `RangeQuery.withRange(from,to)` (and `withNoBounds`), `WindowKeyQuery`, `WindowRangeQuery`, plus the ability to implement your own `Query<R>` — extensibility IQv1 lacked. - **Partition awareness**: `req.withPartitions(set)` or `withAllPartitions()`; the `StateQueryResult` exposes `partitionResults()` (a map of partition → `QueryResult`), so a scatter-gather layer sees exactly which partition produced what, instead of a pre-merged blob. - **`Position` / consistency**: each `QueryResult` carries a `Position` (changelog offsets the store had applied). You can bound a query with `req.withPositionBound(PositionBound.at(position))` to get **read-your-writes**: after producing a record, capture the returned position and require subsequent reads to have at least caught up to it. IQv1 had no such primitive. - **Failure detail**: `QueryResult` distinguishes success from typed failures (e.g. `NOT_UP_TO_BOUND`, `NOT_PRESENT`) per partition. **Why a principal cares.** IQv2 turns IQ from a fixed convenience into a *platform*: custom query types let teams push predicates into the store, partition-level results make the routing/scatter-gather layer precise, and `Position` bounds give a real consistency knob for read-your-writes guarantees across the streams+RPC stack. Windowed/session semantics remain the subtle part regardless of version — the unit of state is `(key, window)`, so query design must always carry a time dimension, and result sizes can be large (many windows per key), which feeds back into RPC payload and pagination decisions. **Migration nuance.** IQv1 and IQv2 coexist; the metadata APIs (`queryMetadataForKey`, `streamsMetadataForStore`) are shared by both for routing — IQv2 changes the *read* call, not the *locate* call.

  • Why can't you use get(key) on a windowed store?
    A windowed store keys data by (key, window), so one application key maps to many values across time windows. You must query with a key plus a time range — fetch(key, timeFrom, timeTo) — which returns an iterator over the matching windows.
  • How does IQv2's Position help with consistency?
    Each QueryResult reports a Position (changelog offsets applied). By capturing the position after a write and using PositionBound on later reads, you require the queried store to be caught up at least to that point, giving read-your-writes semantics IQv1 couldn't express.
  • Do the routing/metadata APIs change between IQv1 and IQv2?
    No. queryMetadataForKey and streamsMetadataForStore are shared for locating partitions/hosts. IQv2 only replaces the local read call (StateQueryRequest/StateQueryResult); your RPC routing layer stays the same.

saying these in an interview costs you the question

  • Trying to get(key) on a windowed store instead of fetch(key, tFrom, tTo)
  • Saying IQv2 replaces the metadata/routing APIs (it doesn't)
  • Claiming IQv1 has read-your-writes / Position bounds (it doesn't)
  • Treating a windowed key as one value rather than many across windows
  • Assuming IQv2 results are pre-merged like IQv1 (they're per-partition)

context