diff --git a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala index 1014f40205dc6..f9a45d4e9d672 100644 --- a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala @@ -86,10 +86,16 @@ private[group] class MemberMetadata(var memberId: String, } def shouldKeepAlive(deadlineMs: Long): Boolean = { - if (isAwaitingJoin) - !isNew || latestHeartbeat + GroupCoordinator.NewMemberJoinTimeoutMs > deadlineMs - else awaitingSyncCallback != null || + if (isNew) { + // New members are expired after the static join timeout + latestHeartbeat + GroupCoordinator.NewMemberJoinTimeoutMs > deadlineMs + } else if (isAwaitingJoin || isAwaitingSync) { + // Don't remove members as long as they have a request in purgatory + true + } else { + // Otherwise check for session expiration latestHeartbeat + sessionTimeoutMs > deadlineMs + } } /** 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 cd85e0a86016d..50dc0db583aa9 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -313,6 +313,36 @@ class GroupCoordinatorTest { assertEquals(Errors.INCONSISTENT_GROUP_PROTOCOL, joinGroupResult.error) } + @Test + def testNewMemberTimeoutCompletion(): Unit = { + val sessionTimeout = GroupCoordinator.NewMemberJoinTimeoutMs + 5000 + val responseFuture = sendJoinGroup(groupId, JoinGroupRequest.UNKNOWN_MEMBER_ID, protocolType, protocols, None, sessionTimeout, DefaultRebalanceTimeout, false) + + timer.advanceClock(GroupInitialRebalanceDelay + 1) + + val joinResult = Await.result(responseFuture, Duration(DefaultRebalanceTimeout + 100, TimeUnit.MILLISECONDS)) + val group = groupCoordinator.groupManager.getGroup(groupId).get + val memberId = joinResult.memberId + + assertEquals(Errors.NONE, joinResult.error) + assertEquals(0, group.allMemberMetadata.count(_.isNew)) + + EasyMock.reset(replicaManager) + val syncGroupResult = syncGroupLeader(groupId, joinResult.generationId, memberId, Map(memberId -> Array[Byte]())) + val syncGroupError = syncGroupResult._2 + + assertEquals(Errors.NONE, syncGroupError) + assertEquals(1, group.size) + + timer.advanceClock(GroupCoordinator.NewMemberJoinTimeoutMs + 100) + + // Make sure the NewMemberTimeout is not still in effect, and the member is not kicked + assertEquals(1, group.size) + + timer.advanceClock(sessionTimeout + 100) + assertEquals(0, group.size) + } + @Test def testNewMemberJoinExpiration(): Unit = { // This tests new member expiration during a protracted rebalance. We first create a @@ -1882,7 +1912,7 @@ class GroupCoordinatorTest { val nextGenerationId = joinResult.generationId - // with no leader SyncGroup, the follower's request should failure with an error indicating + // with no leader SyncGroup, the follower's request should fail with an error indicating // that it should rejoin EasyMock.reset(replicaManager) val followerSyncFuture = sendSyncGroupFollower(groupId, nextGenerationId, otherJoinResult.memberId, None)