[ 
https://issues.apache.org/jira/browse/KAFKA-20748?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18092606#comment-18092606
 ] 

Lucas Brutschy commented on KAFKA-20748:
----------------------------------------

Draft fix: https://github.com/confluentinc/kafka/commit/9f38e384ab

> StoredDescriptionTopologyEpoch can diverge from the topology-description 
> plugin across failures and races
> ---------------------------------------------------------------------------------------------------------
>
>                 Key: KAFKA-20748
>                 URL: https://issues.apache.org/jira/browse/KAFKA-20748
>             Project: Kafka
>          Issue Type: Sub-task
>            Reporter: Lucas Brutschy
>            Assignee: Alieh Saeedi
>            Priority: Major
>
> The broker tracks which topology a Streams group's plugin entry holds in 
> StoredDescriptionTopologyEpoch, and uses that belief to decide whether to 
> solicit a fresh push from the client (TopologyDescriptionRequired), whether 
> the group's plugin entry still needs deleteTopology, and what to serve from 
> the describe API. The plugin call (setTopology/deleteTopology) and the 
> metadata write that records the epoch are two separate, non-atomic steps — an 
> external side effect followed by a committed group-metadata record. A crash 
> or error between the two steps, or a client push that races a concurrent 
> delete, can leave the broker's belief and the plugin permanently disagreeing 
> in one of two ways.
> In the first, a leak: the plugin holds a topology the broker believes is 
> absent (for example, setTopology succeeds but the follow-up epoch write 
> fails), so no cleanup path ever reclaims it and the entry lingers 
> indefinitely. In the second, a loss: the plugin holds nothing but the broker 
> believes a topology is stored (for example, a delete removes the plugin entry 
> but the clearing write fails, or a delete races a push), so the broker never 
> re-solicits and the describe API reports NOT_STORED for a group that should 
> be able to recover.
> This affects every operation that mutates the plugin: a client topology push, 
> the periodic cleanup that deletes entries for naturally expired empty groups, 
> an explicit DeleteGroups, and the conversion of an empty Streams group to a 
> classic group on a classic JoinGroup (which tombstones the Streams metadata 
> and would orphan the plugin entry with no cleanup at all). The cleanup path 
> additionally has no ordering guarantee between deleteTopology and a 
> concurrent setTopology, so a member that rejoins and re-pushes while a delete 
> is in flight can have its topology silently removed.
> *Suggested fix*
> A fix was added to the KIP and needs to be implemented as well. The idea is 
> to make StoredDescriptionTopologyEpoch three-valued by adding a -2 
> "uncertain" marker alongside the existing real epoch (>= 0, the plugin 
> definitely holds that topology) and -1 (the plugin definitely holds nothing). 
> Uncertain means the plugin may or may not hold a topology; it is read like -1 
> for the solicitation decision (so the broker re-solicits) and like a real 
> epoch for the deletion decision (so the group stays reclaimable). That 
> combination is what makes a half-finished operation self-healing rather than 
> stuck: a group left uncertain both asks the client to re-push and remains 
> eligible for cleanup.
> Every operation that is about to disturb the plugin first commits -2 durably 
> and only then calls the plugin; on success it writes the final value (the 
> pushed epoch after a push, -1 after a delete). Any failure after the -2 
> commit leaves the group uncertain, which resolves on its own — the client 
> re-pushes on its next heartbeat, or the periodic cleanup re-attempts the 
> delete. This barrier is applied to all four operations above; for the 
> classic-JoinGroup conversion the plugin delete runs before the group is 
> converted, and on plugin failure the join is rejected with a retriable error 
> so the group stays a Streams group at -2 rather than being converted with an 
> orphaned entry.
> The cleanup path also closes the delete-vs-push race with a smart finalize: 
> after a successful delete it re-checks the stored epoch; if it is still -2 it 
> clears to -1, but if a racing push advanced it to a real epoch it writes -2 
> back to force a re-solicit, since the pushed topology may have been removed 
> by the delete. The only residual cost is one redundant push when the raced 
> topology actually survived.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to