This is an automated email from the ASF dual-hosted git repository.
bbejeck pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 24f03e95576 KAFKA-20798: Add readOnly(IsolationLevel) to
TimestampedWindowStoreWithHeaders and SessionStoreWithHeaders (#22806)
24f03e95576 is described below
commit 24f03e95576e0f9b8d20f7d84d4b09ef623dd764
Author: Lucy Liu <[email protected]>
AuthorDate: Tue Jul 14 07:59:36 2026 -0500
KAFKA-20798: Add readOnly(IsolationLevel) to
TimestampedWindowStoreWithHeaders and SessionStoreWithHeaders (#22806)
## Summary
`MeteredTimestampedWindowStoreWithHeaders` and
`MeteredSessionStoreWithHeaders` never override
`readOnly(IsolationLevel)`. Since Kafka Streams' Interactive Query (IQ)
path always calls `.readOnly(level)` on a store before querying it
(`CompositeReadOnlyWindowStore.readOnlyStores()`), both classes silently
fell back to their superclass's generic, headers-ignorant
ReadOnlyView/ReadOnlyView implementation — which deserializes the key
before the value, instead of value-first-then-key (the order this
feature actually requires, since the key's schema-ID header is only
recoverable from the value's embedded headers). This threw
"SerializationException: Error deserializing schema ID /
IllegalArgumentException: Unknown magic byte!" whenever a headers-aware
windowed/session store was queried via IQ after a restart+restore.
## Files changed
1. MeteredTimestampedWindowStoreWithHeaders.java
- Added `readOnly(IsolationLevel)` override + new HeadersReadOnlyView
inner class — the actual fix.
2. MeteredSessionStoreWithHeaders.java
- Same fix: readOnly(IsolationLevel) override + HeadersReadOnlyView
inner class.
3. MeteredTimestampedWindowStoreWithHeadersTest.java
- Added `shouldUseHeadersFromValueToDeserializeKeyInReadOnlyFetchAll`
4. MeteredSessionStoreWithHeadersTest.java
-
Added`shouldUseHeadersFromValueToDeserializeKeyInReadOnlyFindSessions`
test
Reviewers: Nick Telford <[email protected]>, Bill Bejeck
<[email protected]>
---
.../internals/MeteredSessionStoreWithHeaders.java | 97 +++++++++++++++++++
.../MeteredTimestampedWindowStoreWithHeaders.java | 107 +++++++++++++++++++++
.../MeteredSessionStoreWithHeadersTest.java | 29 ++++++
...teredTimestampedWindowStoreWithHeadersTest.java | 32 ++++++
4 files changed, 265 insertions(+)
diff --git
a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeaders.java
b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeaders.java
index cbddac333cc..65758c49591 100644
---
a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeaders.java
+++
b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeaders.java
@@ -16,6 +16,7 @@
*/
package org.apache.kafka.streams.state.internals;
+import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.apache.kafka.common.serialization.Serde;
@@ -35,6 +36,7 @@ import org.apache.kafka.streams.query.WindowRangeQuery;
import org.apache.kafka.streams.query.internals.InternalQueryResultUtil;
import org.apache.kafka.streams.state.AggregationWithHeaders;
import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.ReadOnlySessionStore;
import org.apache.kafka.streams.state.SessionStore;
import org.apache.kafka.streams.state.SessionStoreWithHeaders;
@@ -360,6 +362,101 @@ public class MeteredSessionStoreWithHeaders<K, AGG>
return queryResult;
}
+ @Override
+ public ReadOnlySessionStore<K, AggregationWithHeaders<AGG>> readOnly(final
IsolationLevel isolationLevel) {
+ Objects.requireNonNull(isolationLevel, "isolationLevel cannot be
null");
+ return new ReadOnlyHeadersView(wrapped().readOnly(isolationLevel));
+ }
+
+ private final class ReadOnlyHeadersView implements ReadOnlySessionStore<K,
AggregationWithHeaders<AGG>> {
+
+ private final ReadOnlySessionStore<Bytes, byte[]> underlying;
+
+ ReadOnlyHeadersView(final ReadOnlySessionStore<Bytes, byte[]>
underlying) {
+ this.underlying = underlying;
+ }
+
+ @Override
+ public AggregationWithHeaders<AGG> fetchSession(
+ final K key, final long earliestSessionEndTime, final long
latestSessionStartTime) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return maybeMeasureLatency(
+ () -> deserializeValue(underlying.fetchSession(
+ serializeKey(key, internalContext.headers()),
earliestSessionEndTime, latestSessionStartTime)),
+ time,
+ fetchSensor
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
fetch(final K key) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return new
MeteredSessionStoreWithHeadersIterator(underlying.fetch(serializeKey(key,
internalContext.headers())));
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
backwardFetch(final K key) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return new
MeteredSessionStoreWithHeadersIterator(underlying.backwardFetch(serializeKey(key,
internalContext.headers())));
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
fetch(final K keyFrom, final K keyTo) {
+ return new MeteredSessionStoreWithHeadersIterator(
+ underlying.fetch(serializeKey(keyFrom,
internalContext.headers()), serializeKey(keyTo, internalContext.headers()))
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
backwardFetch(final K keyFrom, final K keyTo) {
+ return new MeteredSessionStoreWithHeadersIterator(
+ underlying.backwardFetch(serializeKey(keyFrom,
internalContext.headers()), serializeKey(keyTo, internalContext.headers()))
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
findSessions(
+ final K key, final long earliestSessionEndTime, final long
latestSessionStartTime) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return new MeteredSessionStoreWithHeadersIterator(
+ underlying.findSessions(serializeKey(key,
internalContext.headers()), earliestSessionEndTime, latestSessionStartTime)
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
backwardFindSessions(
+ final K key, final long earliestSessionEndTime, final long
latestSessionStartTime) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return new MeteredSessionStoreWithHeadersIterator(
+ underlying.backwardFindSessions(serializeKey(key,
internalContext.headers()), earliestSessionEndTime, latestSessionStartTime)
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
findSessions(
+ final K keyFrom, final K keyTo, final long earliestSessionEndTime,
final long latestSessionStartTime) {
+ return new MeteredSessionStoreWithHeadersIterator(
+ underlying.findSessions(
+ serializeKey(keyFrom, internalContext.headers()),
+ serializeKey(keyTo, internalContext.headers()),
+ earliestSessionEndTime,
+ latestSessionStartTime)
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>
backwardFindSessions(
+ final K keyFrom, final K keyTo, final long earliestSessionEndTime,
final long latestSessionStartTime) {
+ return new MeteredSessionStoreWithHeadersIterator(
+ underlying.backwardFindSessions(
+ serializeKey(keyFrom, internalContext.headers()),
+ serializeKey(keyTo, internalContext.headers()),
+ earliestSessionEndTime,
+ latestSessionStartTime)
+ );
+ }
+ }
+
private class MeteredSessionStoreWithHeadersIterator
implements KeyValueIterator<Windowed<K>, AggregationWithHeaders<AGG>>,
MeteredIterator {
diff --git
a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeaders.java
b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeaders.java
index c1ab307001b..eb02894fbf6 100644
---
a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeaders.java
+++
b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeaders.java
@@ -16,6 +16,7 @@
*/
package org.apache.kafka.streams.state.internals;
+import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.apache.kafka.common.serialization.Serde;
@@ -36,6 +37,7 @@ import org.apache.kafka.streams.query.WindowKeyQuery;
import org.apache.kafka.streams.query.WindowRangeQuery;
import org.apache.kafka.streams.query.internals.InternalQueryResultUtil;
import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.ReadOnlyWindowStore;
import org.apache.kafka.streams.state.TimestampedBytesStore;
import org.apache.kafka.streams.state.TimestampedWindowStoreWithHeaders;
import org.apache.kafka.streams.state.ValueAndTimestamp;
@@ -43,6 +45,7 @@ import org.apache.kafka.streams.state.ValueTimestampHeaders;
import org.apache.kafka.streams.state.WindowStore;
import org.apache.kafka.streams.state.WindowStoreIterator;
+import java.time.Instant;
import java.util.Objects;
import java.util.function.Function;
@@ -347,6 +350,110 @@ public class MeteredTimestampedWindowStoreWithHeaders<K,
V>
);
}
+ @Override
+ public ReadOnlyWindowStore<K, ValueTimestampHeaders<V>> readOnly(final
IsolationLevel isolationLevel) {
+ Objects.requireNonNull(isolationLevel, "isolationLevel cannot be
null");
+ return new ReadOnlyHeadersView(wrapped().readOnly(isolationLevel));
+ }
+
+ private final class ReadOnlyHeadersView implements ReadOnlyWindowStore<K,
ValueTimestampHeaders<V>> {
+
+ private final ReadOnlyWindowStore<Bytes, byte[]> underlying;
+
+ ReadOnlyHeadersView(final ReadOnlyWindowStore<Bytes, byte[]>
underlying) {
+ this.underlying = underlying;
+ }
+
+ @Override
+ public ValueTimestampHeaders<V> fetch(final K key, final long
windowStartTimestamp) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return maybeMeasureLatency(
+ () -> {
+ final byte[] result = underlying.fetch(serializeKey(key,
internalContext.headers()), windowStartTimestamp);
+ return result == null ? null : deserializeValue(result);
+ },
+ time,
+ fetchSensor
+ );
+ }
+
+ @Override
+ public WindowStoreIterator<ValueTimestampHeaders<V>> fetch(
+ final K key, final Instant timeFrom, final Instant timeTo) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return new MeteredWindowStoreIterator<>(
+ underlying.fetch(serializeKey(key, internalContext.headers()),
timeFrom, timeTo),
+ fetchSensor,
+ iteratorDurationSensor,
+
MeteredTimestampedWindowStoreWithHeaders.this::deserializeValue,
+ time,
+ numOpenIterators,
+ openIterators
+ );
+ }
+
+ @Override
+ public WindowStoreIterator<ValueTimestampHeaders<V>> backwardFetch(
+ final K key, final Instant timeFrom, final Instant timeTo) {
+ Objects.requireNonNull(key, "key cannot be null");
+ return new MeteredWindowStoreIterator<>(
+ underlying.backwardFetch(serializeKey(key,
internalContext.headers()), timeFrom, timeTo),
+ fetchSensor,
+ iteratorDurationSensor,
+
MeteredTimestampedWindowStoreWithHeaders.this::deserializeValue,
+ time,
+ numOpenIterators,
+ openIterators
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>> fetch(
+ final K keyFrom, final K keyTo, final Instant timeFrom, final
Instant timeTo) {
+ return new
MeteredTimestampedWindowStoreWithHeadersKeyValueIterator(
+ underlying.fetch(
+ serializeKey(keyFrom, internalContext.headers()),
+ serializeKey(keyTo, internalContext.headers()),
+ timeFrom,
+ timeTo)
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>>
backwardFetch(
+ final K keyFrom, final K keyTo, final Instant timeFrom, final
Instant timeTo) {
+ return new
MeteredTimestampedWindowStoreWithHeadersKeyValueIterator(
+ underlying.backwardFetch(
+ serializeKey(keyFrom, internalContext.headers()),
+ serializeKey(keyTo, internalContext.headers()),
+ timeFrom,
+ timeTo)
+ );
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>> all() {
+ return new
MeteredTimestampedWindowStoreWithHeadersKeyValueIterator(underlying.all());
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>>
backwardAll() {
+ return new
MeteredTimestampedWindowStoreWithHeadersKeyValueIterator(underlying.backwardAll());
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>>
fetchAll(
+ final Instant timeFrom, final Instant timeTo) {
+ return new
MeteredTimestampedWindowStoreWithHeadersKeyValueIterator(underlying.fetchAll(timeFrom,
timeTo));
+ }
+
+ @Override
+ public KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>>
backwardFetchAll(
+ final Instant timeFrom, final Instant timeTo) {
+ return new
MeteredTimestampedWindowStoreWithHeadersKeyValueIterator(underlying.backwardFetchAll(timeFrom,
timeTo));
+ }
+ }
+
private class MeteredTimestampedWindowStoreWithHeadersKeyValueIterator
implements KeyValueIterator<Windowed<K>, ValueTimestampHeaders<V>>,
MeteredIterator {
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeadersTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeadersTest.java
index aafa28d00b2..6b9d602cf8e 100644
---
a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeadersTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreWithHeadersTest.java
@@ -16,6 +16,7 @@
*/
package org.apache.kafka.streams.state.internals;
+import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.Metric;
import org.apache.kafka.common.MetricName;
import org.apache.kafka.common.header.Headers;
@@ -47,6 +48,7 @@ import org.apache.kafka.streams.query.QueryResult;
import org.apache.kafka.streams.query.WindowRangeQuery;
import org.apache.kafka.streams.state.AggregationWithHeaders;
import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.ReadOnlySessionStore;
import org.apache.kafka.streams.state.SessionStore;
import org.apache.kafka.test.KeyValueIteratorStub;
@@ -1050,6 +1052,33 @@ public class MeteredSessionStoreWithHeadersTest {
verify(keySerde.deserializer()).deserialize(any(), eq(HEADERS),
eq(KEY.getBytes()));
}
+ @SuppressWarnings("unchecked")
+ @Test
+ public void
shouldUseHeadersFromValueToDeserializeKeyInReadOnlyFindSessions() {
+ setUp();
+ final Serde<String> keySerde = mock(Serde.class);
+ final MeteredSessionStoreWithHeaders<String, String> store =
createStoreWithMockSerdes(keySerde);
+
+ final ReadOnlySessionStore<Bytes, byte[]> readOnlyInner =
mock(ReadOnlySessionStore.class);
+
when(innerStore.readOnly(IsolationLevel.READ_COMMITTED)).thenReturn(readOnlyInner);
+ when(readOnlyInner.findSessions(any(Bytes.class), eq(0L), eq(100L)))
+ .thenReturn(new KeyValueIteratorStub<>(
+ List.of(KeyValue.pair(WINDOWED_KEY_BYTES,
SERIALIZED_VALUE)).iterator()));
+
+ final KeyValueIterator<Windowed<String>,
AggregationWithHeaders<String>> iterator =
+ store.readOnly(IsolationLevel.READ_COMMITTED).findSessions(KEY,
0L, 100L);
+
+ assertTrue(iterator.hasNext());
+ assertEquals(KEY, iterator.peekNextKey().key());
+ final KeyValue<Windowed<String>, AggregationWithHeaders<String>>
result = iterator.next();
+ assertEquals(KEY, result.key.key());
+ assertEquals(AGG_WITH_HEADERS, result.value);
+ assertFalse(iterator.hasNext());
+ iterator.close();
+
+ verify(keySerde.deserializer()).deserialize(any(), eq(HEADERS),
eq(KEY.getBytes()));
+ }
+
@SuppressWarnings("unchecked")
@Test
public void
shouldUseHeadersFromValueToDeserializeKeyInFindSessionsByTime() {
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeadersTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeadersTest.java
index eb2c6a65f8e..1f0dcbb1fdc 100644
---
a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeadersTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStoreWithHeadersTest.java
@@ -16,6 +16,7 @@
*/
package org.apache.kafka.streams.state.internals;
+import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.apache.kafka.common.metrics.MetricConfig;
import org.apache.kafka.common.metrics.Metrics;
@@ -39,6 +40,7 @@ import
org.apache.kafka.streams.processor.internals.ProcessorRecordContext;
import org.apache.kafka.streams.processor.internals.ProcessorStateManager;
import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl;
import org.apache.kafka.streams.state.KeyValueIterator;
+import org.apache.kafka.streams.state.ReadOnlyWindowStore;
import org.apache.kafka.streams.state.ValueTimestampHeaders;
import org.apache.kafka.streams.state.WindowStore;
import org.apache.kafka.test.InternalMockProcessorContext;
@@ -54,6 +56,7 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.mockito.junit.jupiter.MockitoSettings;
import org.mockito.quality.Strictness;
+import java.time.Instant;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -484,4 +487,33 @@ public class MeteredTimestampedWindowStoreWithHeadersTest {
verify(keyDeserializer).deserialize(any(), eq(HEADERS),
eq(KEY.getBytes()));
}
+
+ @Test
+ public void shouldUseHeadersFromValueToDeserializeKeyInReadOnlyFetchAll() {
+ setUp();
+
+ final Windowed<Bytes> windowedKey = new Windowed<>(KEY_BYTES, new
TimeWindow(0, WINDOW_SIZE_MS));
+ final KeyValue<Windowed<Bytes>, byte[]> testData =
KeyValue.pair(windowedKey, VALUE_TIMESTAMP_HEADERS_BYTES);
+
+ final ReadOnlyWindowStore<Bytes, byte[]> readOnlyInner =
mock(ReadOnlyWindowStore.class);
+
when(innerStoreMock.readOnly(IsolationLevel.READ_COMMITTED)).thenReturn(readOnlyInner);
+ when(readOnlyInner.fetchAll(Instant.ofEpochMilli(0),
Instant.ofEpochMilli(100)))
+ .thenReturn(new
KeyValueIteratorStub<>(List.of(testData).iterator()));
+
+ store = createStoreWithMockSerdes();
+
+ final KeyValueIterator<Windowed<String>,
ValueTimestampHeaders<String>> iterator =
+
store.readOnly(IsolationLevel.READ_COMMITTED).fetchAll(Instant.ofEpochMilli(0),
Instant.ofEpochMilli(100));
+
+ assertTrue(iterator.hasNext());
+ assertEquals(KEY, iterator.peekNextKey().key());
+ final KeyValue<Windowed<String>, ValueTimestampHeaders<String>> result
= iterator.next();
+
+ assertEquals(KEY, result.key.key());
+ assertEquals(VALUE_TIMESTAMP_HEADERS, result.value);
+ assertFalse(iterator.hasNext());
+ iterator.close();
+
+ verify(keyDeserializer).deserialize(any(), eq(HEADERS),
eq(KEY.getBytes()));
+ }
}