[
https://issues.apache.org/jira/browse/KAFKA-20721?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Eswarar Siva updated KAFKA-20721:
---------------------------------
Attachment: (was: task-not-found-repro.patch)
> 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)