lucasbru commented on code in PR #23542:
URL: https://github.com/apache/kafka/pull/23542#discussion_r4070205417


##########
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:
   Thanks for the update, looks good to me.
   
   Unfortunately @mjsax  is out of a couple of weeks. Maybe @squah-confluent 
knows, having reviewed the code. I'm trying to build more context at the 
moment. I suppose the timeline data structure is used here to roll back in case 
of a log-append failure, but as you say it won't help when an exception is 
thrown.
   
   It seems this is "mostly" harmless today:
   
   1) all "routine" exceptions are thrown before the update to the refined 
assignment.
   2) if we do throw after a refined assignment update, the cache is gated by 
the assignment epoch. If we cache a refined assignment that is the result of a 
never-written group change (group epoch bump) or never-written target 
assignment, the refined assignment will in many cases be dropped and recomputed 
on the next heartbeat. But I think not always - there is a case where we cache 
a refined assignment for N+1, fall back to epoch N after an exception, then 
immediately bump again to N+1 and attempt to reuse an invalid cache. We should 
probably invalidate the cache if assignment epoch < refined assignment epoch at 
the beginning of the heartbeat handler. But not sure if this closes all holes.
   
   I'm wondering if it wouldn't be easier to persist the refined assignment. 
This code is not live yet so, so there is no risk ATM.



-- 
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]

Reply via email to