[
https://issues.apache.org/jira/browse/KAFKA-20721?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Eswarar Siva updated KAFKA-20721:
---------------------------------
Description:
Updated 2026-06-22: the original mechanism in this ticket (a race between
drainRestoredActiveTasks and revokeTasksInStateUpdater) was wrong and is
withdrawn. Both run on the StreamThread and cannot interleave. The corrected
root cause and a reproduction are below; see the comments for the history.
A StreamThread dies with:
java.lang.IllegalStateException: Task X_Y was not found in the state updater.
This indicates a bug.
at
org.apache.kafka.streams.processor.internals.TaskManager.waitForFuture(...)
at
org.apache.kafka.streams.processor.internals.TaskManager.addToTasksToClose(...)
at
org.apache.kafka.streams.processor.internals.TaskManager.shutdownStateUpdater(...)
thrown when DefaultStateUpdater.removeTask completes a removal future with null.
Root cause: DefaultStateUpdater.tasks() can return the same TaskId twice.
tasks() snapshots under executeWithQueuesLocked, which holds
tasksAndActionsLock, restoredActiveTasksLock and exceptionsAndFailedTasksLock,
but not the updatingTasks or pausedTasks maps. pauseTask does
pausedTasks.put(id) then updatingTasks.remove(id), and resumeTask does the
reverse, so for a short window a task is in both maps and streamOfTasks() lists
it twice. ReadOnlyTask has no equals/hashCode, so the returned Set keeps both
wrappers. A caller that does one remove per tasks() entry (shutdownStateUpdater
via addToTasksToClose, and handleAssignment) then removes the same task twice.
The first remove succeeds; the second reaches removeTask, finds the task in
none of the four collections, and completes the future with null, so
waitForFuture throws the ISE. The same duplicate also breaks
TaskManager.allTasks(), which uses Collectors.toMap(Task::id, ...) and throws
IllegalStateException: Duplicate key, killing the StreamThread in runLoop.
Reproduction: a standalone Streams application, several KafkaStreams instances
in one JVM, exactly_once_v2, against a real broker, with cold start churn plus
KafkaStreams.pause()/resume(). The duplicate window is normally sub
microsecond, so to make it deterministic a Byteman rule sleeps right after the
put in pauseTask and resumeTask (no logic change, it only widens the gap that
is already there). With that, the ISE and the Duplicate key error both fire
within a couple of minutes. Reproduced on 4.1.2 and 4.3.0; trunk carries the
same code on this path. With pause/resume disabled there are no duplicates and
no crash, even under aggressive churn. The reproducer, the Byteman rule and
DEBUG logs (TaskManager, DefaultStateUpdater, StoreChangelogReader) are
attached.
Affected versions: 4.1.2 and 4.3.0 (reproduced); 4.2.x and trunk by code
inspection.
Related: KAFKA-17402 is the same family (a duplicate from tasks(), seen there
as shouldGetTasksFromRestoredActiveTasks counting 3 instead of 2); its fix made
the restore completion transition atomic but did not touch
pauseTask/resumeTask, which is the path here. The original production
occurrence behind this ticket was on 4.1.2 with no pause/resume and very large
state, which looks closer to KAFKA-20456 (the state updater stalling and
waitForFuture timing out, fixed in 4.3.0), so the production incident and this
pause/resume reproduction may be two different routes to the same exception.
Suggested fix: stop tasks() from returning a duplicate, for example dedupe by
task id before wrapping, or give ReadOnlyTask equals/hashCode based on the task
id, or make the pause/resume transition atomic with respect to tasks(). The ISE
should stay; it is correctly catching the inconsistency.
was:
Updated 2026-06-22: the original mechanism in this ticket (a race between
drainRestoredActiveTasks and
revokeTasksInStateUpdater) was wrong and is withdrawn. Both run on the
StreamThread and cannot interleave. The
corrected root cause and a reproduction are below; see the comments for the
history.
A StreamThread dies with:
java.lang.IllegalStateException: Task X_Y was not found in the state updater.
This indicates a bug.
at
org.apache.kafka.streams.processor.internals.TaskManager.waitForFuture(...)
at
org.apache.kafka.streams.processor.internals.TaskManager.addToTasksToClose(...)
at
org.apache.kafka.streams.processor.internals.TaskManager.shutdownStateUpdater(...)
thrown when DefaultStateUpdater.removeTask completes a removal future with null.
Root cause: DefaultStateUpdater.tasks() can return the same TaskId twice.
tasks() snapshots under
executeWithQueuesLocked, which holds tasksAndActionsLock,
restoredActiveTasksLock and exceptionsAndFailedTasksLock,
but not the updatingTasks or pausedTasks maps. pauseTask does
pausedTasks.put(id) then updatingTasks.remove(id), and
resumeTask does the reverse, so for a short window a task is in both maps and
streamOfTasks() lists it twice.
ReadOnlyTask has no equals/hashCode, so the returned Set keeps both wrappers. A
caller that does one remove per
tasks() entry (shutdownStateUpdater via addToTasksToClose, and
handleAssignment) then removes the same task twice.
The first remove succeeds; the second reaches removeTask, finds the task in
none of the four collections, and
completes the future with null, so waitForFuture throws the ISE. The same
duplicate also breaks
TaskManager.allTasks(), which uses Collectors.toMap(Task::id, ...) and throws
IllegalStateException: Duplicate key,
killing the StreamThread in runLoop.
Reproduction: a standalone Streams application, several KafkaStreams instances
in one JVM, exactly_once_v2, against a
real broker, with cold start churn plus KafkaStreams.pause()/resume(). The
duplicate window is normally sub
microsecond, so to make it deterministic a Byteman rule sleeps right after the
put in pauseTask and resumeTask (no
logic change, it only widens the gap that is already there). With that, the ISE
and the Duplicate key error both fire
within a couple of minutes. Reproduced on 4.1.2 and 4.3.0; trunk carries the
same code on this path. With pause/resume
disabled there are no duplicates and no crash, even under aggressive churn. The
reproducer, the Byteman rule and DEBUG
logs (TaskManager, DefaultStateUpdater, StoreChangelogReader) are attached.
Affected versions: 4.1.2 and 4.3.0 (reproduced); 4.2.x and trunk by code
inspection.
Related: KAFKA-17402 is the same family (a duplicate from tasks(), seen there
as shouldGetTasksFromRestoredActiveTasks
counting 3 instead of 2); its fix made the restore completion transition atomic
but did not touch pauseTask/resumeTask,
which is the path here. The original production occurrence behind this ticket
was on 4.1.2 with no pause/resume and very
large state, which looks closer to KAFKA-20456 (the state updater stalling and
waitForFuture timing out, fixed in
4.3.0), so the production incident and this pause/resume reproduction may be
two different routes to the same exception.
Suggested fix: stop tasks() from returning a duplicate, for example dedupe by
task id before wrapping, or give
ReadOnlyTask equals/hashCode based on the task id, or make the pause/resume
transition atomic with respect to tasks().
The ISE should stay; it is correctly catching the inconsistency.
> Streams: "Task not found in the state updater. This indicates a bug." caused
> by a duplicate TaskId from DefaultStateUpdater.tasks() under pause/resume
> ------------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: KAFKA-20721
> URL: https://issues.apache.org/jira/browse/KAFKA-20721
> Project: Kafka
> Issue Type: Bug
> Components: streams
> Affects Versions: 4.3.0, 4.1.2
> Environment: Kafka Streams, exactly_once_v2, default state updater,
> several stateful tasks. Hit during a cold start with a lot of rebalancing.
> Verified the code at 4.1.2, 4.2.1, 4.3.0 and trunk.
> Reporter: Eswarar Siva
> Assignee: Eswarar Siva
> Priority: Critical
> Attachments: SuBugPoc.java, captured-stacks-4.1.2.txt,
> captured-stacks-4.3.0.txt, poc-debug-excerpt.txt, pom.xml, widen-window.btm
>
>
> Updated 2026-06-22: the original mechanism in this ticket (a race between
> drainRestoredActiveTasks and revokeTasksInStateUpdater) was wrong and is
> withdrawn. Both run on the StreamThread and cannot interleave. The corrected
> root cause and a reproduction are below; see the comments for the history.
> A StreamThread dies with:
> java.lang.IllegalStateException: Task X_Y was not found in the state updater.
> This indicates a bug.
> at
> org.apache.kafka.streams.processor.internals.TaskManager.waitForFuture(...)
> at
> org.apache.kafka.streams.processor.internals.TaskManager.addToTasksToClose(...)
> at
> org.apache.kafka.streams.processor.internals.TaskManager.shutdownStateUpdater(...)
> thrown when DefaultStateUpdater.removeTask completes a removal future with
> null.
> Root cause: DefaultStateUpdater.tasks() can return the same TaskId twice.
> tasks() snapshots under executeWithQueuesLocked, which holds
> tasksAndActionsLock, restoredActiveTasksLock and
> exceptionsAndFailedTasksLock, but not the updatingTasks or pausedTasks maps.
> pauseTask does pausedTasks.put(id) then updatingTasks.remove(id), and
> resumeTask does the reverse, so for a short window a task is in both maps and
> streamOfTasks() lists it twice. ReadOnlyTask has no equals/hashCode, so the
> returned Set keeps both wrappers. A caller that does one remove per tasks()
> entry (shutdownStateUpdater via addToTasksToClose, and handleAssignment) then
> removes the same task twice. The first remove succeeds; the second reaches
> removeTask, finds the task in none of the four collections, and completes the
> future with null, so waitForFuture throws the ISE. The same duplicate also
> breaks TaskManager.allTasks(), which uses Collectors.toMap(Task::id, ...) and
> throws IllegalStateException: Duplicate key, killing the StreamThread in
> runLoop.
> Reproduction: a standalone Streams application, several KafkaStreams
> instances in one JVM, exactly_once_v2, against a real broker, with cold start
> churn plus KafkaStreams.pause()/resume(). The duplicate window is normally
> sub microsecond, so to make it deterministic a Byteman rule sleeps right
> after the put in pauseTask and resumeTask (no logic change, it only widens
> the gap that is already there). With that, the ISE and the Duplicate key
> error both fire within a couple of minutes. Reproduced on 4.1.2 and 4.3.0;
> trunk carries the same code on this path. With pause/resume disabled there
> are no duplicates and no crash, even under aggressive churn. The reproducer,
> the Byteman rule and DEBUG logs (TaskManager, DefaultStateUpdater,
> StoreChangelogReader) are attached.
> Affected versions: 4.1.2 and 4.3.0 (reproduced); 4.2.x and trunk by code
> inspection.
> Related: KAFKA-17402 is the same family (a duplicate from tasks(), seen there
> as shouldGetTasksFromRestoredActiveTasks counting 3 instead of 2); its fix
> made the restore completion transition atomic but did not touch
> pauseTask/resumeTask, which is the path here. The original production
> occurrence behind this ticket was on 4.1.2 with no pause/resume and very
> large state, which looks closer to KAFKA-20456 (the state updater stalling
> and waitForFuture timing out, fixed in 4.3.0), so the production incident and
> this pause/resume reproduction may be two different routes to the same
> exception.
> Suggested fix: stop tasks() from returning a duplicate, for example dedupe by
> task id before wrapping, or give ReadOnlyTask equals/hashCode based on the
> task id, or make the pause/resume transition atomic with respect to tasks().
> The ISE should stay; it is correctly catching the inconsistency.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)