Repository: kafka Updated Branches: refs/heads/0.10.2 29bc8905a -> e2cc88602
HOTFIX: Consumer offsets not properly loaded on coordinator failover Author: Jason Gustafson <[email protected]> Reviewers: Manikumar Reddy <[email protected]>, Ismael Juma <[email protected]> Closes #2436 from hachikuji/hotfix-offset-deletion (cherry picked from commit 8301cc608e0e3c334da96b546bf6d943a1704024) Signed-off-by: Ismael Juma <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/kafka/repo Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/e2cc8860 Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/e2cc8860 Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/e2cc8860 Branch: refs/heads/0.10.2 Commit: e2cc88602fb49f36feec21e965d185d5c3c52d4d Parents: 29bc890 Author: Jason Gustafson <[email protected]> Authored: Thu Jan 26 20:52:54 2017 +0000 Committer: Ismael Juma <[email protected]> Committed: Thu Jan 26 20:53:12 2017 +0000 ---------------------------------------------------------------------- .../scala/kafka/coordinator/GroupMetadata.scala | 15 +-- .../coordinator/GroupMetadataManager.scala | 96 +++++++++----------- .../src/main/scala/kafka/server/KafkaApis.scala | 4 +- .../kafka/api/ConsumerBounceTest.scala | 1 - 4 files changed, 55 insertions(+), 61 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/kafka/blob/e2cc8860/core/src/main/scala/kafka/coordinator/GroupMetadata.scala ---------------------------------------------------------------------- diff --git a/core/src/main/scala/kafka/coordinator/GroupMetadata.scala b/core/src/main/scala/kafka/coordinator/GroupMetadata.scala index eac3151..0e49995 100644 --- a/core/src/main/scala/kafka/coordinator/GroupMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/GroupMetadata.scala @@ -265,6 +265,10 @@ private[coordinator] class GroupMetadata(val groupId: String, initialState: Grou GroupOverview(groupId, protocolType.getOrElse("")) } + def initializeOffsets(offsets: collection.Map[TopicPartition, OffsetAndMetadata]) { + this.offsets ++= offsets + } + def completePendingOffsetWrite(topicPartition: TopicPartition, offset: OffsetAndMetadata) { if (pendingOffsetCommits.contains(topicPartition)) offsets.put(topicPartition, offset) @@ -287,14 +291,11 @@ private[coordinator] class GroupMetadata(val groupId: String, initialState: Grou } def removeOffsets(topicPartitions: Seq[TopicPartition]): immutable.Map[TopicPartition, OffsetAndMetadata] = { - val removedOffsetMap = new mutable.HashMap[TopicPartition, OffsetAndMetadata] - for (topicPartition <- topicPartitions) { + topicPartitions.flatMap { topicPartition => pendingOffsetCommits.remove(topicPartition) - val removedOffsets = offsets.remove(topicPartition) - if (!removedOffsets.isEmpty) - removedOffsetMap.put(topicPartition, removedOffsets.get) - } - removedOffsetMap.toMap + val removedOffset = offsets.remove(topicPartition) + removedOffset.map(topicPartition -> _) + }.toMap } def removeExpiredOffsets(startMs: Long) = { http://git-wip-us.apache.org/repos/asf/kafka/blob/e2cc8860/core/src/main/scala/kafka/coordinator/GroupMetadataManager.scala ---------------------------------------------------------------------- diff --git a/core/src/main/scala/kafka/coordinator/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/GroupMetadataManager.scala index 459297e..45ed77b 100644 --- a/core/src/main/scala/kafka/coordinator/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/GroupMetadataManager.scala @@ -436,13 +436,13 @@ class GroupMetadataManager(val brokerId: Int, } } - val (groupOffsets, noGroupOffsets) = loadedOffsets + val (groupOffsets, noGroupOffsets) = loadedOffsets .groupBy(_._1.group) .mapValues(_.map{ case (groupTopicPartition, offsetAndMetadata) => (groupTopicPartition.topicPartition, offsetAndMetadata)}) .partition(value => loadedGroups.contains(value._1)) loadedGroups.values.foreach { group => - val offsets = groupOffsets.getOrElse(group.groupId, Map.empty) + val offsets = groupOffsets.getOrElse(group.groupId, Map.empty[TopicPartition, OffsetAndMetadata]) loadGroup(group, offsets) onGroupLoaded(group) } @@ -479,30 +479,24 @@ class GroupMetadataManager(val brokerId: Int, } } - private def loadGroup(group: GroupMetadata, offsets: Iterable[(TopicPartition, OffsetAndMetadata)]): Unit = { + private def loadGroup(group: GroupMetadata, offsets: Map[TopicPartition, OffsetAndMetadata]): Unit = { + // offsets are initialized prior to loading the group into the cache to ensure that clients see a consistent + // view of the group's offsets + val loadedOffsets = offsets.mapValues { offsetAndMetadata => + // special handling for version 0: + // set the expiration time stamp as commit time stamp + server default retention time + if (offsetAndMetadata.expireTimestamp == org.apache.kafka.common.requests.OffsetCommitRequest.DEFAULT_TIMESTAMP) + offsetAndMetadata.copy(expireTimestamp = offsetAndMetadata.commitTimestamp + config.offsetsRetentionMs) + else + offsetAndMetadata + } + trace(s"Initialized offsets $loadedOffsets for group ${group.groupId}") + group.initializeOffsets(loadedOffsets) + val currentGroup = addGroup(group) - if (group != currentGroup) { + if (group != currentGroup) debug(s"Attempt to load group ${group.groupId} from log with generation ${group.generationId} failed " + s"because there is already a cached group with generation ${currentGroup.generationId}") - } else { - - offsets.foreach { - case (topicPartition, offsetAndMetadata) => { - val offset = offsetAndMetadata.copy ( - expireTimestamp = { - // special handling for version 0: - // set the expiration time stamp as commit time stamp + server default retention time - if (offsetAndMetadata.expireTimestamp == org.apache.kafka.common.requests.OffsetCommitRequest.DEFAULT_TIMESTAMP) - offsetAndMetadata.commitTimestamp + config.offsetsRetentionMs - else - offsetAndMetadata.expireTimestamp - } - ) - trace("Loaded offset %s for %s.".format(offset, topicPartition)) - group.completePendingOffsetWrite(topicPartition, offset) - } - } - } } /** @@ -548,23 +542,22 @@ class GroupMetadataManager(val brokerId: Int, cleanupGroupMetadata(None) } - def cleanupGroupMetadata(topicPartitions: Option[Seq[TopicPartition]]) { + def cleanupGroupMetadata(deletedTopicPartitions: Option[Seq[TopicPartition]]) { val startMs = time.milliseconds() var offsetsRemoved = 0 groupMetadataCache.foreach { case (groupId, group) => - val (expiredOffsets, groupIsDead, generation) = group synchronized { - // remove expired offsets from the cache - val expiredOffsets = if (topicPartitions.isEmpty) - group.removeExpiredOffsets(startMs) - else - group.removeOffsets(topicPartitions.get) + val (removedOffsets, groupIsDead, generation) = group synchronized { + val removedOffsets = deletedTopicPartitions match { + case Some(topicPartitions) => group.removeOffsets(topicPartitions) + case None => group.removeExpiredOffsets(startMs) + } if (group.is(Empty) && !group.hasOffsets) { info(s"Group $groupId transitioned to Dead in generation ${group.generationId}") group.transitionTo(Dead) } - (expiredOffsets, group.is(Dead), group.generationId) + (removedOffsets, group.is(Dead), group.generationId) } val offsetsPartition = partitionFor(groupId) @@ -573,12 +566,12 @@ class GroupMetadataManager(val brokerId: Int, case Some((magicValue, timestampType, timestamp)) => val partitionOpt = replicaManager.getPartition(appendPartition) partitionOpt.foreach { partition => - val tombstones = expiredOffsets.map { case (topicPartition, offsetAndMetadata) => - trace(s"Removing expired offset and metadata for $groupId, $topicPartition: $offsetAndMetadata") + val tombstones = removedOffsets.map { case (topicPartition, offsetAndMetadata) => + trace(s"Removing expired/deleted offset and metadata for $groupId, $topicPartition: $offsetAndMetadata") val commitKey = GroupMetadataManager.offsetCommitKey(groupId, topicPartition.topic, topicPartition.partition) Record.create(magicValue, timestampType, timestamp, commitKey, null) }.toBuffer - trace(s"Marked ${expiredOffsets.size} offsets in $appendPartition for deletion.") + trace(s"Marked ${removedOffsets.size} offsets in $appendPartition for deletion.") // We avoid writing the tombstone when the generationId is 0, since this group is only using // Kafka for offset storage. @@ -595,11 +588,11 @@ class GroupMetadataManager(val brokerId: Int, // do not need to require acks since even if the tombstone is lost, // it will be appended again in the next purge cycle partition.appendRecordsToLeader(MemoryRecords.withRecords(timestampType, compressionType, tombstones: _*)) - offsetsRemoved += expiredOffsets.size - trace(s"Successfully appended ${tombstones.size} tombstones to $appendPartition for expired offsets and/or metadata for group $groupId") + offsetsRemoved += removedOffsets.size + trace(s"Successfully appended ${tombstones.size} tombstones to $appendPartition for expired/deleted offsets and/or metadata for group $groupId") } catch { case t: Throwable => - error(s"Failed to append ${tombstones.size} tombstones to $appendPartition for expired offsets and/or metadata for group $groupId.", t) + error(s"Failed to append ${tombstones.size} tombstones to $appendPartition for expired/deleted offsets and/or metadata for group $groupId.", t) // ignore and continue } } @@ -888,26 +881,25 @@ object GroupMetadataManager { value.set(PROTOCOL_KEY, groupMetadata.protocol) value.set(LEADER_KEY, groupMetadata.leaderId) - val memberArray = groupMetadata.allMemberMetadata.map { - case memberMetadata => - val memberStruct = value.instance(MEMBERS_KEY) - memberStruct.set(MEMBER_ID_KEY, memberMetadata.memberId) - memberStruct.set(CLIENT_ID_KEY, memberMetadata.clientId) - memberStruct.set(CLIENT_HOST_KEY, memberMetadata.clientHost) - memberStruct.set(SESSION_TIMEOUT_KEY, memberMetadata.sessionTimeoutMs) + val memberArray = groupMetadata.allMemberMetadata.map { memberMetadata => + val memberStruct = value.instance(MEMBERS_KEY) + memberStruct.set(MEMBER_ID_KEY, memberMetadata.memberId) + memberStruct.set(CLIENT_ID_KEY, memberMetadata.clientId) + memberStruct.set(CLIENT_HOST_KEY, memberMetadata.clientHost) + memberStruct.set(SESSION_TIMEOUT_KEY, memberMetadata.sessionTimeoutMs) - if (version > 0) - memberStruct.set(REBALANCE_TIMEOUT_KEY, memberMetadata.rebalanceTimeoutMs) + if (version > 0) + memberStruct.set(REBALANCE_TIMEOUT_KEY, memberMetadata.rebalanceTimeoutMs) - val metadata = memberMetadata.metadata(groupMetadata.protocol) - memberStruct.set(SUBSCRIPTION_KEY, ByteBuffer.wrap(metadata)) + val metadata = memberMetadata.metadata(groupMetadata.protocol) + memberStruct.set(SUBSCRIPTION_KEY, ByteBuffer.wrap(metadata)) - val memberAssignment = assignment(memberMetadata.memberId) - assert(memberAssignment != null) + val memberAssignment = assignment(memberMetadata.memberId) + assert(memberAssignment != null) - memberStruct.set(ASSIGNMENT_KEY, ByteBuffer.wrap(memberAssignment)) + memberStruct.set(ASSIGNMENT_KEY, ByteBuffer.wrap(memberAssignment)) - memberStruct + memberStruct } value.set(MEMBERS_KEY, memberArray.toArray) http://git-wip-us.apache.org/repos/asf/kafka/blob/e2cc8860/core/src/main/scala/kafka/server/KafkaApis.scala ---------------------------------------------------------------------- diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 3ce27db..7bb3ed5 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -197,7 +197,9 @@ class KafkaApis(val requestChannel: RequestChannel, val updateMetadataResponse = if (authorize(request.session, ClusterAction, Resource.ClusterResource)) { val deletedPartitions = replicaManager.maybeUpdateMetadataCache(correlationId, updateMetadataRequest, metadataCache) - coordinator.handleDeletedPartitions(deletedPartitions) + if (deletedPartitions.nonEmpty) + coordinator.handleDeletedPartitions(deletedPartitions) + if (adminManager.hasDelayedTopicOperations) { updateMetadataRequest.partitionStates.keySet.asScala.map(_.topic).foreach { topic => adminManager.tryCompleteDelayedTopicOperations(topic) http://git-wip-us.apache.org/repos/asf/kafka/blob/e2cc8860/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala ---------------------------------------------------------------------- diff --git a/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala b/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala index f98716a..fba3239 100644 --- a/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala +++ b/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala @@ -161,7 +161,6 @@ class ConsumerBounceTest extends IntegrationTestHarness with Logging { } } - @Test def testClose() { val numRecords = 10
