[ 
https://issues.apache.org/jira/browse/KAFKA-20721?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18093304#comment-18093304
 ] 

Nikita Shupletsov commented on KAFKA-20721:
-------------------------------------------

yeah, I found a separate problem. I  created another ticket for that: 
https://issues.apache.org/jira/browse/KAFKA-20764



for this problem I agree that we should switch from timing out to killing the 
stream thread when the state updater thread dies

> Streams: "Task not found in the state updater. This indicates a bug." after a 
> state updater remove times out during cold start rebalancing
> ------------------------------------------------------------------------------------------------------------------------------------------
>
>                 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