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

Eswarar Siva updated KAFKA-20721:
---------------------------------
    Description: 
Note (2026-06-22): an earlier revision of this ticket described a pause/resume 
duplicate listing mechanism. That was wrong for this report; the original 
production occurrence had no pause/resume. The pause/resume duplicate listing 
is tracked separately in KAFKA-20724. This description is reverted to the 
original crash, which matches KAFKA-20456.

A StreamThread dies during cold start rebalancing with:

java.lang.IllegalStateException: Task X_Y was not found in the state updater. 
This indicates a bug.
    at TaskManager.waitForFuture(TaskManager.java:711)
    at TaskManager.addToTasksToClose(TaskManager.java:685)
    at TaskManager.handleTasksInStateUpdater(TaskManager.java:640)
    at TaskManager.handleAssignment(TaskManager.java:378)
    -> KafkaException: User rebalance callback throws an error -> 
SHUTDOWN_CLIENT

Environment: Kafka Streams 4.1.2, exactly_once_v2, several stateful instances, 
large RocksDB state, cold start with heavy rebalancing. The application does 
not call pause() or resume().

Mechanism (matches KAFKA-20456): during the rebalance, a remove issued from 
revokeTasksInStateUpdater waits in TaskManager.waitForFuture, which has a 
hardcoded 5 minute timeout. The state updater is stalled (restore under load), 
so the get() times out and waitForFuture returns null after logging "The state 
updater wasn't able to remove task X in time. The state updater thread may be 
dead." The revoked partitions are already removed from
remainingRevokedPartitions before that wait, so the task is not suspended or re 
tracked. The task is still visible in stateUpdater.tasks() because a pending 
REMOVE action is not included in the snapshot (only pending ADD is). The 
following onAssignment then enumerates tasks(), finds the task absent from the 
new assignment, and enqueues a second remove via handleTasksInStateUpdater. 
When the stalled state updater finally drains its action queue (FIFO), the 
first remove succeeds and the second finds the task in none of its collections, 
so DefaultStateUpdater.removeTask
completes the future with null and logs "Task X could not be removed from the 
state updater because the state updater does not own this task." waitForFuture 
in addToTasksToClose receives that null and throws the ISE.

Production evidence: four occurrences on 4.1.2, all the same order, the timeout 
WARN first, then 13 seconds to 2 minutes later the "does not own this task" 
WARN, then the ISE. No STANDBY_UPDATING, no Duplicate key, no pause or resume 
in any of them.

Relationship to 4.3.0: TaskManager.waitForFuture is identical in 4.1.2, 4.3.0 
and trunk (same 5 minute timeout, same throw on a null completion). The 
KAFKA-20456 changes present in 4.3.0 are the two RocksDB ones (lighter offsets 
column family, avoid persisting closed state), which reduce the stalls that 
trigger the timeout. The timeout bound and the leaked task cleanup are not in 
4.3.0 or trunk. So 4.3.0 reduces the likelihood but does not remove the timeout 
to duplicate remove to ISE path. No recurrence on 4.3.0 so far, but the 
observation window is short.

Caveat on the trigger: we did not capture a thread dump or GC log at the stall, 
so we can confirm the state updater stalled past 5 minutes but not that the 
staller was RocksDB specifically rather than, for example, a long GC pause.

Affects version: 4.1.2 (observed). The same code path exists in 4.2.x, 4.3.0 
and trunk.

Suggested direction: bound the waitForFuture timeout and clean up the leaked 
task when the remove times out, so a slow remove does not leave a task that is 
later removed twice. 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
>
>
> Note (2026-06-22): an earlier revision of this ticket described a 
> pause/resume duplicate listing mechanism. That was wrong for this report; the 
> original production occurrence had no pause/resume. The pause/resume 
> duplicate listing is tracked separately in KAFKA-20724. This description is 
> reverted to the original crash, which matches KAFKA-20456.
> A StreamThread dies during cold start rebalancing with:
> java.lang.IllegalStateException: Task X_Y was not found in the state updater. 
> This indicates a bug.
>     at TaskManager.waitForFuture(TaskManager.java:711)
>     at TaskManager.addToTasksToClose(TaskManager.java:685)
>     at TaskManager.handleTasksInStateUpdater(TaskManager.java:640)
>     at TaskManager.handleAssignment(TaskManager.java:378)
>     -> KafkaException: User rebalance callback throws an error -> 
> SHUTDOWN_CLIENT
> Environment: Kafka Streams 4.1.2, exactly_once_v2, several stateful 
> instances, large RocksDB state, cold start with heavy rebalancing. The 
> application does not call pause() or resume().
> Mechanism (matches KAFKA-20456): during the rebalance, a remove issued from 
> revokeTasksInStateUpdater waits in TaskManager.waitForFuture, which has a 
> hardcoded 5 minute timeout. The state updater is stalled (restore under 
> load), so the get() times out and waitForFuture returns null after logging 
> "The state updater wasn't able to remove task X in time. The state updater 
> thread may be dead." The revoked partitions are already removed from
> remainingRevokedPartitions before that wait, so the task is not suspended or 
> re tracked. The task is still visible in stateUpdater.tasks() because a 
> pending REMOVE action is not included in the snapshot (only pending ADD is). 
> The following onAssignment then enumerates tasks(), finds the task absent 
> from the new assignment, and enqueues a second remove via 
> handleTasksInStateUpdater. When the stalled state updater finally drains its 
> action queue (FIFO), the first remove succeeds and the second finds the task 
> in none of its collections, so DefaultStateUpdater.removeTask
> completes the future with null and logs "Task X could not be removed from the 
> state updater because the state updater does not own this task." 
> waitForFuture in addToTasksToClose receives that null and throws the ISE.
> Production evidence: four occurrences on 4.1.2, all the same order, the 
> timeout WARN first, then 13 seconds to 2 minutes later the "does not own this 
> task" WARN, then the ISE. No STANDBY_UPDATING, no Duplicate key, no pause or 
> resume in any of them.
> Relationship to 4.3.0: TaskManager.waitForFuture is identical in 4.1.2, 4.3.0 
> and trunk (same 5 minute timeout, same throw on a null completion). The 
> KAFKA-20456 changes present in 4.3.0 are the two RocksDB ones (lighter 
> offsets column family, avoid persisting closed state), which reduce the 
> stalls that trigger the timeout. The timeout bound and the leaked task 
> cleanup are not in 4.3.0 or trunk. So 4.3.0 reduces the likelihood but does 
> not remove the timeout to duplicate remove to ISE path. No recurrence on 
> 4.3.0 so far, but the observation window is short.
> Caveat on the trigger: we did not capture a thread dump or GC log at the 
> stall, so we can confirm the state updater stalled past 5 minutes but not 
> that the staller was RocksDB specifically rather than, for example, a long GC 
> pause.
> Affects version: 4.1.2 (observed). The same code path exists in 4.2.x, 4.3.0 
> and trunk.
> Suggested direction: bound the waitForFuture timeout and clean up the leaked 
> task when the remove times out, so a slow remove does not leave a task that 
> is later removed twice. 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