From f4c98df8e8b772ceeda210d93fc4b7e74ad403cc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=A8=B1=E5=85=83=E8=B1=AA?= <146086744+edenfunf@users.noreply.github.com> Date: Thu, 9 Jul 2026 12:54:19 +0800 Subject: [PATCH 1/2] KAFKA-20597: Fix prefixScan overflow for prefixes with no upper bound in in-memory state stores When a prefix has no lexicographically larger successor (e.g. it is all 0xFF bytes), incrementing it to compute the exclusive upper bound of the scan range overflows. RocksDBStore already handled this by treating the overflow as an unbounded range, but InMemoryKeyValueStore, MemoryNavigableLRUCache and CachingKeyValueStore called ByteUtils.increment() directly and threw IndexOutOfBoundsException. Centralize the overflow-safe increment as ByteUtils.incrementWithoutOverflow(), which returns null on overflow, and reuse it from the RocksDB code paths that previously defined a private copy. The in-memory stores now treat a null upper bound as "scan to the end of the keyspace" (tailMap), matching RocksDB. Add a regression test to AbstractKeyValueStoreTest so it runs against every KeyValueStore implementation, using IntegerSerializer's encoding of -1 (0xFFFFFFFF) to reproduce the overflow, plus unit tests for ByteUtils.incrementWithoutOverflow(). --- .../common/utils/internals/ByteUtils.java | 17 +++++++++++++ .../common/utils/internals/ByteUtilsTest.java | 15 ++++++++++++ .../state/internals/CachingKeyValueStore.java | 4 +++- .../internals/DualColumnFamilyAccessor.java | 2 +- .../internals/InMemoryKeyValueStore.java | 9 +++++-- .../internals/LogicalKeyValueSegment.java | 2 +- .../internals/MemoryNavigableLRUCache.java | 10 ++++++-- .../streams/state/internals/RocksDBStore.java | 17 +------------ .../internals/AbstractKeyValueStoreTest.java | 24 ++++++++++++++++++- 9 files changed, 76 insertions(+), 24 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/internals/ByteUtils.java b/clients/src/main/java/org/apache/kafka/common/utils/internals/ByteUtils.java index 9a20b8cb91fa2..ad58d661a442f 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/internals/ByteUtils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/internals/ByteUtils.java @@ -66,6 +66,23 @@ public static Bytes increment(Bytes input) throws IndexOutOfBoundsException { } } + /** + * Same as {@link #increment(Bytes)} but returns {@code null} instead of throwing + * {@code IndexOutOfBoundsException} when incrementing would overflow the byte array + * (i.e. every byte is {@code 0xFF}). A {@code null} result represents an unbounded + * upper range, which callers performing prefix scans should treat as "scan to the end". + * + * @param input the byte array to increment + * @return A new copy of the incremented byte array, or {@code null} if incrementing would overflow + */ + public static Bytes incrementWithoutOverflow(final Bytes input) { + try { + return increment(input); + } catch (final IndexOutOfBoundsException e) { + return null; + } + } + /** * A byte array comparator based on lexicographic ordering. */ diff --git a/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java b/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java index 139e3efdd4b53..43bd07d7c8b7f 100644 --- a/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java +++ b/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java @@ -41,6 +41,7 @@ import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -77,6 +78,20 @@ public void testIncrementUpperBoundary() { byte[] input = new byte[]{(byte) 0xFF, (byte) 0xFF, (byte) 0xFF}; assertThrows(IndexOutOfBoundsException.class, () -> ByteUtils.increment(Bytes.wrap(input))); } + + @Test + public void testIncrementWithoutOverflow() { + byte[] input = new byte[]{(byte) 0xAB, (byte) 0xCD, (byte) 0xFF}; + byte[] expected = new byte[]{(byte) 0xAB, (byte) 0xCE, (byte) 0x00}; + Bytes output = ByteUtils.incrementWithoutOverflow(Bytes.wrap(input)); + assertArrayEquals(expected, output.get()); + } + + @Test + public void testIncrementWithoutOverflowReturnsNullOnOverflow() { + byte[] input = new byte[]{(byte) 0xFF, (byte) 0xFF, (byte) 0xFF}; + assertNull(ByteUtils.incrementWithoutOverflow(Bytes.wrap(input))); + } @Test public void testIncrementWithSubmap() { final NavigableMap map = new TreeMap<>(); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java index b6d855900911a..4ca3135ed87a1 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java @@ -439,7 +439,9 @@ private , P> KeyValueIterator prefixScan validateStoreOpen(); final KeyValueIterator storeIterator = underlying.prefixScan(prefix, prefixKeySerializer); final Bytes from = Bytes.wrap(prefixKeySerializer.serialize(null, prefix)); - final Bytes to = ByteUtils.increment(from); + // A null upper bound means the prefix has no lexicographic successor (e.g. it is all 0xFF bytes), + // so the scan is unbounded and must run to the end of the keyspace. + final Bytes to = ByteUtils.incrementWithoutOverflow(from); final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = internalContext.cache().range(cacheName, from, to, false); return new MergedSortedCacheKeyValueBytesStoreIterator(cacheIterator, storeIterator, true); } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java index c41aaf4a357ff..55b1de554d7d1 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java @@ -35,7 +35,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; -import static org.apache.kafka.streams.state.internals.RocksDBStore.incrementWithoutOverflow; +import static org.apache.kafka.common.utils.internals.ByteUtils.incrementWithoutOverflow; /** * A generic implementation of {@link RocksDBStore.ColumnFamilyAccessor} that supports dual-column-family diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java index fcf570a101a49..59dfd3990e8fc 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java @@ -239,13 +239,18 @@ public , P> KeyValueIterator prefixScan( private , P> KeyValueIterator prefixScan(final P prefix, final PS prefixKeySerializer, final IsolationLevel isolationLevel) { final Bytes from = Bytes.wrap(prefixKeySerializer.serialize(null, prefix)); - final Bytes to = ByteUtils.increment(from); + // A null upper bound means the prefix has no lexicographic successor (e.g. it is all 0xFF bytes), + // so the scan is unbounded and must run to the end of the keyspace. + final Bytes to = ByteUtils.incrementWithoutOverflow(from); if (transactionBuffer != null) { return transactionBuffer.range(from, to, true, false, isolationLevel); } synchronized (this) { - return new InMemoryKeyValueIterator(map.subMap(from, true, to, false).keySet(), true); + final NavigableMap subMap = to == null + ? map.tailMap(from, true) + : map.subMap(from, true, to, false); + return new InMemoryKeyValueIterator(subMap.keySet(), true); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/LogicalKeyValueSegment.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/LogicalKeyValueSegment.java index 2aa5182fc2dec..0aaea383295c2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/LogicalKeyValueSegment.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/LogicalKeyValueSegment.java @@ -45,7 +45,7 @@ import java.util.function.Function; import java.util.stream.Collectors; -import static org.apache.kafka.streams.state.internals.RocksDBStore.incrementWithoutOverflow; +import static org.apache.kafka.common.utils.internals.ByteUtils.incrementWithoutOverflow; /** * This "logical segment" is a segment which shares its underlying physical store with other diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java index 0406f217aecae..171dc891e6b46 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java @@ -31,6 +31,7 @@ import java.util.Iterator; import java.util.Map; +import java.util.NavigableMap; import java.util.Objects; import java.util.TreeMap; @@ -90,13 +91,18 @@ private Iterator getIterator(final TreeMap treeMap, final public , P> KeyValueIterator prefixScan(final P prefix, final PS prefixKeySerializer) { final Bytes from = Bytes.wrap(prefixKeySerializer.serialize(null, prefix)); - final Bytes to = ByteUtils.increment(from); + // A null upper bound means the prefix has no lexicographic successor (e.g. it is all 0xFF bytes), + // so the scan is unbounded and must run to the end of the keyspace. + final Bytes to = ByteUtils.incrementWithoutOverflow(from); final TreeMap treeMap = toTreeMap(); + final NavigableMap subMap = to == null + ? treeMap.tailMap(from, true) + : treeMap.subMap(from, true, to, false); return new DelegatingPeekingKeyValueIterator<>( name(), - new MemoryNavigableLRUCache.CacheIterator(treeMap.subMap(from, true, to, false).keySet().iterator(), treeMap) + new MemoryNavigableLRUCache.CacheIterator(subMap.keySet().iterator(), treeMap) ); } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java index e7b2eef3a4cdb..9671afba57c44 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java @@ -1478,7 +1478,7 @@ public ManagedKeyValueIterator all(final DBAccessor accessor, fin @Override public ManagedKeyValueIterator prefixScan(final DBAccessor accessor, final Bytes prefix) { - final Bytes to = incrementWithoutOverflow(prefix); + final Bytes to = ByteUtils.incrementWithoutOverflow(prefix); return accessor.prefixScan(columnFamily, name, prefix, to); } @@ -1552,19 +1552,4 @@ public Position getPosition() { } } - /** - * Same as {@link ByteUtils#increment(Bytes)} but {@code null} is returned instead of throwing - * {@code IndexOutOfBoundsException} in the event of overflow. - * - * @param input bytes to increment - * @return A new copy of the incremented byte array, or {@code null} if incrementing would - * result in overflow. - */ - static Bytes incrementWithoutOverflow(final Bytes input) { - try { - return ByteUtils.increment(input); - } catch (final IndexOutOfBoundsException e) { - return null; - } - } } \ No newline at end of file diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractKeyValueStoreTest.java index 3727add492d13..ab0e99235d8ab 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractKeyValueStoreTest.java @@ -660,5 +660,27 @@ public void prefixScanShouldNotThrowConcurrentModificationException() { iter.next(); } } - } + } + + @Test + public void prefixScanShouldReturnMatchesWhenPrefixHasNoUpperBound() { + // A prefix that serializes to all 0xFF bytes has no lexicographically larger successor, + // so incrementing it would overflow. The scan must treat this as an unbounded upper range + // and run to the end of the keyspace rather than throwing IndexOutOfBoundsException. + // IntegerSerializer encodes -1 as 0xFFFFFFFF, which reproduces the overflow case (KAFKA-20597). + store.put(0, "zero"); + store.put(1, "one"); + store.put(2, "two"); + store.put(-1, "minus-one"); + + final List> result = new ArrayList<>(); + try (final KeyValueIterator iter = store.prefixScan(-1, new IntegerSerializer())) { + while (iter.hasNext()) { + result.add(iter.next()); + } + } + + // Only the key whose bytes are 0xFFFFFFFF matches the prefix; the lower-valued keys are excluded. + assertEquals(Collections.singletonList(KeyValue.pair(-1, "minus-one")), result); + } } From 2dc3288627b57f4f55a31a18963e3e0e58f35e84 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=A8=B1=E5=85=83=E8=B1=AA?= <146086744+edenfunf@users.noreply.github.com> Date: Thu, 9 Jul 2026 17:51:23 +0800 Subject: [PATCH 2/2] Clean up redundant comments and formatting --- .../org/apache/kafka/common/utils/internals/ByteUtilsTest.java | 1 + .../kafka/streams/state/internals/CachingKeyValueStore.java | 2 -- .../kafka/streams/state/internals/InMemoryKeyValueStore.java | 2 -- .../kafka/streams/state/internals/MemoryNavigableLRUCache.java | 2 -- .../org/apache/kafka/streams/state/internals/RocksDBStore.java | 3 +-- 5 files changed, 2 insertions(+), 8 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java b/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java index 43bd07d7c8b7f..051ededbb040c 100644 --- a/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java +++ b/clients/src/test/java/org/apache/kafka/common/utils/internals/ByteUtilsTest.java @@ -92,6 +92,7 @@ public void testIncrementWithoutOverflowReturnsNullOnOverflow() { byte[] input = new byte[]{(byte) 0xFF, (byte) 0xFF, (byte) 0xFF}; assertNull(ByteUtils.incrementWithoutOverflow(Bytes.wrap(input))); } + @Test public void testIncrementWithSubmap() { final NavigableMap map = new TreeMap<>(); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java index 4ca3135ed87a1..e342849f14667 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java @@ -439,8 +439,6 @@ private , P> KeyValueIterator prefixScan validateStoreOpen(); final KeyValueIterator storeIterator = underlying.prefixScan(prefix, prefixKeySerializer); final Bytes from = Bytes.wrap(prefixKeySerializer.serialize(null, prefix)); - // A null upper bound means the prefix has no lexicographic successor (e.g. it is all 0xFF bytes), - // so the scan is unbounded and must run to the end of the keyspace. final Bytes to = ByteUtils.incrementWithoutOverflow(from); final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = internalContext.cache().range(cacheName, from, to, false); return new MergedSortedCacheKeyValueBytesStoreIterator(cacheIterator, storeIterator, true); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java index 59dfd3990e8fc..c676c1ab39d7b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryKeyValueStore.java @@ -239,8 +239,6 @@ public , P> KeyValueIterator prefixScan( private , P> KeyValueIterator prefixScan(final P prefix, final PS prefixKeySerializer, final IsolationLevel isolationLevel) { final Bytes from = Bytes.wrap(prefixKeySerializer.serialize(null, prefix)); - // A null upper bound means the prefix has no lexicographic successor (e.g. it is all 0xFF bytes), - // so the scan is unbounded and must run to the end of the keyspace. final Bytes to = ByteUtils.incrementWithoutOverflow(from); if (transactionBuffer != null) { diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java index 171dc891e6b46..b57846fc06811 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MemoryNavigableLRUCache.java @@ -91,8 +91,6 @@ private Iterator getIterator(final TreeMap treeMap, final public , P> KeyValueIterator prefixScan(final P prefix, final PS prefixKeySerializer) { final Bytes from = Bytes.wrap(prefixKeySerializer.serialize(null, prefix)); - // A null upper bound means the prefix has no lexicographic successor (e.g. it is all 0xFF bytes), - // so the scan is unbounded and must run to the end of the keyspace. final Bytes to = ByteUtils.incrementWithoutOverflow(from); final TreeMap treeMap = toTreeMap(); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java index 9671afba57c44..dfb4775af1f2e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java @@ -1551,5 +1551,4 @@ public Position getPosition() { return position.copy().merge(dbAccessor.uncommittedPositionDeltas()); } } - -} \ No newline at end of file +}