diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala index 58a68a2c18b16..4d283bc8d62e8 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala @@ -365,7 +365,7 @@ private[group] class GroupMetadata(val groupId: String, initialState: GroupState case None => clientId + GroupMetadata.MemberIdDelimiter + UUID.randomUUID().toString case Some(instanceId) => - instanceId + GroupMetadata.MemberIdDelimiter + currentStateTimestamp.get + instanceId + GroupMetadata.MemberIdDelimiter + UUID.randomUUID().toString } } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala index 0a7bcc74e2804..d72bfc8f5cebc 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -873,6 +873,23 @@ class GroupCoordinatorTest { assertEquals(Errors.FENCED_INSTANCE_ID, invalidHeartbeatResult) } + @Test + def shouldGetDifferentStaticMemberIdAfterEachRejoin(): Unit = { + val initialResult = staticMembersJoinAndRebalance(leaderInstanceId, followerInstanceId) + + val timeAdvance = 1 + var lastMemberId = initialResult.leaderId + for (_ <- 1 to 5) { + EasyMock.reset(replicaManager) + + val joinGroupResult = staticJoinGroup(groupId, JoinGroupRequest.UNKNOWN_MEMBER_ID, + leaderInstanceId, protocolType, protocols, clockAdvance = timeAdvance) + assertTrue(joinGroupResult.memberId.startsWith(leaderInstanceId.get)) + assertNotEquals(lastMemberId, joinGroupResult.memberId) + lastMemberId = joinGroupResult.memberId + } + } + @Test def testOffsetCommitDeadGroup() { val memberId = "memberId"