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