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

Reply via email to