diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index c0f6797eddd07..a4bbcefaad0c1 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -27,6 +27,7 @@ import kafka.utils._ import org.apache.kafka.clients.{CommonClientConfigs, admin} import org.apache.kafka.clients.admin._ import org.apache.kafka.clients.consumer.{ConsumerConfig, KafkaConsumer, OffsetAndMetadata} +import org.apache.kafka.common.errors.CoordinatorNotAvailableException import org.apache.kafka.common.serialization.StringDeserializer import org.apache.kafka.common.utils.Utils import org.apache.kafka.common.{KafkaException, Node, TopicPartition} diff --git a/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala b/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala index eaa0c853e66e9..a92be96fd53dc 100644 --- a/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala @@ -20,8 +20,10 @@ import joptsimple.OptionException import kafka.admin.ConsumerGroupCommand.ConsumerGroupService import kafka.server.KafkaConfig import kafka.utils.TestUtils +import org.apache.kafka.clients.consumer.OffsetAndMetadata import org.apache.kafka.clients.producer.ProducerRecord import org.apache.kafka.common.TopicPartition +import org.apache.kafka.common.errors.CoordinatorNotAvailableException import org.junit.Assert._ import org.junit.Test @@ -86,7 +88,7 @@ class ResetConsumerGroupOffsetTest extends ConsumerGroupCommandTest { val args = Array("--bootstrap-server", brokerList, "--reset-offsets", "--group", "missing.group", "--all-topics", "--to-current", "--execute") val consumerGroupCommand = getConsumerGroupService(args) - val resetOffsets = consumerGroupCommand.resetOffsets() + val resetOffsets = tryResetOffsets(consumerGroupCommand) assertEquals(Map.empty, resetOffsets) assertEquals(resetOffsets, committedOffsets(group = "missing.group")) } @@ -389,7 +391,21 @@ class ResetConsumerGroupOffsetTest extends ConsumerGroupCommandTest { } private def resetOffsets(consumerGroupService: ConsumerGroupService): Map[TopicPartition, Long] = { - consumerGroupService.resetOffsets().mapValues(_.offset) + tryResetOffsets(consumerGroupService).mapValues(_.offset) } + private def tryResetOffsets(consumerGroupService: ConsumerGroupService): Map[TopicPartition, OffsetAndMetadata] = { + var retries = 0 + var gotResult = false + + while (retries < 3) { + try { + return consumerGroupService.resetOffsets() + } catch { + case _: CoordinatorNotAvailableException => retries += 1 + } + } + + fail(s"Could not reset offsets. $retries retries resulted in ${classOf[CoordinatorNotAvailableException].getName}") + } }