[
https://issues.apache.org/jira/browse/KAFKA-20870?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18100615#comment-18100615
]
sanghyeok An commented on KAFKA-20870:
--------------------------------------
[~mjsax]Hi!
If there is no one working on this, May I take a look this?
> Fix owned-task reports in StreamsGroupHeartbeat
> -----------------------------------------------
>
> Key: KAFKA-20870
> URL: https://issues.apache.org/jira/browse/KAFKA-20870
> Project: Kafka
> Issue Type: Bug
> Components: group-coordinator, streams
> Reporter: Matthias J. Sax
> Priority: Major
>
> This ticket is partially a bug report, and mainly an improvement idea.
> **Summary:** Broker discards the entire StreamsGroupHeartbeat owned-task
> report unless all three task lists are non-null
> ## Description
> ### What the protocol documents
> `StreamsGroupHeartbeatRequest` has three **independently** nullable
> owned-task lists, each documented per field:
> - `ActiveTasks` — *"Currently owned active tasks for this client. Null if
> unchanged since last heartbeat."*
> - `StandbyTasks` — same wording
> - `WarmupTasks` — same wording
> Nullable arrays distinguish null from empty on the wire (compact encoding:
> `0x00` = null, `0x01` = empty), and both values are meaningful here. An empty
> list is how a member reports that it has released a role's tasks; null is
> documented to mean "unchanged". Nothing in the schema says the three fields
> must be sent together.
> ### What the broker does
> In the streams reconciliation path, `GroupMetadataManager` converts the
> request into a `TasksTuple` only when all three lists are non-null:
> ```java
> TasksTuple ownedTasks = null;
> if (ownedActiveTasks != null && ownedStandbyTasks != null && ownedWarmupTasks
> != null) {
> ownedTasks = TasksTuple.fromHeartbeatRequest(ownedActiveTasks,
> ownedStandbyTasks, ownedWarmupTasks);
> }
> ```
> `ownedTasks` is put into an `Optional`.
> Otherwise the whole report is dropped (an empty `Optional` is create, even if
> some field are not-null), and `CurrentAssignmentBuilder.hasNotReleased` reads
> an absent report conservatively:
> ```java
> return ownedTasks.isEmpty() ||
> ownedTasks.get().containsAny(tasksPendingRevocation);
> ```
> (`ownedTasks` is the `Optional<TasksTuple>`; `isEmpty()` here means "no
> report", not "owns nothing".)
> So a heartbeat reporting two of the three lists is treated identically to one
> reporting none — the lists that *were* sent are discarded.
> ### Consequence
> A member that reports a subset of the three lists never has **any**
> revocation acknowledged, including for the roles it did report. The
> acknowledgement never lands, so the member stays in its revoking state and
> the group never leaves rebalancing. No error is returned and nothing is
> logged, so the group simply
> stops making progress.
> The failure mode is **liveness only, never safety**: the conservative reading
> means the broker never concludes a task was released when it wasn't, so it
> cannot grant a task to a new owner while the previous owner still runs it.
> A second site makes the same assumption:
> `throwIfStreamsGroupMemberEpochIsInvalid`. In the stale-epoch branch
> (`receivedMemberEpoch < member.memberEpoch()`),
> `areOwnedTasksContainedInAssignedTasks` returns `false` for a null list, so a
> member retrying with its previous epoch and any null owned-task list is
> fenced with `FencedMemberEpochException`.
> ### Root cause
> `TasksTuple`'s record constructor requires all three maps to be non-null, so
> "this role was not reported" has no representation in the type the broker
> converts the request into. Three independent absences on the wire are
> collapsed into a single `Optional<TasksTuple>`, whose only lossless inputs
> are all-present and all-absent; everything in between is rounded to absent.
> `TasksTuple` is the right type for an *assignment* (where all three roles are
> always fully known) and the wrong type for a wire-side *report*.
> ### Not reachable with the Apache Java client
> The client is itself all-or-nothing:
> - The JOINING heartbeat sets all three lists to empty lists (an explicit "I
> own nothing").
> - Otherwise all three are set together, and only when the reconciled
> assignment changed.
> - `HeartbeatState.reset()` clears the recorded assignment and is called
> unconditionally on every error response and on every transport failure, so a
> full report is resent after a coordinator change.
> This is therefore latent, and it affects third-party protocol implementations
> — which the per-field wording actively invites. An implementer who reads the
> schema and resends only the list that changed loses all revocation
> acknowledgements, not just the omitted role's.
> ### Contrast with two places that get this right
> - **KIP-848.** `ConsumerGroupHeartbeatRequest` has a *single* nullable
> `TopicPartitions` field with the same convention (*"null if it didn't change
> since the last heartbeat"*), and the consumer-side `CurrentAssignmentBuilder`
> applies the same conservative reading (`if (ownedTopicPartitions == null)
> return true;`). With one field the partial state is unrepresentable, so this
> cannot arise. The streams protocol implements the same mechanism, split
> across three fields.
> - **`ProcessId`**, in this same request, carries the identical *"Null if
> unchanged since last heartbeat"* wording and **is** honoured per field
> (`maybeUpdateProcessId(Optional.ofNullable(processId))`) — because that value
> is persisted in the member record and so survives a coordinator failover. The
> convention works for persisted fields and breaks for transient ones; owned
> tasks are transient.
> ---
> ## Options considered
> Relevant compatibility fact: `StreamsGroupHeartbeatRequest` is
> `validVersions: 0-1` with no `latestVersionUnstable`, so both versions are
> shipped and stable (talking AK 4.4 already even if not release yet).
> **A. Persist partial revocation progress.** Make the acknowledgement
> per-role: subtract the released tasks of each reported role from
> `tasksPendingRevocation` and persist the shrunken set. Failover is then a
> non-issue because progress is in the log. Note that today a partial release
> records *nothing* — `CurrentAssignmentBuilder` returns the same member object
> while `hasNotReleased` is true, and `maybeReconcile` only appends a record
> when the member actually changed — so this adds current-assignment record
> writes on the heartbeat path, up to one per role per revocation instead of
> one per revocation.
> Steady-state cost increase; likely the hardest to land.
> **B. Accumulate partial reports in memory and persist only on completion.**
> Avoids A's extra writes, but the accumulated progress is transient and is
> lost on coordinator failover, with no event the client can observe (its own
> assignment has not changed, so it has no reason to resend). Recovering
> therefore needs a
> broker→client "resend the full report" signal — the shape already used by
> `TopologyDescriptionRequired` in the heartbeat response. Adds transient
> broker state *and* protocol surface to serve a case no client produces.
> **C. Document the all-or-nothing contract and reject violations.** Reword the
> three owned-task fields to say they are reported together, all or none, and
> have the broker reject a partial report with an explicit error instead of
> absorbing it into a rebalance that never finishes. No format change; makes
> the prose match the implementation and both clients; turns a silent hang into
> something diagnosable. Leaves a wire format that still *permits* the illegal
> state, i.e. the contract is prose rather than structure.
> **D. Change the request shape in a new version.** Replace the three lists
> with a single nullable outer field containing three non-nullable inner lists,
> so "nothing changed" is the outer null and a partial report is
> unrepresentable — restoring the property KIP-848 has. This is the correct
> shape, but since v0 and v1 are stable it is a **v2 addition**: the flat
> three-field path has to be retained for v0/v1 clients, so the rounding
> behaviour survives for them regardless.
> C and D compose: fix the wording and reject now; if the request shape is ever
> revised for another reason, adopt the nested form and C's prose becomes
> redundant.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)