lucasbru commented on code in PR #23542:
URL: https://github.com/apache/kafka/pull/23542#discussion_r4069305655
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java:
##########
@@ -3503,44 +3583,70 @@ private StreamsGroupMember
getOrMaybeCreateStaticStreamsGroupMember(
boolean memberIsJoining,
List<CoordinatorRecord> records
) {
- StreamsGroupMember existingStaticMemberOrNull =
group.staticMember(instanceId);
- if (memberIsJoining) {
- // A new static member joins or the existing static member rejoins.
- if (existingStaticMemberOrNull == null) {
- // New static member.
- StreamsGroupMember newMember =
group.getOrCreateDefaultMember(memberId);
- log.info("[GroupId {}][MemberId {}] Static member {} with
instance id {} joins the streams group.",
- group.groupId(), memberId, memberId, instanceId);
- return newMember;
- } else {
- throwIfInstanceIdIsUnreleased(existingStaticMemberOrNull,
group.groupId(), memberId, instanceId);
-
- // Copy the member but with its new member id.
- StreamsGroupMember newMember = new
StreamsGroupMember.Builder(existingStaticMemberOrNull, memberId)
- .setMemberEpoch(0)
- .setPreviousMemberEpoch(0)
- .build();
-
- replaceStreamsMember(records, group,
existingStaticMemberOrNull, newMember);
-
- log.info("[GroupId {}][MemberId {}] Static member with
instance id {} re-joins the streams group " +
- "using the streams protocol. Created a new member {}
to replace the existing member {}.",
- group.groupId(), memberId, instanceId, memberId,
existingStaticMemberOrNull.memberId());
-
- return newMember;
- }
- } else {
- throwIfStaticMemberIsUnknown(existingStaticMemberOrNull,
instanceId);
- throwIfInstanceIdIsFenced(existingStaticMemberOrNull,
group.groupId(), memberId, instanceId);
+ if (!memberIsJoining) {
+ StreamsGroupMember staticMember = group.staticMember(instanceId);
+ throwIfStaticMemberIsUnknown(staticMember, instanceId);
+ throwIfInstanceIdIsFenced(staticMember, group.groupId(), memberId,
instanceId);
throwIfStreamsGroupMemberEpochIsInvalid(
- existingStaticMemberOrNull,
+ staticMember,
memberEpoch,
ownedActiveTasks,
ownedStandbyTasks,
ownedWarmupTasks
);
- return existingStaticMemberOrNull;
+ return staticMember;
}
+
+ StreamsGroupMember existingMemberOrNull =
group.members().get(memberId);
+ if (existingMemberOrNull != null) {
+ // The member id is known. A member id must never acquire a
different instance id, so
+ // the member must be the static member owning the instance id.
+ String existingInstanceId = existingMemberOrNull.instanceId() ==
null ?
+ null : existingMemberOrNull.instanceId().orElse(null);
+ throwIfMemberIdHasDifferentInstanceId(group.groupId(), memberId,
existingInstanceId, instanceId);
Review Comment:
I'm only learning about the "relabeling" code as well, but this rejection
happens after the caller has already called
`group.relabelRefinedAssignment(oldMemberId, memberId)`. If member1 joins with
an instance id that used to belong to a departed member2 who still has a cached
refined-assignment entry, the caller relabels member2's assignment to member1
before this method throws for member1 already owning a different instance id.
It seems this could happen before this PR as well, but here we are widening the
error surface. I think the relabel needs to happen after this validation.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]