Matthias J. Sax created KAFKA-20808:
---------------------------------------

             Summary: 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


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