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

Lucas Brutschy commented on KAFKA-20721:
----------------------------------------

Hey, while digging into you original bug report, I also noticed the 
pause/resume bug, and filed https://issues.apache.org/jira/browse/KAFKA-20721


But I created a separate ticket for it. Pause/resume is a rarely used feature, 
triggering it at the right time to trigger this bug should be very very rare. 
Still worth a fix, but not the problem that you originally reported (no 
pause/resume was every mentioned).

 

I don't agree that we should just repurpose you original report for something 
that you didn't really do in (pausing/resuming). 

 

Can you try reproducing your problem with 4.3.0? This may be a recurrence of 
https://issues.apache.org/jira/browse/KAFKA-20456. I don't think the 5 minute 
timeout is correct here, as it can lead to duplication of tasks -> which in 
turn leads to duplication of removals later on, triggering this exception.

> 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