[ 
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)

Reply via email to