Repository: kafka
Updated Branches:
  refs/heads/trunk 254e3b77d -> 8301cc608


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


Project: http://git-wip-us.apache.org/repos/asf/kafka/repo
Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/8301cc60
Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/8301cc60
Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/8301cc60

Branch: refs/heads/trunk
Commit: 8301cc608e0e3c334da96b546bf6d943a1704024
Parents: 254e3b7
Author: Jason Gustafson <[email protected]>
Authored: Thu Jan 26 20:52:54 2017 +0000
Committer: Ismael Juma <[email protected]>
Committed: Thu Jan 26 20:52:54 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/8301cc60/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/8301cc60/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/8301cc60/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/8301cc60/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

Reply via email to