From f5769ebb498c9fd119ea000376e4bc4c150cdba2 Mon Sep 17 00:00:00 2001 From: Aleksandr Savonin Date: Mon, 14 Sep 2026 23:16:07 +0200 Subject: [PATCH 1/3] [FLINK-40584][state/rocksdb] Check MapState iterator validity before advancing Cache refill can seek past the last entry after MapState.remove() deletes the resume key. That removal does not update the cached deletion flag, so advancing the exhausted iterator violates the RocksDB API contract. --- .../flink/state/rocksdb/RocksDBMapState.java | 14 ++- .../state/rocksdb/RocksIteratorWrapper.java | 1 + .../state/rocksdb/RocksDBMapStateTest.java | 103 ++++++++++++++++++ 3 files changed, 114 insertions(+), 4 deletions(-) create mode 100644 flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java index c811fc8caac2c..6282000afad8e 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java @@ -18,6 +18,7 @@ package org.apache.flink.state.rocksdb; +import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.api.common.state.MapState; import org.apache.flink.api.common.state.State; import org.apache.flink.api.common.state.StateDescriptor; @@ -70,6 +71,9 @@ class RocksDBMapState extends AbstractRocksDBState userKeySerializer; @@ -562,8 +566,6 @@ public boolean equals(Object o) { /** An auxiliary utility to scan all entries under the given key. */ private abstract class RocksDBMapIterator implements Iterator { - private static final int CACHE_SIZE_LIMIT = 128; - /** The db where data resides. */ private final RocksDB db; @@ -675,8 +677,12 @@ private void loadCache() { /* * If the entry pointing to the current position is not removed, it will be the first entry in the * new iterating. Skip it to avoid redundant access in such cases. + * + * Removing the current entry through MapState does not update its cached 'deleted' flag. + * A resumed seek can therefore return an invalid iterator even if 'deleted' is false. + * RocksDB requires a valid iterator before calling next(). */ - if (currentEntry != null && !currentEntry.deleted) { + if (currentEntry != null && !currentEntry.deleted && iterator.isValid()) { iterator.next(); } @@ -687,7 +693,7 @@ private void loadCache() { break; } - if (cacheEntries.size() >= CACHE_SIZE_LIMIT) { + if (cacheEntries.size() >= ITERATOR_CACHE_SIZE) { break; } diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java index 3422bf7ca35a8..c03a82ef50b0a 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java @@ -94,6 +94,7 @@ public void seekForPrev(ByteBuffer target) { @Override public void next() { + assert isValid() : "Iterator must be valid before calling next()"; iterator.next(); } diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java new file mode 100644 index 0000000000000..fb62664af42f8 --- /dev/null +++ b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java @@ -0,0 +1,103 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.state.rocksdb; + +import org.apache.flink.api.common.state.MapState; +import org.apache.flink.api.common.state.MapStateDescriptor; +import org.apache.flink.api.common.typeutils.base.IntSerializer; +import org.apache.flink.runtime.state.KeyGroupRange; +import org.apache.flink.util.IOUtils; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Path; +import java.util.Collections; +import java.util.Iterator; +import java.util.Map; + +import static org.apache.flink.state.rocksdb.RocksDBMapState.ITERATOR_CACHE_SIZE; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests iteration of RocksDB map state. */ +class RocksDBMapStateTest { + + private static final MapStateDescriptor STATE_DESCRIPTOR = + new MapStateDescriptor<>("map", IntSerializer.INSTANCE, IntSerializer.INSTANCE); + + @TempDir private Path temporaryDirectory; + + private RocksDBKeyedStateBackend backend; + + @BeforeEach + void setUp() throws Exception { + backend = + RocksDBTestUtils.builderForTestDefaults( + temporaryDirectory.toFile(), + IntSerializer.INSTANCE, + 1, + new KeyGroupRange(0, 0), + Collections.emptyList()) + .build(); + } + + @AfterEach + void tearDown() throws Exception { + if (backend != null) { + IOUtils.closeQuietly(backend); + backend.dispose(); + } + } + + /** Isolates the validity guard without requiring prefix bounds or neighboring tombstones. */ + @Test + void testIteratorCacheReloadDoesNotAdvanceAfterExhaustedSeek() throws Exception { + assertRocksIteratorAssertionsEnabled(); + final MapState state = mapState(0, 0); + + // One entry beyond the cache ensures that exhaustion requires another seek. + for (int i = 0; i <= ITERATOR_CACHE_SIZE; i++) { + state.put(i, i); + } + final Iterator> iterator = state.iterator(); + for (int i = 0; i < ITERATOR_CACHE_SIZE; i++) { + assertThat(iterator.next().getKey()).isEqualTo(i); + } + + // MapState.remove() leaves the cached entry's 'deleted' flag false. Remove both the + // resume entry and the only later entry so the refill seek is exhausted. + state.remove(ITERATOR_CACHE_SIZE - 1); + state.remove(ITERATOR_CACHE_SIZE); + + assertThat(iterator).isExhausted(); + } + + private static void assertRocksIteratorAssertionsEnabled() { + assertThat(RocksIteratorWrapper.class.desiredAssertionStatus()) + .as("This test requires Java assertions for RocksIteratorWrapper") + .isTrue(); + } + + private MapState mapState(int key, int namespace) throws Exception { + backend.setCurrentKey(key); + return backend.getPartitionedState(namespace, IntSerializer.INSTANCE, STATE_DESCRIPTOR); + } +} From 9992c8e1cc40cd071eb08bc079b5fcc99a3e7ff0 Mon Sep 17 00:00:00 2001 From: Aleksandr Savonin Date: Tue, 15 Sep 2026 19:11:59 +0200 Subject: [PATCH 2/3] [FLINK-40584][state/rocksdb] Bound MapState and timer prefix iterators Stop native MapState and timer seeks at the end of their key/namespace or key-group prefix instead of traversing unrelated tombstones during state access and timer queue initialization. Generated-by: Claude Fable 5.1 --- .../RocksDBCachingPriorityQueueSet.java | 3 +- .../flink/state/rocksdb/RocksDBMapState.java | 17 +- .../state/rocksdb/RocksDBOperationUtils.java | 71 ++++ .../state/rocksdb/RocksIteratorWrapper.java | 24 +- ...onedPriorityQueueWithRocksDBStoreTest.java | 114 +++++- .../state/rocksdb/RocksDBMapStateTest.java | 160 +++++++- .../rocksdb/RocksDBPrefixIteratorTest.java | 385 ++++++++++++++++++ 7 files changed, 744 insertions(+), 30 deletions(-) create mode 100644 flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBCachingPriorityQueueSet.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBCachingPriorityQueueSet.java index 2bfcb4d3111f4..7ec6820403d0c 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBCachingPriorityQueueSet.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBCachingPriorityQueueSet.java @@ -290,7 +290,8 @@ public int size() { private RocksBytesIterator orderedBytesIterator() { flushWriteBatch(); return new RocksBytesIterator( - new RocksIteratorWrapper(db.newIterator(columnFamilyHandle, readOptions))); + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamilyHandle, readOptions, groupPrefixBytes)); } /** Ensures that recent writes are flushed and reflect in the RocksDB instance. */ diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java index 6282000afad8e..8f33bfe6aad03 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBMapState.java @@ -282,8 +282,8 @@ public boolean isEmpty() { final byte[] prefixBytes = serializeCurrentKeyWithGroupAndNamespace(); try (RocksIteratorWrapper iterator = - RocksDBOperationUtils.getRocksIterator( - backend.db, columnFamily, backend.getReadOptions())) { + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + backend.db, columnFamily, backend.getReadOptions(), prefixBytes)) { iterator.seek(prefixBytes); @@ -293,16 +293,19 @@ public boolean isEmpty() { @Override public void clear() { + final byte[] keyPrefixBytes = serializeCurrentKeyWithGroupAndNamespace(); try (RocksIteratorWrapper iterator = - RocksDBOperationUtils.getRocksIterator( - backend.db, columnFamily, backend.getReadOptions()); + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + backend.db, + columnFamily, + backend.getReadOptions(), + keyPrefixBytes); RocksDBWriteBatchWrapper rocksDBWriteBatchWrapper = new RocksDBWriteBatchWrapper( backend.db, backend.getWriteOptions(), backend.getWriteBatchSize())) { - final byte[] keyPrefixBytes = serializeCurrentKeyWithGroupAndNamespace(); iterator.seek(keyPrefixBytes); while (iterator.isValid()) { @@ -658,8 +661,8 @@ private void loadCache() { // exception // occurred in the below code block. try (RocksIteratorWrapper iterator = - RocksDBOperationUtils.getRocksIterator( - db, columnFamily, backend.getReadOptions())) { + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, backend.getReadOptions(), keyPrefixBytes)) { /* * The iteration starts from the prefix bytes at the first loading. After #nextEntry() is called, diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java index 0f47d9dafacfa..d340464122744 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java @@ -17,6 +17,7 @@ package org.apache.flink.state.rocksdb; +import org.apache.flink.annotation.VisibleForTesting; import org.apache.flink.configuration.ConfigConstants; import org.apache.flink.core.fs.ICloseableRegistry; import org.apache.flink.runtime.execution.CancelTaskException; @@ -38,6 +39,7 @@ import org.rocksdb.ReadOptions; import org.rocksdb.RocksDB; import org.rocksdb.RocksDBException; +import org.rocksdb.Slice; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -113,6 +115,75 @@ public static RocksIteratorWrapper getRocksIterator( return new RocksIteratorWrapper(db.newIterator(columnFamilyHandle, readOptions)); } + /** + * Creates an iterator that stops at the end of the key range starting with {@code prefix}. + * + *

Without the bound, a seek that finds no live key under the prefix keeps reading into + * neighboring key ranges, possibly through many tombstones, before it reports the end. For a + * map with prefix {@code [0x12, 0x34]}, the bound {@code [0x12, 0x35]} stops the iterator + * before the next map's keys. The prefix only sets the bound; callers still position the + * iterator with {@code seek}. + * + *

The iterator uses a private copy of {@code readOptions} with the upper bound set and + * auto-prefix mode enabled. The caller's options are not modified. If the prefix has no + * exclusive upper bound (see {@code getPrefixEnd}), the caller's options are used as they are + * and the iterator is unbounded: no key outside such a prefix can follow it, so nothing extra + * is visited. + */ + static RocksIteratorWrapper getRocksIteratorBoundedByPrefix( + RocksDB db, + ColumnFamilyHandle columnFamilyHandle, + ReadOptions readOptions, + byte[] prefix) { + final byte[] prefixEnd = getPrefixEnd(prefix); + if (prefixEnd == null) { + return getRocksIterator(db, columnFamilyHandle, readOptions); + } + + final Slice upperBound = new Slice(prefixEnd); + ReadOptions boundedReadOptions = null; + try { + boundedReadOptions = new ReadOptions(readOptions); + boundedReadOptions.setIterateUpperBound(upperBound); + // Auto-prefix mode is needed because, with a prefix extractor configured, the default + // seek mode honors the upper bound only if it shares the seek key's extracted prefix. + boundedReadOptions.setAutoPrefixMode(true); + return new RocksIteratorWrapper( + db.newIterator(columnFamilyHandle, boundedReadOptions), + boundedReadOptions, + upperBound); + } catch (RuntimeException | Error e) { + IOUtils.closeQuietly(boundedReadOptions); + IOUtils.closeQuietly(upperBound); + throw e; + } + } + + /** + * Returns the smallest key that is greater than every key starting with {@code prefix}, or + * {@code null} if there is none. + * + *

Keys are ordered as unsigned bytes. Incrementing the last byte of the prefix moves past + * every key that starts with it. A {@code 0xFF} byte cannot be incremented, so trailing {@code + * 0xFF} bytes are dropped and the byte before them is incremented instead. For example, {@code + * [0x12, 0xFE, 0xFF]} becomes {@code [0x12, 0xFF]}. + * + *

No such key exists for an empty prefix or for a prefix made only of {@code 0xFF} bytes. + * Every key at or after such a prefix still starts with it. + */ + @Nullable + @VisibleForTesting + static byte[] getPrefixEnd(byte[] prefix) { + for (int i = prefix.length - 1; i >= 0; --i) { + if (prefix[i] != (byte) 0xFF) { + final byte[] prefixEnd = Arrays.copyOf(prefix, i + 1); + ++prefixEnd[i]; + return prefixEnd; + } + } + return null; + } + public static void registerKvStateInformation( Map kvStateInformation, RocksDBNativeMetricMonitor nativeMetricMonitor, diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java index c03a82ef50b0a..2131c80670646 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksIteratorWrapper.java @@ -20,11 +20,14 @@ import org.apache.flink.util.FlinkRuntimeException; +import org.rocksdb.ReadOptions; import org.rocksdb.RocksDBException; import org.rocksdb.RocksIterator; import org.rocksdb.RocksIteratorInterface; +import org.rocksdb.Slice; import javax.annotation.Nonnull; +import javax.annotation.Nullable; import java.io.Closeable; import java.nio.ByteBuffer; @@ -48,8 +51,21 @@ public class RocksIteratorWrapper implements RocksIteratorInterface, Closeable { private RocksIterator iterator; + @Nullable private final ReadOptions ownedReadOptions; + + @Nullable private final Slice ownedUpperBound; + public RocksIteratorWrapper(@Nonnull RocksIterator iterator) { + this(iterator, null, null); + } + + RocksIteratorWrapper( + @Nonnull RocksIterator iterator, + @Nullable ReadOptions ownedReadOptions, + @Nullable Slice ownedUpperBound) { this.iterator = iterator; + this.ownedReadOptions = ownedReadOptions; + this.ownedUpperBound = ownedUpperBound; } @Override @@ -128,6 +144,12 @@ public byte[] value() { @Override public void close() { - iterator.close(); + // The native iterator keeps a pointer to the upper-bound slice, so the iterator is closed + // first and the owned read options and slice afterwards. RocksObject#close() disposes a + // native handle only once, which makes closing this wrapper twice safe. + try (ownedUpperBound; + ownedReadOptions) { + iterator.close(); + } } } diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/KeyGroupPartitionedPriorityQueueWithRocksDBStoreTest.java b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/KeyGroupPartitionedPriorityQueueWithRocksDBStoreTest.java index 309b81faeaaf8..0b6b4ed2cfb75 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/KeyGroupPartitionedPriorityQueueWithRocksDBStoreTest.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/KeyGroupPartitionedPriorityQueueWithRocksDBStoreTest.java @@ -23,9 +23,17 @@ import org.apache.flink.runtime.state.CompositeKeySerializationUtils; import org.apache.flink.runtime.state.InternalPriorityQueue; import org.apache.flink.runtime.state.InternalPriorityQueueTestBase; +import org.apache.flink.runtime.state.KeyGroupRange; import org.apache.flink.runtime.state.heap.KeyGroupPartitionedPriorityQueue; +import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; +import org.rocksdb.FlushOptions; +import org.rocksdb.MutableColumnFamilyOptions; +import org.rocksdb.RocksDB; +import org.rocksdb.TableProperties; + +import static org.assertj.core.api.Assertions.assertThat; /** * Test of {@link KeyGroupPartitionedPriorityQueue} powered by a {@link @@ -35,6 +43,74 @@ class KeyGroupPartitionedPriorityQueueWithRocksDBStoreTest extends InternalPrior @RegisterExtension public final RocksDBExtension rocksDBExtension = new RocksDBExtension(); + @Test + void testQueueInitializationDoesNotScanTombstonesOutsideKeyGroupRange() throws Exception { + // Deleted entries of the key group right after the range, as left behind by rescaling. + createTombstonesInKeyGroup(128, 512); + rocksDBExtension.getReadOptions().setMaxSkippableInternalKeys(8); + final InternalPriorityQueue queue = + new KeyGroupPartitionedPriorityQueue<>( + KEY_EXTRACTOR_FUNCTION, + TEST_ELEMENT_PRIORITY_COMPARATOR, + newFactory(), + new KeyGroupRange(0, 127), + 512); + + assertThat(queue.isEmpty()).isTrue(); + } + + @Test + void testCacheRefillKeepsBoundsOfKeyGroup() throws Exception { + final RocksDBCachingPriorityQueueSet queue = + newPriorityQueueForKeyGroup(255, 512, 1); + final TestElement first = new TestElement(1, 1); + final TestElement second = new TestElement(2, 2); + queue.add(first); + queue.add(second); + assertThat(queue.poll()).isEqualTo(first); + createTombstonesInKeyGroup(256, 512); + + rocksDBExtension.getReadOptions().setMaxSkippableInternalKeys(8); + assertThat(queue.poll()).isEqualTo(second); + assertThat(queue.poll()).isNull(); + } + + private void createTombstonesInKeyGroup(int keyGroupId, int numKeyGroups) throws Exception { + final RocksDB db = rocksDBExtension.getRocksDB(); + db.setOptions( + rocksDBExtension.getDefaultColumnFamily(), + MutableColumnFamilyOptions.builder().setDisableAutoCompactions(true).build()); + + final byte[][] keys = new byte[32][]; + final DataOutputSerializer output = new DataOutputSerializer(32); + for (int i = 0; i < keys.length; i++) { + output.clear(); + CompositeKeySerializationUtils.writeKeyGroup( + keyGroupId, + CompositeKeySerializationUtils.computeRequiredBytesInKeyGroupPrefix( + numKeyGroups), + output); + TestElementSerializer.INSTANCE.serialize(new TestElement(i, i), output); + keys[i] = output.getCopyOfBuffer(); + db.put(rocksDBExtension.getDefaultColumnFamily(), keys[i], new byte[0]); + } + + try (FlushOptions flushOptions = new FlushOptions().setWaitForFlush(true)) { + rocksDBExtension.getBatchWrapper().flush(); + db.flush(flushOptions); + for (byte[] key : keys) { + db.delete(rocksDBExtension.getDefaultColumnFamily(), key); + } + db.flush(flushOptions); + } + + assertThat( + db.getPropertiesOfAllTables().values().stream() + .mapToLong(TableProperties::getNumDeletions) + .sum()) + .isGreaterThanOrEqualTo(keys.length); + } + @Override protected InternalPriorityQueue newPriorityQueue(int initialCapacity) { return new KeyGroupPartitionedPriorityQueue<>( @@ -54,24 +130,24 @@ protected boolean testSetSemanticsAgainstDuplicateElements() { TestElement, RocksDBCachingPriorityQueueSet> newFactory() { - return (keyGroupId, numKeyGroups, keyExtractorFunction, elementComparator) -> { - DataOutputSerializer outputStreamWithPos = new DataOutputSerializer(128); - DataInputDeserializer inputStreamWithPos = new DataInputDeserializer(); - int keyGroupPrefixBytes = - CompositeKeySerializationUtils.computeRequiredBytesInKeyGroupPrefix( - numKeyGroups); - TreeOrderedSetCache orderedSetCache = new TreeOrderedSetCache(32); - return new RocksDBCachingPriorityQueueSet<>( - keyGroupId, - keyGroupPrefixBytes, - rocksDBExtension.getRocksDB(), - rocksDBExtension.getReadOptions(), - rocksDBExtension.getDefaultColumnFamily(), - TestElementSerializer.INSTANCE, - outputStreamWithPos, - inputStreamWithPos, - rocksDBExtension.getBatchWrapper(), - orderedSetCache); - }; + return (keyGroupId, numKeyGroups, keyExtractorFunction, elementComparator) -> + newPriorityQueueForKeyGroup(keyGroupId, numKeyGroups, 32); + } + + private RocksDBCachingPriorityQueueSet newPriorityQueueForKeyGroup( + int keyGroupId, int numKeyGroups, int cacheSize) { + final int keyGroupPrefixBytes = + CompositeKeySerializationUtils.computeRequiredBytesInKeyGroupPrefix(numKeyGroups); + return new RocksDBCachingPriorityQueueSet<>( + keyGroupId, + keyGroupPrefixBytes, + rocksDBExtension.getRocksDB(), + rocksDBExtension.getReadOptions(), + rocksDBExtension.getDefaultColumnFamily(), + TestElementSerializer.INSTANCE, + new DataOutputSerializer(128), + new DataInputDeserializer(), + rocksDBExtension.getBatchWrapper(), + new TreeOrderedSetCache(cacheSize)); } } diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java index fb62664af42f8..2ee46e73b0dec 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBMapStateTest.java @@ -28,24 +28,45 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.ValueSource; +import org.rocksdb.ColumnFamilyHandle; +import org.rocksdb.FlushOptions; +import org.rocksdb.MutableColumnFamilyOptions; +import org.rocksdb.TableProperties; import java.nio.file.Path; +import java.util.ArrayList; import java.util.Collections; import java.util.Iterator; +import java.util.List; import java.util.Map; +import java.util.stream.Collectors; +import java.util.stream.IntStream; import static org.apache.flink.state.rocksdb.RocksDBMapState.ITERATOR_CACHE_SIZE; import static org.assertj.core.api.Assertions.assertThat; -/** Tests iteration of RocksDB map state. */ +/** Tests iteration and prefix isolation of RocksDB map state operations. */ class RocksDBMapStateTest { + // Two cache loads plus one entry: a fully refilled cache followed by a partial refill. + private static final int MAP_ENTRIES = 2 * ITERATOR_CACHE_SIZE + 1; + + // A deleted entry can leave both its value and a tombstone. + // Allow both records per entry, plus one extra skip. + private static final int SKIPPABLE_INTERNAL_KEY_LIMIT = 2 * MAP_ENTRIES + 1; + + // An unbounded seek into the neighboring map must exceed the skip limit. + private static final int NEIGHBOR_TOMBSTONES = SKIPPABLE_INTERNAL_KEY_LIMIT + 1; private static final MapStateDescriptor STATE_DESCRIPTOR = new MapStateDescriptor<>("map", IntSerializer.INSTANCE, IntSerializer.INSTANCE); @TempDir private Path temporaryDirectory; private RocksDBKeyedStateBackend backend; + private ColumnFamilyHandle columnFamily; @BeforeEach void setUp() throws Exception { @@ -57,16 +78,85 @@ void setUp() throws Exception { new KeyGroupRange(0, 0), Collections.emptyList()) .build(); + mapState(0, 0); + columnFamily = backend.getColumnFamilyHandle(STATE_DESCRIPTOR.getName()); + backend.db.setOptions( + columnFamily, + MutableColumnFamilyOptions.builder().setDisableAutoCompactions(true).build()); + // Allow the entire map to be deleted, but fail if a seek scans the neighbor's tombstones. + backend.getReadOptions().setMaxSkippableInternalKeys(SKIPPABLE_INTERNAL_KEY_LIMIT); } @AfterEach - void tearDown() throws Exception { + void tearDown() { if (backend != null) { IOUtils.closeQuietly(backend); backend.dispose(); } } + @ParameterizedTest + @EnumSource(EmptyMapOperation.class) + void testEmptyMapDoesNotScanDeletedEntriesOfNextKey(EmptyMapOperation operation) + throws Exception { + createNeighborTombstones(); + final MapState state = mapState(0, 0); + + switch (operation) { + case ENTRIES: + assertThat(state.entries()).isEmpty(); + break; + case IS_EMPTY: + assertThat(state.isEmpty()).isTrue(); + break; + case CLEAR: + state.clear(); + break; + default: + throw new IllegalArgumentException("Unknown map operation: " + operation); + } + + assertNeighborPreserved(); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testIteratorCacheReloadDoesNotScanDeletedEntriesOfNextKey(boolean removeEntries) + throws Exception { + createNeighborTombstones(); + final MapState state = mapState(0, 0); + populateMap(state); + + final List keys = new ArrayList<>(); + final Iterator> iterator = state.iterator(); + while (iterator.hasNext()) { + final Map.Entry entry = iterator.next(); + keys.add(entry.getKey()); + assertThat(entry.getValue()).isEqualTo(entry.getKey()); + if (removeEntries) { + iterator.remove(); + } + } + + assertThat(keys) + .containsExactlyElementsOf( + IntStream.range(0, MAP_ENTRIES).boxed().collect(Collectors.toList())); + assertThat(state.isEmpty()).isEqualTo(removeEntries); + assertNeighborPreserved(); + } + + @Test + void testClearDoesNotScanDeletedEntriesOfNextKey() throws Exception { + createNeighborTombstones(); + final MapState state = mapState(0, 0); + populateMap(state); + + state.clear(); + + assertThat(state.isEmpty()).isTrue(); + assertNeighborPreserved(); + } + /** Isolates the validity guard without requiring prefix bounds or neighboring tombstones. */ @Test void testIteratorCacheReloadDoesNotAdvanceAfterExhaustedSeek() throws Exception { @@ -90,14 +180,80 @@ void testIteratorCacheReloadDoesNotAdvanceAfterExhaustedSeek() throws Exception assertThat(iterator).isExhausted(); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testIteratorCacheReloadAfterEntriesWereRemoved(boolean removeLastReturnedEntry) + throws Exception { + assertRocksIteratorAssertionsEnabled(); + createNeighborTombstones(); + final MapState state = mapState(0, 0); + populateMap(state); + final Iterator> iterator = state.iterator(); + for (int i = 0; i < ITERATOR_CACHE_SIZE; i++) { + assertThat(iterator.next().getKey()).isEqualTo(i); + } + + // Removing the last returned entry exhausts the resumed seek at the prefix bound. + // Keeping it exercises the valid seek followed by an advance to the bound. + final int firstRemovedKey = + removeLastReturnedEntry ? ITERATOR_CACHE_SIZE - 1 : ITERATOR_CACHE_SIZE; + for (int i = firstRemovedKey; i < MAP_ENTRIES; i++) { + state.remove(i); + } + + assertThat(iterator).isExhausted(); + assertNeighborPreserved(); + } + private static void assertRocksIteratorAssertionsEnabled() { assertThat(RocksIteratorWrapper.class.desiredAssertionStatus()) .as("This test requires Java assertions for RocksIteratorWrapper") .isTrue(); } + private void createNeighborTombstones() throws Exception { + final MapState neighbor = mapState(1, 0); + for (int i = 0; i < NEIGHBOR_TOMBSTONES; i++) { + neighbor.put(i, i); + } + flush(); + for (int i = 0; i < NEIGHBOR_TOMBSTONES; i++) { + neighbor.remove(i); + } + flush(); + + assertThat( + backend.db.getPropertiesOfAllTables(columnFamily).values().stream() + .mapToLong(TableProperties::getNumDeletions) + .sum()) + .isEqualTo(NEIGHBOR_TOMBSTONES); + mapState(2, 0).put(NEIGHBOR_TOMBSTONES, NEIGHBOR_TOMBSTONES); + } + + private void populateMap(MapState state) throws Exception { + for (int i = 0; i < MAP_ENTRIES; i++) { + state.put(i, i); + } + } + + private void assertNeighborPreserved() throws Exception { + assertThat(mapState(2, 0).get(NEIGHBOR_TOMBSTONES)).isEqualTo(NEIGHBOR_TOMBSTONES); + } + private MapState mapState(int key, int namespace) throws Exception { backend.setCurrentKey(key); return backend.getPartitionedState(namespace, IntSerializer.INSTANCE, STATE_DESCRIPTOR); } + + private void flush() throws Exception { + try (FlushOptions options = new FlushOptions().setWaitForFlush(true)) { + backend.db.flush(options, columnFamily); + } + } + + private enum EmptyMapOperation { + ENTRIES, + IS_EMPTY, + CLEAR + } } diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java new file mode 100644 index 0000000000000..9552a9adb5329 --- /dev/null +++ b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java @@ -0,0 +1,385 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.state.rocksdb; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.RegisterExtension; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; +import org.junit.jupiter.params.provider.ValueSource; +import org.rocksdb.BlockBasedTableConfig; +import org.rocksdb.BloomFilter; +import org.rocksdb.ColumnFamilyDescriptor; +import org.rocksdb.ColumnFamilyHandle; +import org.rocksdb.ColumnFamilyOptions; +import org.rocksdb.FlushOptions; +import org.rocksdb.ReadOptions; +import org.rocksdb.RocksDB; +import org.rocksdb.RocksDBException; +import org.rocksdb.RocksIterator; +import org.rocksdb.Slice; +import org.rocksdb.Snapshot; +import org.rocksdb.Statistics; +import org.rocksdb.TickerType; + +import java.util.stream.Stream; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for prefix bounds and resource ownership of RocksDB iterators. */ +class RocksDBPrefixIteratorTest { + + @RegisterExtension final RocksDBExtension rocksDBExtension = new RocksDBExtension(true); + + @ParameterizedTest + @MethodSource("prefixEnds") + void testPrefixEnd(byte[] prefix, byte[] expectedEnd) { + assertThat(RocksDBOperationUtils.getPrefixEnd(prefix)).isEqualTo(expectedEnd); + } + + private static Stream prefixEnds() { + return Stream.of( + Arguments.of(bytes(0), bytes(1)), + Arguments.of(bytes(1), bytes(2)), + // The successor of 0x7F is 0x80 under unsigned ordering. + Arguments.of(bytes(127), bytes(128)), + // The incremented byte may itself become 0xFF. + Arguments.of(bytes(254), bytes(255)), + Arguments.of(bytes(0x12, 0xFE, 0xFF), bytes(0x12, 0xFF)), + // Trailing 0xFF bytes are dropped before the increment. + Arguments.of(bytes(1, 255, 255), bytes(2)), + Arguments.of(bytes(1, 255, 0), bytes(1, 255, 1)), + // No finite bound exists for all-0xFF or empty prefixes. + Arguments.of(bytes(255), null), + Arguments.of(bytes(255, 255), null), + Arguments.of(bytes(), null)); + } + + @ParameterizedTest + @MethodSource("prefixRanges") + void testPrefixRange(byte[] prefix, byte[][] keys, byte[][] expectedKeys) + throws RocksDBException { + putKeys(keys); + + try (RocksIteratorWrapper iterator = + createIterator(rocksDBExtension.getReadOptions(), prefix)) { + assertKeys(iterator, prefix, expectedKeys); + } + } + + private static Stream prefixRanges() { + return Stream.of( + Arguments.of( + bytes(1), + new byte[][] {bytes(0), bytes(1), bytes(1, 0), bytes(1, 255), bytes(2)}, + new byte[][] {bytes(1), bytes(1, 0), bytes(1, 255)}), + Arguments.of( + bytes(1, 255, 255), + new byte[][] { + bytes(1, 255, 254), bytes(1, 255, 255), bytes(1, 255, 255, 0), bytes(2) + }, + new byte[][] {bytes(1, 255, 255), bytes(1, 255, 255, 0)}), + Arguments.of( + bytes(127), + new byte[][] {bytes(126), bytes(127), bytes(127, 255), bytes(128)}, + new byte[][] {bytes(127), bytes(127, 255)}), + Arguments.of( + bytes(255, 255), + new byte[][] { + bytes(255, 254), + bytes(255, 255), + bytes(255, 255, 0), + bytes(255, 255, 255) + }, + new byte[][] {bytes(255, 255), bytes(255, 255, 0), bytes(255, 255, 255)}), + Arguments.of( + bytes(), + new byte[][] {bytes(), bytes(0), bytes(127), bytes(128), bytes(255)}, + new byte[][] {bytes(), bytes(0), bytes(127), bytes(128), bytes(255)})); + } + + @Test + void testConcurrentIteratorsDoNotModifyOrCloseSharedOptions() throws RocksDBException { + putKeys(bytes(1, 0), bytes(2, 0), bytes(3, 0)); + final ReadOptions sharedOptions = rocksDBExtension.getReadOptions(); + sharedOptions.setFillCache(false); + + try (RocksIteratorWrapper second = createIterator(sharedOptions, bytes(2))) { + try (RocksIteratorWrapper first = createIterator(sharedOptions, bytes(1))) { + assertKeys(first, bytes(1), bytes(1, 0)); + assertThat(sharedOptions.iterateUpperBound()).isNull(); + assertThat(sharedOptions.fillCache()).isFalse(); + } + + assertKeys(second, bytes(2), bytes(2, 0)); + } + + assertThat(sharedOptions.isOwningHandle()).isTrue(); + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIterator( + rocksDBExtension.getRocksDB(), + rocksDBExtension.getDefaultColumnFamily(), + sharedOptions)) { + assertKeys(iterator, bytes(), bytes(1, 0), bytes(2, 0), bytes(3, 0)); + } + } + + @Test + void testWrapperClosesOwnedResources() { + final ReadOptions sharedOptions = rocksDBExtension.getReadOptions(); + try (TrackingSlice upperBound = new TrackingSlice(bytes(2)); + ReadOptions ownedOptions = + new ReadOptions(sharedOptions).setIterateUpperBound(upperBound); + RocksIterator nativeIterator = + rocksDBExtension + .getRocksDB() + .newIterator( + rocksDBExtension.getDefaultColumnFamily(), ownedOptions); + RocksIteratorWrapper iterator = + new RocksIteratorWrapper(nativeIterator, ownedOptions, upperBound)) { + iterator.close(); + iterator.close(); + + assertThat(nativeIterator.isOwningHandle()).isFalse(); + assertThat(ownedOptions.isOwningHandle()).isFalse(); + assertThat(upperBound.isClosed()).isTrue(); + assertThat(sharedOptions.isOwningHandle()).isTrue(); + assertThat(sharedOptions.iterateUpperBound()).isNull(); + } + } + + @Test + void testConfiguredSnapshotIsPreserved() throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + db.put(rocksDBExtension.getDefaultColumnFamily(), bytes(1, 0), bytes(10)); + final Snapshot snapshot = db.getSnapshot(); + + try (ReadOptions readOptions = new ReadOptions().setSnapshot(snapshot)) { + db.put(rocksDBExtension.getDefaultColumnFamily(), bytes(1, 0), bytes(20)); + db.put(rocksDBExtension.getDefaultColumnFamily(), bytes(1, 1), bytes(30)); + + try (RocksIteratorWrapper iterator = createIterator(readOptions, bytes(1))) { + iterator.seek(bytes(1)); + assertThat(iterator.isValid()).isTrue(); + assertThat(iterator.key()).isEqualTo(bytes(1, 0)); + assertThat(iterator.value()).isEqualTo(bytes(10)); + iterator.next(); + assertThat(iterator.isValid()).isFalse(); + } + + assertThat(readOptions.isOwningHandle()).isTrue(); + } finally { + db.releaseSnapshot(snapshot); + } + } + + @Test + void testUpperBoundWithPrefixExtractor() throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + try (BloomFilter filter = new BloomFilter(10); + ColumnFamilyOptions options = + new ColumnFamilyOptions() + .useFixedLengthPrefixExtractor(1) + .setTableFormatConfig(prefixBloomTableConfig(filter)); + ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + FlushOptions flushOptions = new FlushOptions().setWaitForFlush(true)) { + db.put(columnFamily, bytes(0, 0), bytes()); + db.put(columnFamily, bytes(1, 0), bytes()); + db.put(columnFamily, bytes(1, 255), bytes()); + db.put(columnFamily, bytes(2, 0), bytes()); + db.flush(flushOptions, columnFamily); + + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, rocksDBExtension.getReadOptions(), bytes(1))) { + assertKeys(iterator, bytes(1), bytes(1, 0), bytes(1, 255)); + } + } + } + + @Test + void testUpperBoundShorterThanPrefixExtractorDomain() throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + try (BloomFilter filter = new BloomFilter(10); + ColumnFamilyOptions options = + new ColumnFamilyOptions() + .useFixedLengthPrefixExtractor(2) + .setTableFormatConfig(prefixBloomTableConfig(filter)); + ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + FlushOptions flushOptions = new FlushOptions().setWaitForFlush(true)) { + db.put(columnFamily, bytes(0), bytes()); + db.put(columnFamily, bytes(0, 254, 0), bytes()); + db.put(columnFamily, bytes(0, 255, 0), bytes()); + db.put(columnFamily, bytes(0, 255, 255), bytes()); + db.put(columnFamily, bytes(1), bytes()); + db.put(columnFamily, bytes(1, 0, 0), bytes()); + db.flush(flushOptions, columnFamily); + + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, rocksDBExtension.getReadOptions(), bytes(0, 255))) { + assertKeys(iterator, bytes(0, 255), bytes(0, 255, 0), bytes(0, 255, 255)); + } + } + } + + @ParameterizedTest + @MethodSource("prefixBloomRanges") + void testPrefixBloomFiltersRespectBoundsAndConfiguredTotalOrderSeek( + byte[] prefix, int extractorLength, boolean filterCompatible, boolean totalOrderSeek) + throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + try (BloomFilter filter = new BloomFilter(10); + ColumnFamilyOptions options = + new ColumnFamilyOptions() + .useFixedLengthPrefixExtractor(extractorLength) + .setTableFormatConfig(prefixBloomTableConfig(filter)); + ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + FlushOptions flushOptions = new FlushOptions().setWaitForFlush(true); + ReadOptions readOptions = new ReadOptions().setTotalOrderSeek(totalOrderSeek); + // Each call owns a separate native shared_ptr wrapper that must be closed. + Statistics statistics = rocksDBExtension.getDbOptions().statistics()) { + db.put(columnFamily, bytes(0, 0, 0), bytes()); + db.put(columnFamily, bytes(2, 0, 0), bytes()); + db.flush(flushOptions, columnFamily); + + final long filterLookupsBefore = filterLookups(statistics); + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, readOptions, prefix)) { + assertKeys(iterator, prefix); + } + + // Filter block accesses prove that the SST prefix Bloom filter was consulted. + assertThat(filterLookups(statistics) - filterLookupsBefore) + .isEqualTo(filterCompatible && !totalOrderSeek ? 1 : 0); + assertThat(readOptions.totalOrderSeek()).isEqualTo(totalOrderSeek); + assertThat(readOptions.autoPrefixMode()).isFalse(); + assertThat(readOptions.iterateUpperBound()).isNull(); + } + } + + private static Stream prefixBloomRanges() { + return Stream.of( + // The bound shares the seek key's extracted prefix. + Arguments.of(bytes(1, 0), 1, true, false), + Arguments.of(bytes(1, 0), 1, true, true), + // The bound is the immediate successor of the extracted prefix. + Arguments.of(bytes(1), 1, true, false), + Arguments.of(bytes(1), 1, true, true), + // The bound is shorter than the extractor, so no Bloom filter can be consulted. + Arguments.of(bytes(0, 255), 2, false, false)); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testPrefixRangeSpansMultipleExtractedPrefixes(boolean capped) throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + try (BloomFilter filter = new BloomFilter(10); + ColumnFamilyOptions options = + new ColumnFamilyOptions() + .setTableFormatConfig(prefixBloomTableConfig(filter))) { + if (capped) { + options.useCappedPrefixExtractor(2); + } else { + options.useFixedLengthPrefixExtractor(2); + } + try (ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + FlushOptions flushOptions = new FlushOptions().setWaitForFlush(true)) { + final byte[][] expectedKeys = + new byte[][] {bytes(1), bytes(1, 0), bytes(1, 1, 0), bytes(1, 255)}; + db.put(columnFamily, bytes(0), bytes()); + for (byte[] key : expectedKeys) { + db.put(columnFamily, key, bytes()); + } + db.put(columnFamily, bytes(2), bytes()); + db.flush(flushOptions, columnFamily); + + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, rocksDBExtension.getReadOptions(), bytes(1))) { + assertKeys(iterator, bytes(1), expectedKeys); + } + } + } + } + + private static BlockBasedTableConfig prefixBloomTableConfig(BloomFilter filter) { + return new BlockBasedTableConfig() + .setFilterPolicy(filter) + .setWholeKeyFiltering(false) + .setCacheIndexAndFilterBlocks(true); + } + + private static long filterLookups(Statistics statistics) { + return statistics.getTickerCount(TickerType.BLOCK_CACHE_FILTER_HIT) + + statistics.getTickerCount(TickerType.BLOCK_CACHE_FILTER_MISS); + } + + private RocksIteratorWrapper createIterator(ReadOptions readOptions, byte[] prefix) { + return RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + rocksDBExtension.getRocksDB(), + rocksDBExtension.getDefaultColumnFamily(), + readOptions, + prefix); + } + + private void putKeys(byte[]... keys) throws RocksDBException { + for (byte[] key : keys) { + rocksDBExtension + .getRocksDB() + .put(rocksDBExtension.getDefaultColumnFamily(), key, bytes()); + } + } + + private static void assertKeys( + RocksIteratorWrapper iterator, byte[] seekKey, byte[]... expectedKeys) { + iterator.seek(seekKey); + for (byte[] expectedKey : expectedKeys) { + assertThat(iterator.isValid()).isTrue(); + assertThat(iterator.key()).isEqualTo(expectedKey); + iterator.next(); + } + assertThat(iterator.isValid()).isFalse(); + } + + private static byte[] bytes(int... values) { + final byte[] bytes = new byte[values.length]; + for (int i = 0; i < values.length; i++) { + bytes[i] = (byte) values[i]; + } + return bytes; + } + + private static final class TrackingSlice extends Slice { + + private TrackingSlice(byte[] data) { + super(data); + } + + private boolean isClosed() { + return !isOwningHandle(); + } + } +} From a4ccbc857f56dd7ce951dc91174f31596a0cd515 Mon Sep 17 00:00:00 2001 From: Aleksandr Savonin Date: Tue, 15 Sep 2026 19:12:16 +0200 Subject: [PATCH 3/3] [FLINK-40584][state/rocksdb] Clear inherited seek restrictions on bounded iterators The bounded iterator copies the backend's ReadOptions, so a prefixSameAsStart or iterate_lower_bound set globally by an options factory would carry over and silently drop keys inside the bound: the first stops at the next extracted prefix, the second clamps the seek. Clear both on the private copy, the shared options stay untouched. --- .../state/rocksdb/RocksDBOperationUtils.java | 18 +++-- .../rocksdb/RocksDBPrefixIteratorTest.java | 79 +++++++++++++++++++ 2 files changed, 92 insertions(+), 5 deletions(-) diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java index d340464122744..8c8e5cdb196d1 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/state/rocksdb/RocksDBOperationUtils.java @@ -124,11 +124,11 @@ public static RocksIteratorWrapper getRocksIterator( * before the next map's keys. The prefix only sets the bound; callers still position the * iterator with {@code seek}. * - *

The iterator uses a private copy of {@code readOptions} with the upper bound set and - * auto-prefix mode enabled. The caller's options are not modified. If the prefix has no - * exclusive upper bound (see {@code getPrefixEnd}), the caller's options are used as they are - * and the iterator is unbounded: no key outside such a prefix can follow it, so nothing extra - * is visited. + *

The iterator uses a private copy of {@code readOptions} with the upper bound set, + * auto-prefix mode enabled, and {@code prefixSameAsStart} and any lower bound cleared. The + * caller's options are not modified. If the prefix has no exclusive upper bound (see {@code + * getPrefixEnd}), the caller's options are used as they are and the iterator is unbounded: no + * key outside such a prefix can follow it, so nothing extra is visited. */ static RocksIteratorWrapper getRocksIteratorBoundedByPrefix( RocksDB db, @@ -148,6 +148,14 @@ static RocksIteratorWrapper getRocksIteratorBoundedByPrefix( // Auto-prefix mode is needed because, with a prefix extractor configured, the default // seek mode honors the upper bound only if it shares the seek key's extracted prefix. boundedReadOptions.setAutoPrefixMode(true); + // prefix_same_as_start would invalidate the iterator as soon as the extracted prefix + // changes. With an extractor longer than the seek prefix that happens inside the map + // or key group, or immediately for the initial seek from the bare prefix. Neither + // auto_prefix_mode nor total_order_seek disables it, so it is cleared on the copy. + boundedReadOptions.setPrefixSameAsStart(false); + // The seek key is the lower end of every bounded scan. A configured lower bound above + // it would silently skip entries, so the copy carries none. + boundedReadOptions.setIterateLowerBound(null); return new RocksIteratorWrapper( db.newIterator(columnFamilyHandle, boundedReadOptions), boundedReadOptions, diff --git a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java index 9552a9adb5329..c920c532473b1 100644 --- a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java +++ b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/state/rocksdb/RocksDBPrefixIteratorTest.java @@ -242,6 +242,85 @@ void testUpperBoundShorterThanPrefixExtractorDomain() throws RocksDBException { } } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testResumedSeekDoesNotStopAtExtractedPrefix(boolean totalOrderSeek) + throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + try (ColumnFamilyOptions options = + new ColumnFamilyOptions().useFixedLengthPrefixExtractor(2); + ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + ReadOptions readOptions = + new ReadOptions() + .setPrefixSameAsStart(true) + .setTotalOrderSeek(totalOrderSeek)) { + db.put(columnFamily, bytes(1, 0), bytes()); + db.put(columnFamily, bytes(1, 1), bytes()); + db.put(columnFamily, bytes(2, 0), bytes()); + + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, readOptions, bytes(1))) { + // The map prefix is [1], but the resume key's extracted prefix is [1, 0]. + assertKeys(iterator, bytes(1, 0), bytes(1, 0), bytes(1, 1)); + } + + assertThat(readOptions.prefixSameAsStart()).isTrue(); + } + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testInitialSeekDoesNotStopAtExtractedPrefix(boolean totalOrderSeek) + throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + // RocksDB enforces prefixSameAsStart regardless of totalOrderSeek, so both are covered. + try (ColumnFamilyOptions options = new ColumnFamilyOptions().useCappedPrefixExtractor(2); + ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + ReadOptions readOptions = + new ReadOptions() + .setPrefixSameAsStart(true) + .setTotalOrderSeek(totalOrderSeek)) { + db.put(columnFamily, bytes(1, 0), bytes()); + db.put(columnFamily, bytes(1, 1), bytes()); + db.put(columnFamily, bytes(2, 0), bytes()); + + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, readOptions, bytes(1))) { + // The bare prefix [1] is its own extracted prefix, which no stored key shares. + assertKeys(iterator, bytes(1), bytes(1, 0), bytes(1, 1)); + } + + assertThat(readOptions.prefixSameAsStart()).isTrue(); + } + } + + @Test + void testConfiguredLowerBoundDoesNotClampSeek() throws RocksDBException { + final RocksDB db = rocksDBExtension.getRocksDB(); + try (ColumnFamilyOptions options = new ColumnFamilyOptions(); + ColumnFamilyHandle columnFamily = + db.createColumnFamily(new ColumnFamilyDescriptor(bytes(42), options)); + Slice lowerBound = new Slice(bytes(1, 1)); + ReadOptions readOptions = new ReadOptions().setIterateLowerBound(lowerBound)) { + db.put(columnFamily, bytes(1, 0), bytes()); + db.put(columnFamily, bytes(1, 1), bytes()); + db.put(columnFamily, bytes(2, 0), bytes()); + + try (RocksIteratorWrapper iterator = + RocksDBOperationUtils.getRocksIteratorBoundedByPrefix( + db, columnFamily, readOptions, bytes(1))) { + // An inherited lower bound would clamp the seek to [1, 1]. + assertKeys(iterator, bytes(1, 0), bytes(1, 0), bytes(1, 1)); + } + + assertThat(readOptions.iterateLowerBound().data()).isEqualTo(bytes(1, 1)); + } + } + @ParameterizedTest @MethodSource("prefixBloomRanges") void testPrefixBloomFiltersRespectBoundsAndConfiguredTotalOrderSeek(