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