[
https://issues.apache.org/jira/browse/KAFKA-20808?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Matthias J. Sax reassigned KAFKA-20808:
---------------------------------------
Assignee: Matthias J. Sax
> StreamThread dies with fatal InvalidStateStoreException after
> TaskCorruptedException
> ------------------------------------------------------------------------------------
>
> Key: KAFKA-20808
> URL: https://issues.apache.org/jira/browse/KAFKA-20808
> Project: Kafka
> Issue Type: Bug
> Components: streams
> Affects Versions: 4.3.0
> Reporter: Matthias J. Sax
> Assignee: Matthias J. Sax
> Priority: Critical
>
> First KafkaStreams hits an error opening a store (this code was added via
> KIP-1035, and triggers only via EOS)
> {code:java}
> Caused by: org.apache.kafka.streams.errors.ProcessorStateException: State
> store example-window-store-v1.1700000000000 didn't find a valid state, since
> under EOS it has the risk of getting uncommitted data in stores ... 37
> common frames omitted
> Caused by: org.apache.kafka.streams.errors.ProcessorStateException:
> Invalid state during store open. Expected state to be either empty or closed
> at
> org.apache.kafka.streams.state.internals.AbstractColumnFamilyAccessor.open(AbstractColumnFamilyAccessor.java:93)
> at
> org.apache.kafka.streams.state.internals.RocksDBStore.openDB(RocksDBStore.java:263){code}
> The error is converted into a TaskCorruptedException and KafkaStreams try to
> close the task dirty, and re-initialize it, but fails with
> {code:java}
> org.apache.kafka.streams.errors.InvalidStateStoreException: Store
> example-window-store-v1.1700000000000 is currently closed
> at
> org.apache.kafka.streams.state.internals.RocksDBStore.validateStoreOpen(RocksDBStore.java:472)
> at
> org.apache.kafka.streams.state.internals.RocksDBStore.committedOffset(RocksDBStore.java:745)
> at
> org.apache.kafka.streams.state.internals.AbstractSegments.committedOffset(AbstractSegments.java:196)
> at
> org.apache.kafka.streams.state.internals.AbstractRocksDBSegmentedBytesStore.committedOffset(AbstractRocksDBSegmentedBytesStore.java:329)
> at
> org.apache.kafka.streams.state.internals.WrappedStateStore.committedOffset(WrappedStateStore.java:153)
> at
> org.apache.kafka.streams.state.internals.WrappedStateStore.committedOffset(WrappedStateStore.java:153)
> at
> org.apache.kafka.streams.state.internals.WrappedStateStore.committedOffset(WrappedStateStore.java:153)
> at
> org.apache.kafka.streams.processor.internals.ProcessorStateManager.initializeStoreOffsets(ProcessorStateManager.java:308)
> at
> org.apache.kafka.streams.processor.internals.StateManagerUtil.registerStateStores(StateManagerUtil.java:150)
> at
> org.apache.kafka.streams.processor.internals.StandbyTask.initializeIfNeeded(StandbyTask.java:112)
> at
> org.apache.kafka.streams.processor.internals.TaskManager.addTaskToStateUpdater(TaskManager.java:949)
> at
> org.apache.kafka.streams.processor.internals.TaskManager.addTasksToStateUpdater(TaskManager.java:934)
> at
> org.apache.kafka.streams.processor.internals.TaskManager.checkStateUpdater(TaskManager.java:848)
> at
> org.apache.kafka.streams.processor.internals.StreamThread.checkStateUpdater(StreamThread.java:1409)
> at
> org.apache.kafka.streams.processor.internals.StreamThread.runOnceWithoutProcessingThreads(StreamThread.java:1235)
> at
> org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:956)
> at
> org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:916)"}
> {code}
>
> Root cause:
> When a segmented store fails to open a segment during openExisting under EOS,
> the task is closed dirty and revived as intended, but the retry dies fatally
> instead of recovering. Root cause is a stale, unopened segment left in the
> in-memory AbstractSegments.segments map that the dirty-close never clears.
> Mechanism
> 1. AbstractSegments.getOrCreateSegment inserts the segment into the map (:95)
> before openSegmentDB (:99) — long-standing (KAFKA-3522).
> 2. Under EOS the open now throws TaskCorruptedException via the OPEN/closed
> status-key check (AbstractColumnFamilyAccessor.open:107 →
> RocksDBStore.openDB:279). This fires inside
> AbstractRocksDBSegmentedBytesStore.init at segments.openExisting (:313) —
> before stateStoreContext.register(...) (:319).
> 3. Because the store was never registered, ProcessorStateManager.close()
> (closes only registered stores, :666) skips it, so
> AbstractSegments.close()/segments.clear() (:205/:209) never runs. The
> unopened segment (open==false) survives closeDirty.
> 4. On revive+retry, openExisting short-circuits on segments.containsKey
> (:90-91) and returns the stale segment without reopening; init completes and
> registers the store.
> 5. ProcessorStateManager.initializeStoreOffsets then calls committedOffset
> (:345), which iterates segments (AbstractSegments:196) and calls
> validateStoreOpen on the stale one → fatal InvalidStateStoreException
> (RocksDBStore:894→:478) → caught as generic StreamsException
> (StreamThread:1039, not the TaskCorruptedException branch at :1004) →
> SHUTDOWN_CLIENT.
> While the code in AbstractSegments.getOrCreateSegment is old, it could not
> trip before, because no TaskCorruptedException was throw on this code path.
> This only change with KIP-1035.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)