diff --git a/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala b/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala index 838444ce748f1..c7b14e4d2fb13 100644 --- a/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ResetConsumerGroupOffsetTest.scala @@ -413,8 +413,11 @@ class ResetConsumerGroupOffsetTest extends ConsumerGroupCommandTest { val cgcArgs = buildArgsForGroups(Seq(group1, group2), "--all-topics", "--to-offset", "2", "--export") val consumerGroupCommand = getConsumerGroupService(cgcArgs) - produceConsumeAndShutdown(topic = topic1, group = group1, totalMessages = 100, numConsumers = 2) - produceConsumeAndShutdown(topic = topic2, group = group2, totalMessages = 100, numConsumers = 5) + produceConsumeAndShutdown(topic = topic1, group = group1, totalMessages = 100) + produceConsumeAndShutdown(topic = topic2, group = group2, totalMessages = 100) + + awaitConsumerGroupInactive(consumerGroupCommand, group1) + awaitConsumerGroupInactive(consumerGroupCommand, group2) val file = File.createTempFile("reset", ".csv") file.deleteOnExit() @@ -470,6 +473,13 @@ class ResetConsumerGroupOffsetTest extends ConsumerGroupCommandTest { s"Expected offset: $count. Actual offset: ${committedOffsets(topic, group).values.sum}") } + private def awaitConsumerGroupInactive(consumerGroupService: ConsumerGroupService, group: String): Unit = { + TestUtils.waitUntilTrue(() => { + val state = consumerGroupService.collectGroupState(group).state + state == "Empty" || state == "Dead" + }, s"Expected that consumer group is inactive. Actual state: ${consumerGroupService.collectGroupState(group).state}") + } + private def resetAndAssertOffsets(args: Array[String], expectedOffset: Long, dryRun: Boolean = false,