squah-confluent commented on code in PR #23542:
URL: https://github.com/apache/kafka/pull/23542#discussion_r4075404377


##########
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:
   > 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.
   
   For assignment offloading, there is a similar problem except with in-flight 
assignor runs for epoch N+1. Right now my plan is to cancel any in-flight runs 
whose epoch >= new epoch at the point we bump the epoch. This is really 
annoying of course because we must be careful to run the check at every place 
we bump the epoch.



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