From 61a2acfd9d124f944cee14166b49ed74efa8c5a1 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Tue, 26 Nov 2019 18:38:06 -0800 Subject: [PATCH 1/8] return true when isNew --- .../main/scala/kafka/coordinator/group/MemberMetadata.scala | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala index 1014f40205dc6..b3918f718e806 100644 --- a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala @@ -88,7 +88,9 @@ private[group] class MemberMetadata(var memberId: String, def shouldKeepAlive(deadlineMs: Long): Boolean = { if (isAwaitingJoin) !isNew || latestHeartbeat + GroupCoordinator.NewMemberJoinTimeoutMs > deadlineMs - else awaitingSyncCallback != null || + else if (isNew || isAwaitingSync) + true + else latestHeartbeat + sessionTimeoutMs > deadlineMs } From 74a9fe92913a03e27d46eb1f60bb1f090a8f3fa8 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Wed, 27 Nov 2019 12:14:26 -0800 Subject: [PATCH 2/8] reorder calls in onCompleteJoin --- .../scala/kafka/coordinator/group/GroupCoordinator.scala | 4 ++-- .../scala/kafka/coordinator/group/MemberMetadata.scala | 9 ++++++--- 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 24a17807d9d31..309d884592fb0 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -1091,9 +1091,9 @@ class GroupCoordinator(val brokerId: Int, leaderId = group.leaderOrNull, error = Errors.NONE) - group.maybeInvokeJoinCallback(member, joinResult) - completeAndScheduleNextHeartbeatExpiration(group, member) member.isNew = false + completeAndScheduleNextHeartbeatExpiration(group, member) + group.maybeInvokeJoinCallback(member, joinResult) } } } diff --git a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala index b3918f718e806..b229456fffb1c 100644 --- a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala @@ -86,11 +86,14 @@ private[group] class MemberMetadata(var memberId: String, } def shouldKeepAlive(deadlineMs: Long): Boolean = { - if (isAwaitingJoin) - !isNew || latestHeartbeat + GroupCoordinator.NewMemberJoinTimeoutMs > deadlineMs - else if (isNew || isAwaitingSync) + 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 } From 839227369ccd8958650ac7983b4a9c3284059693 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Wed, 11 Dec 2019 15:16:48 -0800 Subject: [PATCH 3/8] Revert "reorder calls in onCompleteJoin" This reverts commit 74a9fe92913a03e27d46eb1f60bb1f090a8f3fa8. --- .../scala/kafka/coordinator/group/GroupCoordinator.scala | 4 ++-- .../scala/kafka/coordinator/group/MemberMetadata.scala | 9 +++------ 2 files changed, 5 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 309d884592fb0..24a17807d9d31 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -1091,9 +1091,9 @@ class GroupCoordinator(val brokerId: Int, leaderId = group.leaderOrNull, error = Errors.NONE) - member.isNew = false - completeAndScheduleNextHeartbeatExpiration(group, member) group.maybeInvokeJoinCallback(member, joinResult) + completeAndScheduleNextHeartbeatExpiration(group, member) + member.isNew = false } } } diff --git a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala index b229456fffb1c..b3918f718e806 100644 --- a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala @@ -86,14 +86,11 @@ private[group] class MemberMetadata(var memberId: String, } def shouldKeepAlive(deadlineMs: Long): Boolean = { - 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 + if (isAwaitingJoin) + !isNew || latestHeartbeat + GroupCoordinator.NewMemberJoinTimeoutMs > deadlineMs + else if (isNew || isAwaitingSync) true else - // Otherwise check for session expiration latestHeartbeat + sessionTimeoutMs > deadlineMs } From 151477fc2fa0d8550831fed8778c5da16bba4247 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Wed, 11 Dec 2019 15:17:02 -0800 Subject: [PATCH 4/8] typo --- .../unit/kafka/coordinator/group/GroupCoordinatorTest.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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..23909181f4b41 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -1882,7 +1882,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) From 321fb1c785308eb63097708a926ea366897827da Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Wed, 18 Dec 2019 14:24:23 -0800 Subject: [PATCH 5/8] should fall to new member timeout --- .../scala/kafka/coordinator/group/MemberMetadata.scala | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala index b3918f718e806..b229456fffb1c 100644 --- a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala @@ -86,11 +86,14 @@ private[group] class MemberMetadata(var memberId: String, } def shouldKeepAlive(deadlineMs: Long): Boolean = { - if (isAwaitingJoin) - !isNew || latestHeartbeat + GroupCoordinator.NewMemberJoinTimeoutMs > deadlineMs - else if (isNew || isAwaitingSync) + 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 } From 18d441dabff2e08f4bc404243a26c2b9cafa8d73 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 19 Dec 2019 16:47:07 -0800 Subject: [PATCH 6/8] add unit test --- .../coordinator/group/GroupCoordinatorTest.scala | 15 +++++++++++++++ 1 file changed, 15 insertions(+) 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 23909181f4b41..324fec095918d 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,21 @@ class GroupCoordinatorTest { assertEquals(Errors.INCONSISTENT_GROUP_PROTOCOL, joinGroupResult.error) } + @Test + def testCompleteNewMemberTimeoutHeartbeat(): Unit = { + // New member joins the group + var responseFuture = sendJoinGroup(groupId, JoinGroupRequest.UNKNOWN_MEMBER_ID, protocolType, protocols, None, DefaultSessionTimeout, DefaultRebalanceTimeout, false) + timer.advanceClock(GroupInitialRebalanceDelay + 1) + + val joinResult = Await.result(responseFuture, Duration(DefaultRebalanceTimeout + 100, TimeUnit.MILLISECONDS)) + val group = groupCoordinator.groupManager.getGroup(groupId).get + assertEquals(Errors.NONE, joinResult.error) + assertEquals(0, group.allMemberMetadata.count(_.isNew)) + + // Make sure the delayed heartbeat based on the new member timeout has been completed + assertEquals(1, groupCoordinator.heartbeatPurgatory.numDelayed) + } + @Test def testNewMemberJoinExpiration(): Unit = { // This tests new member expiration during a protracted rebalance. We first create a From d4d920978537f855afb1756bf3539c71a98c07f6 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 19 Dec 2019 16:47:52 -0800 Subject: [PATCH 7/8] indentation --- .../kafka/coordinator/group/MemberMetadata.scala | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala index b229456fffb1c..f9a45d4e9d672 100644 --- a/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/MemberMetadata.scala @@ -86,15 +86,16 @@ private[group] class MemberMetadata(var memberId: String, } def shouldKeepAlive(deadlineMs: Long): Boolean = { - if (isNew) - // New members are expired after the static join timeout + 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 + } else if (isAwaitingJoin || isAwaitingSync) { + // Don't remove members as long as they have a request in purgatory true - else - // Otherwise check for session expiration + } else { + // Otherwise check for session expiration latestHeartbeat + sessionTimeoutMs > deadlineMs + } } /** From 69477f8f5aa319bcf6993b9bb6f35babcd5eb14b Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 20 Dec 2019 14:44:28 -0800 Subject: [PATCH 8/8] refactor test --- .../group/GroupCoordinatorTest.scala | 25 +++++++++++++++---- 1 file changed, 20 insertions(+), 5 deletions(-) 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 324fec095918d..50dc0db583aa9 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -314,18 +314,33 @@ class GroupCoordinatorTest { } @Test - def testCompleteNewMemberTimeoutHeartbeat(): Unit = { - // New member joins the group - var responseFuture = sendJoinGroup(groupId, JoinGroupRequest.UNKNOWN_MEMBER_ID, protocolType, protocols, None, DefaultSessionTimeout, DefaultRebalanceTimeout, false) + 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)) - // Make sure the delayed heartbeat based on the new member timeout has been completed - assertEquals(1, groupCoordinator.heartbeatPurgatory.numDelayed) + 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