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)