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