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

Reply via email to