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 88a5222a9bc KAFKA-20777: Make a closed StagedMergeIterator throw
InvalidStateStoreException (#22797)
88a5222a9bc is described below
commit 88a5222a9bcc3eaa4509d3ce33c152dcc0d53101
Author: Alan Lau <[email protected]>
AuthorDate: Mon Jul 13 16:35:25 2026 -0400
KAFKA-20777: Make a closed StagedMergeIterator throw
InvalidStateStoreException (#22797)
Jira: https://issues.apache.org/jira/browse/KAFKA-20777
When transactional state stores are enabled
(enable.transactional.statestores=true), an open StagedMergeIterator
that is used after the store is closed throws
java.lang.IllegalStateException("Iterator has already been closed.")
from StagedMergeIterator.hasNext(), for consistency with the other
managed store iterators (e.g. RocksDbIterator), it should throw
InvalidStateStoreException.
Reviewers: Bill Bejeck <[email protected]>
---
.../streams/state/internals/StagedMergeIterator.java | 3 ++-
.../state/internals/StagedMergeIteratorTest.java | 18 ++++++++++++++++++
2 files changed, 20 insertions(+), 1 deletion(-)
diff --git
a/streams/src/main/java/org/apache/kafka/streams/state/internals/StagedMergeIterator.java
b/streams/src/main/java/org/apache/kafka/streams/state/internals/StagedMergeIterator.java
index 76bd3abc5b2..8477ef6dad9 100644
---
a/streams/src/main/java/org/apache/kafka/streams/state/internals/StagedMergeIterator.java
+++
b/streams/src/main/java/org/apache/kafka/streams/state/internals/StagedMergeIterator.java
@@ -17,6 +17,7 @@
package org.apache.kafka.streams.state.internals;
import org.apache.kafka.streams.KeyValue;
+import org.apache.kafka.streams.errors.InvalidStateStoreException;
import org.apache.kafka.streams.state.KeyValueIterator;
import java.util.Iterator;
@@ -77,7 +78,7 @@ class StagedMergeIterator<K extends Comparable<K>, V>
implements ManagedKeyValue
@Override
public boolean hasNext() {
if (closed) {
- throw new IllegalStateException("Iterator has already been
closed.");
+ throw new InvalidStateStoreException("Iterator has already been
closed.");
}
if (prefetched != null) {
return true;
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/internals/StagedMergeIteratorTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/internals/StagedMergeIteratorTest.java
index 0d99e8cebce..84db1756f8a 100644
---
a/streams/src/test/java/org/apache/kafka/streams/state/internals/StagedMergeIteratorTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/state/internals/StagedMergeIteratorTest.java
@@ -18,6 +18,7 @@ package org.apache.kafka.streams.state.internals;
import org.apache.kafka.common.utils.Bytes;
import org.apache.kafka.streams.KeyValue;
+import org.apache.kafka.streams.errors.InvalidStateStoreException;
import org.apache.kafka.streams.state.KeyValueIterator;
import org.junit.jupiter.api.Test;
@@ -87,6 +88,23 @@ public class StagedMergeIteratorTest {
}
}
+ @Test
+ public void shouldThrowInvalidStateStoreExceptionOnHasNextAfterClose() {
+ final NavigableMap<Bytes, Optional<byte[]>> staging = new TreeMap<>();
+ staging.put(key("a"), Optional.of(val("staged-a")));
+
+ final KeyValueIterator<Bytes, byte[]> iter =
+ new StagedMergeIterator(staging, new ListIterator(List.of()));
+ iter.close();
+
+ // A closed store iterator must signal via InvalidStateStoreException
-- the type
+ // SegmentIterator catches to close open iterators gracefully -- not a
bare
+ // IllegalStateException, which would escape that handling.
+ assertThrows(InvalidStateStoreException.class, iter::hasNext);
+ assertThrows(InvalidStateStoreException.class, iter::next);
+ assertThrows(InvalidStateStoreException.class, iter::peekNextKey);
+ }
+
@Test
public void shouldSkipTombstones() {
final NavigableMap<Bytes, Optional<byte[]>> staging = new TreeMap<>();