From 22e7349db547a75356fa9c3f7c88dc0ea27d45ce Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 9 Apr 2020 14:28:23 +0200 Subject: [PATCH 1/5] KAFKA-9844; Maximum number of members within a group is not always enforced due to a race condition in join group --- .../kafka/coordinator/group/GroupCoordinator.scala | 11 +++-------- .../coordinator/group/GroupMetadataManager.scala | 10 +++++++--- 2 files changed, 10 insertions(+), 11 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 772d51d19f5c4..d10efcf659c01 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -145,17 +145,12 @@ class GroupCoordinator(val brokerId: Int, responseCallback(JoinGroupResult(memberId, Errors.INVALID_SESSION_TIMEOUT)) } else { val isUnknownMember = memberId == JoinGroupRequest.UNKNOWN_MEMBER_ID - groupManager.getGroup(groupId) match { + groupManager.getGroup(groupId, isUnknownMember) match { case None => // only try to create the group if the group is UNKNOWN AND // the member id is UNKNOWN, if member is specified but group does not // exist we should reject the request. - if (isUnknownMember) { - val group = groupManager.addGroup(new GroupMetadata(groupId, Empty, time)) - doUnknownJoinGroup(group, groupInstanceId, requireKnownMemberId, clientId, clientHost, rebalanceTimeoutMs, sessionTimeoutMs, protocolType, protocols, responseCallback) - } else { - responseCallback(JoinGroupResult(memberId, Errors.UNKNOWN_MEMBER_ID)) - } + responseCallback(JoinGroupResult(memberId, Errors.UNKNOWN_MEMBER_ID)) case Some(group) => group.inLock { if ((groupIsOverCapacity(group) @@ -175,9 +170,9 @@ class GroupCoordinator(val brokerId: Int, joinPurgatory.checkAndComplete(GroupKey(group.groupId)) } } - } } } + } private def doUnknownJoinGroup(group: GroupMetadata, groupInstanceId: Option[String], diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 9f14e13435b19..a543c074856b5 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -209,10 +209,14 @@ class GroupMetadataManager(brokerId: Int, false } /** - * Get the group associated with the given groupId, or null if not found + * Get the group associated with the given groupId - the group is created if createIfNotExist + * is true - or null if not found */ - def getGroup(groupId: String): Option[GroupMetadata] = { - Option(groupMetadataCache.get(groupId)) + def getGroup(groupId: String, createIfNotExist: Boolean = false): Option[GroupMetadata] = { + if (createIfNotExist) + Option(groupMetadataCache.getAndMaybePut(groupId, new GroupMetadata(groupId, Empty, time))) + else + Option(groupMetadataCache.get(groupId)) } /** From 1dab5f151f0017ee82d6de6839941f86899c9775 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 9 Apr 2020 14:36:05 +0200 Subject: [PATCH 2/5] fixup --- .../scala/kafka/coordinator/group/GroupCoordinator.scala | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index d10efcf659c01..7606b142bffc8 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -145,11 +145,10 @@ class GroupCoordinator(val brokerId: Int, responseCallback(JoinGroupResult(memberId, Errors.INVALID_SESSION_TIMEOUT)) } else { val isUnknownMember = memberId == JoinGroupRequest.UNKNOWN_MEMBER_ID + // group is created if it does not exist and the member id is UNKNOWN. if member + // is specified but group does not exist, request is rejected with UNKNOWN_MEMBER_ID groupManager.getGroup(groupId, isUnknownMember) match { case None => - // only try to create the group if the group is UNKNOWN AND - // the member id is UNKNOWN, if member is specified but group does not - // exist we should reject the request. responseCallback(JoinGroupResult(memberId, Errors.UNKNOWN_MEMBER_ID)) case Some(group) => group.inLock { From 16bb5e402e8ba1838b13cc53d2ef1f7736936dd0 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Mon, 13 Apr 2020 19:50:09 +0200 Subject: [PATCH 3/5] add unit test --- .../GroupCoordinatorConcurrencyTest.scala | 33 +++++++++++++++++-- 1 file changed, 31 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala index 50f0f5e7334c3..ff673899e7f54 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala @@ -69,6 +69,8 @@ class GroupCoordinatorConcurrencyTest extends AbstractCoordinatorConcurrencyTest new LeaveGroupOperation ) + var heartbeatPurgatory: DelayedOperationPurgatory[DelayedHeartbeat] = _ + var joinPurgatory: DelayedOperationPurgatory[DelayedJoin] = _ var groupCoordinator: GroupCoordinator = _ @Before @@ -86,8 +88,8 @@ class GroupCoordinatorConcurrencyTest extends AbstractCoordinatorConcurrencyTest val config = KafkaConfig.fromProps(serverProps) - val heartbeatPurgatory = new DelayedOperationPurgatory[DelayedHeartbeat]("Heartbeat", timer, config.brokerId, reaperEnabled = false) - val joinPurgatory = new DelayedOperationPurgatory[DelayedJoin]("Rebalance", timer, config.brokerId, reaperEnabled = false) + heartbeatPurgatory = new DelayedOperationPurgatory[DelayedHeartbeat]("Heartbeat", timer, config.brokerId, reaperEnabled = false) + joinPurgatory = new DelayedOperationPurgatory[DelayedJoin]("Rebalance", timer, config.brokerId, reaperEnabled = false) groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, timer.time, new Metrics()) groupCoordinator.startup(false) @@ -124,6 +126,33 @@ class GroupCoordinatorConcurrencyTest extends AbstractCoordinatorConcurrencyTest verifyConcurrentRandomSequences(createGroupMembers, allOperationsWithTxn) } + @Test + def testConcurrentJoinGroupEnforceGroupMaxSize(): Unit = { + val groupMaxSize = 1 + serverProps.put(KafkaConfig.GroupMaxSizeProp, groupMaxSize.toString) + val config = KafkaConfig.fromProps(serverProps) + + if (groupCoordinator != null) + groupCoordinator.shutdown() + groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, heartbeatPurgatory, + joinPurgatory, timer.time, new Metrics()) + groupCoordinator.startup(false) + + val members = new Group(s"group", nMembersPerGroup, groupCoordinator, replicaManager) + .members + val joinOp = new JoinGroupOperation() + + verifyConcurrentActions(members.toSet.map(joinOp.actionNoVerify)) + + val errors = members.map { member => + val joinGroupResult = joinOp.await(member, DefaultRebalanceTimeout) + joinGroupResult.error + } + + assertEquals(groupMaxSize, errors.count(_ == Errors.NONE)) + assertEquals(members.size-groupMaxSize, errors.count(_ == Errors.GROUP_MAX_SIZE_REACHED)) + } + abstract class GroupOperation[R, C] extends Operation { val responseFutures = new ConcurrentHashMap[GroupMember, Future[R]]() From dce9dd22a861abc182b5f64974e2f9442305ff07 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Tue, 14 Apr 2020 09:08:47 +0200 Subject: [PATCH 4/5] fixup --- .../coordinator/group/GroupCoordinatorConcurrencyTest.scala | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala index ff673899e7f54..a16aef7ea14a1 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala @@ -17,6 +17,7 @@ package kafka.coordinator.group +import java.util.Properties import java.util.concurrent.{ConcurrentHashMap, TimeUnit} import kafka.common.OffsetAndMetadata @@ -129,8 +130,9 @@ class GroupCoordinatorConcurrencyTest extends AbstractCoordinatorConcurrencyTest @Test def testConcurrentJoinGroupEnforceGroupMaxSize(): Unit = { val groupMaxSize = 1 - serverProps.put(KafkaConfig.GroupMaxSizeProp, groupMaxSize.toString) - val config = KafkaConfig.fromProps(serverProps) + val newProperties = new Properties + newProperties.put(KafkaConfig.GroupMaxSizeProp, groupMaxSize.toString) + val config = KafkaConfig.fromProps(serverProps, newProperties) if (groupCoordinator != null) groupCoordinator.shutdown() From c08ba9576f94a0bba6639bc2233e81bc0756a3c5 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Fri, 17 Apr 2020 11:33:00 +0200 Subject: [PATCH 5/5] Address review --- .../kafka/coordinator/group/GroupCoordinator.scala | 2 +- .../kafka/coordinator/group/GroupMetadataManager.scala | 10 +++++++++- 2 files changed, 10 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 7606b142bffc8..f178dc73fe7fe 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -147,7 +147,7 @@ class GroupCoordinator(val brokerId: Int, val isUnknownMember = memberId == JoinGroupRequest.UNKNOWN_MEMBER_ID // group is created if it does not exist and the member id is UNKNOWN. if member // is specified but group does not exist, request is rejected with UNKNOWN_MEMBER_ID - groupManager.getGroup(groupId, isUnknownMember) match { + groupManager.getOrMaybeCreateGroup(groupId, isUnknownMember) match { case None => responseCallback(JoinGroupResult(memberId, Errors.UNKNOWN_MEMBER_ID)) case Some(group) => diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index a543c074856b5..230e30d079e7d 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -208,11 +208,19 @@ class GroupMetadataManager(brokerId: Int, case None => false } + + /** + * Get the group associated with the given groupId or null if not found + */ + def getGroup(groupId: String): Option[GroupMetadata] = { + Option(groupMetadataCache.get(groupId)) + } + /** * Get the group associated with the given groupId - the group is created if createIfNotExist * is true - or null if not found */ - def getGroup(groupId: String, createIfNotExist: Boolean = false): Option[GroupMetadata] = { + def getOrMaybeCreateGroup(groupId: String, createIfNotExist: Boolean): Option[GroupMetadata] = { if (createIfNotExist) Option(groupMetadataCache.getAndMaybePut(groupId, new GroupMetadata(groupId, Empty, time))) else