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

Eswarar Siva updated KAFKA-20721:
---------------------------------
    Attachment:     (was: captured-stacktrace.txt)

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