This is an automated email from the ASF dual-hosted git repository.

mjsax 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 edc0800b480 KAFKA-20808: StreamThread dies with fatal 
InvalidStateStoreException after TaskCorruptedException (#22855)
edc0800b480 is described below

commit edc0800b480a5996c969730d889f01cb838d7fbb
Author: Matthias J. Sax <[email protected]>
AuthorDate: Fri Jul 17 12:38:55 2026 -0700

    KAFKA-20808: StreamThread dies with fatal InvalidStateStoreException after 
TaskCorruptedException (#22855)
    
    This PR ensures that an unsuccessfully opened segment is remove from the
    in-memory cache for open segments.
    
    Reviewers: Nick Telford <[email protected]>, Bill Bejeck
     <[email protected]>
---
 .../streams/state/internals/AbstractSegments.java  |  7 ++++++-
 .../state/internals/AbstractSegmentsTest.java      | 24 ++++++++++++++++++++++
 2 files changed, 30 insertions(+), 1 deletion(-)

diff --git 
a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractSegments.java
 
b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractSegments.java
index a46291175e2..8379230e02a 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractSegments.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractSegments.java
@@ -96,7 +96,12 @@ public abstract class AbstractSegments<S extends Segment> 
implements Segments<S>
                 throw new 
IllegalStateException(newSegment.getClass().getSimpleName() + " already exists. 
Possible concurrent access.");
             }
 
-            openSegmentDB(newSegment, context);
+            try {
+                openSegmentDB(newSegment, context);
+            } catch (final Exception openException) {
+                segments.remove(segmentId);
+                throw openException;
+            }
             return newSegment;
         }
     }
diff --git 
a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSegmentsTest.java
 
b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSegmentsTest.java
index e4108de78bf..499af8653b3 100644
--- 
a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSegmentsTest.java
+++ 
b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractSegmentsTest.java
@@ -112,6 +112,30 @@ abstract class AbstractSegmentsTest<S extends Segments> {
         assertEquals("Tasks [0_0] are corrupted and hence need to be 
re-initialized", thrown.getMessage());
     }
 
+    @Test
+    public void shouldNotRetainFailedSegmentSoRecoveryStaysRecoverable() 
throws Exception {
+        context = getEOSProcessorContext();
+        segments = getSegments();
+        segments.openExisting(context, 0L);
+        final Segment segment = segments.getOrCreateSegmentIfLive(0, context, 
1L);
+
+        assertTrue(segment.isOpen());
+        segments.close();
+
+        simulateUncleanShutdownForAllSegments();
+        segments = getSegments();
+
+        // The first open fails with a (recoverable) TaskCorruptedException 
because the on-disk state is dirty.
+        assertThrows(TaskCorruptedException.class, () -> 
segments.openExisting(context, 0L));
+
+        // A segment that failed to open must not be retained in the segments 
map. Otherwise a subsequent
+        // open would return the stale, unopened segment instead of retrying 
to open it, and querying its
+        // committed offset during task re-initialization would then fail with 
a fatal
+        // InvalidStateStoreException. Re-opening must instead stay 
recoverable, i.e. keep throwing
+        // TaskCorruptedException until the dirty local state is wiped.
+        assertThrows(TaskCorruptedException.class, () -> 
segments.openExisting(context, 0L));
+    }
+
     private void simulateUncleanShutdownForAllSegments() throws Exception {
         for (final File dbDir : 
Objects.requireNonNull(context.stateDir().listFiles())) {
             for (final File storeDir : 
Objects.requireNonNull(dbDir.listFiles())) {

Reply via email to