[
https://issues.apache.org/jira/browse/KAFKA-20721?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18091788#comment-18091788
]
Lucas Brutschy commented on KAFKA-20721:
----------------------------------------
[~zoro30102000] I think you proposal makes sense. If we can complete all
futures during shutdown, we should not block here when we close state updater.
I thought about the timeout. It is fairly difficult to cleanly deal with tasks
that we cannot remove. We have to make sure that the rocksdb locks are
released. I wonder if it would just introduce extra complications to have the
timeout. It must be a reasonable to expect the state updater to not just live
lock. If the state updater does block for more than the rebalance timeout due
to rocksdb, this can be fixed by increasing the rebalance timeout or tuning the
rocksdb configurations. But we'd cleanly solve it by rejoining the group. if
the state updater crashes, it must take the stream thread down with it.
[~nikita-shupletsov] these races are valid but they only affect the
pausing/resuming functionality and are much less severe than the bug in this
ticket. I created a separate ticket for them, see above.
Full disclosure I will be on leave starting next week and may not be able to
follow up, but given the severity, I am sure [~mjsax] or [~bbejeck] would be
available to review a fix PR.
> Streams: "Task not found in the state updater. This indicates a bug." after a
> state updater remove times out during cold start rebalancing
> ------------------------------------------------------------------------------------------------------------------------------------------
>
> 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)