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()));
+    }
 }

Reply via email to