[ 
https://issues.apache.org/jira/browse/KAFKA-20808?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Matthias J. Sax updated KAFKA-20808:
------------------------------------
    Description: 
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.

  was:
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.


> 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
>             Fix For: 4.4.0, 4.3.2
>
>
> 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