diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index e0e6cb6cffa00..3d735bfbe6a77 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -182,7 +182,6 @@ import org.apache.kafka.common.requests.FindCoordinatorResponse; import org.apache.kafka.common.requests.IncrementalAlterConfigsRequest; import org.apache.kafka.common.requests.IncrementalAlterConfigsResponse; -import org.apache.kafka.common.requests.JoinGroupRequest; import org.apache.kafka.common.requests.LeaveGroupRequest; import org.apache.kafka.common.requests.LeaveGroupResponse; import org.apache.kafka.common.requests.ListGroupsRequest; @@ -3612,6 +3611,25 @@ private boolean dependsOnSpecificNode(ConfigResource resource) { || resource.type() == ConfigResource.Type.BROKER_LOGGER; } + private List getMembersFromGroup(String groupId) { + Collection members; + try { + members = describeConsumerGroups(Collections.singleton(groupId)).describedGroups().get(groupId).get().members(); + } catch (Exception ex) { + throw new KafkaException("Encounter exception when trying to get members from group: " + groupId, ex); + } + + List membersToRemove = new ArrayList<>(); + for (final MemberDescription member : members) { + if (member.groupInstanceId().isPresent()) { + membersToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get())); + } else { + membersToRemove.add(new MemberIdentity().setMemberId(member.consumerId())); + } + } + return membersToRemove; + } + @Override public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(String groupId, RemoveMembersFromConsumerGroupOptions options) { @@ -3623,22 +3641,24 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context = new ConsumerGroupOperationContext<>(groupId, options, deadline, future); - Call findCoordinatorCall = getFindCoordinatorCall(context, - () -> getRemoveMembersFromGroupCall(context)); + List members; + if (options.removeAll()) { + members = getMembersFromGroup(groupId); + } else { + members = options.members().stream().map(MemberToRemove::toMemberIdentity).collect(Collectors.toList()); + } + Call findCoordinatorCall = getFindCoordinatorCall(context, () -> getRemoveMembersFromGroupCall(context, members)); runnable.call(findCoordinatorCall, startFindCoordinatorMs); return new RemoveMembersFromConsumerGroupResult(future, options.members()); } - private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context) { - return new Call("leaveGroup", - context.deadline(), - new ConstantNodeIdProvider(context.node().get().id())) { + private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context, + List members) { + return new Call("leaveGroup", context.deadline(), new ConstantNodeIdProvider(context.node().get().id())) { @Override LeaveGroupRequest.Builder createRequest(int timeoutMs) { - return new LeaveGroupRequest.Builder(context.groupId(), - context.options().members().stream().map( - MemberToRemove::toMemberIdentity).collect(Collectors.toList())); + return new LeaveGroupRequest.Builder(context.groupId(), members); } @Override @@ -3647,7 +3667,7 @@ void handleResponse(AbstractResponse abstractResponse) { // If coordinator changed since we fetched it, retry if (ConsumerGroupOperationContext.hasCoordinatorMoved(response)) { - Call call = getRemoveMembersFromGroupCall(context); + Call call = getRemoveMembersFromGroupCall(context, members); rescheduleFindCoordinatorTask(context, () -> call, this); return; } @@ -3657,10 +3677,8 @@ void handleResponse(AbstractResponse abstractResponse) { final Map memberErrors = new HashMap<>(); for (MemberResponse memberResponse : response.memberResponses()) { - // We set member.id to empty here explicitly, so that the lookup will succeed as user doesn't - // know the exact member.id. memberErrors.put(new MemberIdentity() - .setMemberId(JoinGroupRequest.UNKNOWN_MEMBER_ID) + .setMemberId(memberResponse.memberId()) .setGroupInstanceId(memberResponse.groupInstanceId()), Errors.forCode(memberResponse.errorCode())); } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java index dc346f7c3a1be..322beec39e419 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.annotation.InterfaceStability; import java.util.Collection; +import java.util.Collections; import java.util.HashSet; import java.util.Set; @@ -34,10 +35,21 @@ public class RemoveMembersFromConsumerGroupOptions extends AbstractOptions members; public RemoveMembersFromConsumerGroupOptions(Collection members) { + if (members.isEmpty()) { + throw new IllegalArgumentException("Invalid empty members has been provided"); + } this.members = new HashSet<>(members); } + public RemoveMembersFromConsumerGroupOptions() { + this.members = Collections.emptySet(); + } + public Set members() { return members; } + + public boolean removeAll() { + return members.isEmpty(); + } } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java index 405973b73f226..3845e2f6aac62 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.admin; +import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.internals.KafkaFutureImpl; import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity; @@ -51,9 +52,21 @@ public KafkaFuture all() { if (throwable != null) { result.completeExceptionally(throwable); } else { - for (MemberToRemove memberToRemove : memberInfos) { - if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { - return; + if (removeAll()) { + for (Map.Entry entry: memberErrors.entrySet()) { + Exception exception = entry.getValue().exception(); + if (exception != null) { + Throwable ex = new KafkaException("Encounter exception when trying to remove: " + + entry.getKey(), exception); + result.completeExceptionally(ex); + return; + } + } + } else { + for (MemberToRemove memberToRemove : memberInfos) { + if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { + return; + } } } result.complete(null); @@ -66,6 +79,9 @@ public KafkaFuture all() { * Returns the selected member future. */ public KafkaFuture memberResult(MemberToRemove member) { + if (removeAll()) { + throw new IllegalArgumentException("The method: memberResult is not applicable in 'removeAll' mode"); + } if (!memberInfos.contains(member)) { throw new IllegalArgumentException("Member " + member + " was not included in the original request"); } @@ -93,4 +109,8 @@ private boolean maybeCompleteExceptionally(Map memberErr return false; } } + + private boolean removeAll() { + return memberInfos.isEmpty(); + } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 2510ea727d43d..c3ce1b4728e1b 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -114,6 +114,7 @@ import org.apache.kafka.common.requests.ElectLeadersResponse; import org.apache.kafka.common.requests.FindCoordinatorResponse; import org.apache.kafka.common.requests.IncrementalAlterConfigsResponse; +import org.apache.kafka.common.requests.JoinGroupRequest; import org.apache.kafka.common.requests.LeaveGroupResponse; import org.apache.kafka.common.requests.ListGroupsResponse; import org.apache.kafka.common.requests.ListOffsetResponse; @@ -379,6 +380,23 @@ private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED); } + private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(String groupId, + List groupInstances, + List topicPartitions) { + final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); + List describedGroupMembers = groupInstances.stream().map(groupInstance -> DescribeGroupsResponse.groupMember(JoinGroupRequest.UNKNOWN_MEMBER_ID, + groupInstance, "clientId0", "clientHost", new byte[memberAssignment.remaining()], null)).collect(Collectors.toList()); + DescribeGroupsResponseData data = new DescribeGroupsResponseData(); + data.groups().add(DescribeGroupsResponse.groupMetadata( + groupId, + Errors.NONE, + "", + ConsumerProtocol.PROTOCOL_TYPE, + "", + describedGroupMembers, + Collections.emptySet())); + return data; + } /** * Test that the client properly times out when we don't receive any metadata. */ @@ -2411,6 +2429,50 @@ public void testRemoveMembersFromGroup() throws Exception { assertNull(noErrorResult.all().get()); assertNull(noErrorResult.memberResult(memberOne).get()); assertNull(noErrorResult.memberResult(memberTwo).get()); + + // Test the "removeAll" scenario + final List topicPartitions = Arrays.asList(1, 2, 3).stream().map(partition -> new TopicPartition("my_topic", partition)) + .collect(Collectors.toList()); + // construct the DescribeGroupsResponse + DescribeGroupsResponseData data = prepareDescribeGroupsResponseData(groupId, Arrays.asList(instanceOne, instanceTwo), topicPartitions); + + // Return with partial failure for "removeAll" scenario + // 1 prepare response for AdminClient.describeConsumerGroups + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new DescribeGroupsResponse(data)); + + // 2 KafkaAdminClient encounter partial failure when trying to delete all members + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(responseOne, responseTwo)) + )); + final RemoveMembersFromConsumerGroupResult partialFailureResults = env.adminClient().removeMembersFromConsumerGroup( + groupId, + new RemoveMembersFromConsumerGroupOptions() + ); + ExecutionException exception = assertThrows(ExecutionException.class, () -> partialFailureResults.all().get()); + assertTrue(exception.getCause() instanceof KafkaException); + assertTrue(exception.getCause().getCause() instanceof UnknownMemberIdException); + + // Return with success for "removeAll" scenario + // 1 prepare response for AdminClient.describeConsumerGroups + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new DescribeGroupsResponse(data)); + + // 2. KafkaAdminClient should delete all members correctly + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(responseTwo, + new MemberResponse().setGroupInstanceId(instanceOne).setErrorCode(Errors.NONE.code()) + )) + )); + final RemoveMembersFromConsumerGroupResult successResult = env.adminClient().removeMembersFromConsumerGroup( + groupId, + new RemoveMembersFromConsumerGroupOptions() + ); + assertNull(successResult.all().get()); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java index 41aa386a0590a..92b37e77df596 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java @@ -21,6 +21,7 @@ import java.util.Collections; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; public class RemoveMembersFromConsumerGroupOptionsTest { @@ -31,5 +32,8 @@ public void testConstructor() { assertEquals(Collections.singleton( new MemberToRemove("instance-1")), options.members()); + + // Construct will fail if illegal empty members provided + assertThrows(IllegalArgumentException.class, () -> new RemoveMembersFromConsumerGroupOptions(Collections.emptyList())); } } diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 574e9c66e291a..65ca7b2ee3141 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -27,6 +27,7 @@ import org.apache.kafka.clients.admin.DescribeConsumerGroupsOptions; import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; import org.apache.kafka.clients.admin.MemberDescription; +import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupOptions; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -106,6 +107,7 @@ public class StreamsResetter { private static OptionSpec versionOption; private static OptionSpecBuilder executeOption; private static OptionSpec commandConfigOption; + private static OptionSpec forceOption; private static String usage = "This tool helps to quickly reset an application in order to reprocess " + "its data from scratch.\n" @@ -119,7 +121,11 @@ public class StreamsResetter { + "* This tool will not clean up the local state on the stream application instances (the persisted " + "stores used to cache aggregation results).\n" + "You need to call KafkaStreams#cleanUp() in your application or manually delete them from the " - + "directory specified by \"state.dir\" configuration (/tmp/kafka-streams/ by default).\n\n" + + "directory specified by \"state.dir\" configuration (/tmp/kafka-streams/ by default).\n" + + "* When long session timeout has been configured, active members could take longer to get expired on the " + + "broker thus blocking the reset job to complete. Use the \"--force\" option could remove those left-over " + + "members immediately. Make sure to stop all stream applications when this option is specified " + + "to avoid unexpected disruptions.\n\n" + "*** Important! You will get wrong output if you don't clean up the local stores after running the " + "reset tool!\n\n"; @@ -149,7 +155,7 @@ public int run(final String[] args, properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, options.valueOf(bootstrapServerOption)); adminClient = Admin.create(properties); - validateNoActiveConsumers(groupId, adminClient); + maybeDeleteActiveConsumers(groupId, adminClient); allTopics.clear(); allTopics.addAll(adminClient.listTopics().names().get(60, TimeUnit.SECONDS)); @@ -176,8 +182,8 @@ public int run(final String[] args, return exitCode; } - private void validateNoActiveConsumers(final String groupId, - final Admin adminClient) + private void maybeDeleteActiveConsumers(final String groupId, + final Admin adminClient) throws ExecutionException, InterruptedException { final DescribeConsumerGroupsResult describeResult = adminClient.describeConsumerGroups( @@ -186,9 +192,19 @@ private void validateNoActiveConsumers(final String groupId, final List members = new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); if (!members.isEmpty()) { - throw new IllegalStateException("Consumer group '" + groupId + "' is still active " - + "and has following members: " + members + ". " - + "Make sure to stop all running application instances before running the reset tool."); + if (options.has(forceOption)) { + System.out.println("Force deleting all active members in the group: " + groupId); + try { + adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions()).all().get(); + } catch (Exception e) { + throw e; + } + } else { + throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + + "and has following members: " + members + ". " + + "Make sure to stop all running application instances before running the reset tool." + + " You can use option '--force' to remove active members from the group."); + } } } @@ -236,6 +252,8 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); + forceOption = optionParser.accepts("force", "Force the removal of members of the consumer group (intended to remove stopped members if a long session timeout was used). " + + "Make sure to shut down all stream applications when this option is specified to avoid unexpected rebalances."); executeOption = optionParser.accepts("execute", "Execute the command."); dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); helpOption = optionParser.accepts("help", "Print usage information.").forHelp(); diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index f014481951290..90b6572616d3f 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1017,10 +1017,16 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertTrue(0 == list1.errors().get().size()) assertTrue(0 == list1.valid().get().size()) val testTopicName = "test_topic" + val testTopicName1 = testTopicName + "1" + val testTopicName2 = testTopicName + "2" val testNumPartitions = 2 - client.createTopics(Collections.singleton( - new NewTopic(testTopicName, testNumPartitions, 1.toShort))).all().get() - waitForTopics(client, List(testTopicName), List()) + + client.createTopics(util.Arrays.asList( + new NewTopic(testTopicName, testNumPartitions, 1.toShort), + new NewTopic(testTopicName1, testNumPartitions, 1.toShort), + new NewTopic(testTopicName2, testNumPartitions, 1.toShort) + )).all().get() + waitForTopics(client, List(testTopicName, testTopicName1, testTopicName2), List()) val producer = createProducer() try { @@ -1028,36 +1034,54 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { } finally { Utils.closeQuietly(producer, "producer") } + + val EMPTY_GROUP_INSTANCE_ID = "" val testGroupId = "test_group_id" val testClientId = "test_client_id" - val testInstanceId = "test_instance_id" + val testInstanceId1 = "test_instance_id_1" + val testInstanceId2 = "test_instance_id_2" val fakeGroupId = "fake_group_id" - val newConsumerConfig = new Properties(consumerConfig) - newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) - newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) - newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, testInstanceId) - val consumer = createConsumer(configOverrides = newConsumerConfig) - val latch = new CountDownLatch(1) + + def createProperties(groupInstanceId: String): Properties = { + val newConsumerConfig = new Properties(consumerConfig) + newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) + newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) + if (groupInstanceId != EMPTY_GROUP_INSTANCE_ID) { + newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, groupInstanceId) + } + newConsumerConfig + } + + // contains two static members and one dynamic member + val groupInstanceSet = Set(testInstanceId1, testInstanceId2, EMPTY_GROUP_INSTANCE_ID) + val consumerSet = groupInstanceSet.map { groupInstanceId => createConsumer(configOverrides = createProperties(groupInstanceId))} + val topicSet = Set(testTopicName, testTopicName1, testTopicName2) + + val latch = new CountDownLatch(consumerSet.size) try { - // Start a consumer in a thread that will subscribe to a new group. - val consumerThread = new Thread { - override def run : Unit = { - consumer.subscribe(Collections.singleton(testTopicName)) - - try { - while (true) { - consumer.poll(JDuration.ofSeconds(5)) - if (!consumer.assignment.isEmpty && latch.getCount > 0L) - latch.countDown() - consumer.commitSync() + def createConsumerThread[K,V](consumer: KafkaConsumer[K,V], topic: String): Thread = { + new Thread { + override def run : Unit = { + consumer.subscribe(Collections.singleton(topic)) + try { + while (true) { + consumer.poll(JDuration.ofSeconds(5)) + if (!consumer.assignment.isEmpty && latch.getCount > 0L) + latch.countDown() + consumer.commitSync() + } + } catch { + case _: InterruptException => // Suppress the output to stderr } - } catch { - case _: InterruptException => // Suppress the output to stderr } } } + + // Start consumers in a thread that will subscribe to a new group. + val consumerThreads = consumerSet.zip(topicSet).map(zipped => createConsumerThread(zipped._1, zipped._2)) + try { - consumerThread.start + consumerThreads.foreach(_.start()) assertTrue(latch.await(30000, TimeUnit.MILLISECONDS)) // Test that we can list the new group. TestUtils.waitUntilTrue(() => { @@ -1075,13 +1099,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId()) assertFalse(testGroupDescription.isSimpleConsumerGroup) - assertEquals(1, testGroupDescription.members().size()) - val member = testGroupDescription.members().iterator().next() - assertEquals(testClientId, member.clientId()) - val topicPartitions = member.assignment().topicPartitions() - assertEquals(testNumPartitions, topicPartitions.size()) - assertEquals(testNumPartitions, topicPartitions.asScala. - count(tp => tp.topic().equals(testTopicName))) + assertEquals(groupInstanceSet.size, testGroupDescription.members().size()) + val members = testGroupDescription.members() + members.asScala.foreach(member => assertEquals(testClientId, member.clientId())) + val topicPartitionsByTopic = members.asScala.flatMap(_.assignment().topicPartitions().asScala).groupBy(_.topic()) + topicSet.foreach { topic => + val topicPartitions = topicPartitionsByTopic.getOrElse(topic, List.empty) + assertEquals(testNumPartitions, topicPartitions.size) + } + val expectedOperations = AclEntry.supportedOperations(ResourceType.GROUP).asJava assertEquals(expectedOperations, testGroupDescription.authorizedOperations()) @@ -1129,16 +1155,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertFutureExceptionTypeEquals(deleteResult.deletedGroups().get(testGroupId), classOf[GroupNotEmptyException]) - // Test delete correct member + // Test delete one correct static member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(testInstanceId)) + Collections.singleton(new MemberToRemove(testInstanceId1)) )) assertNull(removeMembersResult.all().get()) - val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) + val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId1)) assertNull(validMemberFuture.get()) - // The group should contain no member now. val describeTestGroupResult = client.describeConsumerGroups(Seq(testGroupId).asJava, new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) assertEquals(1, describeTestGroupResult.describedGroups().size()) @@ -1147,6 +1172,16 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId) assertFalse(testGroupDescription.isSimpleConsumerGroup) + assertEquals(consumerSet.size - 1, testGroupDescription.members().size()) + + // Delete all active members remaining (a static member + a dynamic member) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions()) + assertNull(removeMembersResult.all().get()) + + // The group should contain no members now. + testGroupDescription = client.describeConsumerGroups(Seq(testGroupId).asJava, + new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) + .describedGroups().get(testGroupId).get() assertTrue(testGroupDescription.members().isEmpty) // Consumer group deletion on empty group should succeed @@ -1155,12 +1190,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertTrue(deleteResult.deletedGroups().containsKey(testGroupId)) assertNull(deleteResult.deletedGroups().get(testGroupId).get()) - } finally { - consumerThread.interrupt() - consumerThread.join() + } finally { + consumerThreads.foreach { + case consumerThread => + consumerThread.interrupt() + consumerThread.join() } + } } finally { - Utils.closeQuietly(consumer, "consumer") + consumerSet.zip(groupInstanceSet).foreach(zipped => Utils.closeQuietly(zipped._1, zipped._2)) } } finally { Utils.closeQuietly(client, "adminClient") diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java index 32a75f4b0f28e..9652db206709f 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java @@ -62,6 +62,7 @@ import java.util.Properties; import static java.time.Duration.ofMillis; +import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.isEmptyConsumerGroup; import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.waitForEmptyConsumerGroup; import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; @@ -261,6 +262,41 @@ public void shouldNotAllowToResetWhenIntermediateTopicAbsent() throws Exception Assert.assertEquals(1, exitCode); } + public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Exception { + appID = testId + "-with-force-option"; + streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, appID); + streamsConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + STREAMS_CONSUMER_TIMEOUT * 100); + + // Run + streams = new KafkaStreams(setupTopologyWithoutIntermediateUserTopic(), streamsConfig); + streams.start(); + final List> result = IntegrationTestUtils.waitUntilMinKeyValueRecordsReceived(resultConsumerConfig, OUTPUT_TOPIC, 10); + + streams.close(); + + // RESET + streams = new KafkaStreams(setupTopologyWithoutIntermediateUserTopic(), streamsConfig); + streams.cleanUp(); + + // Reset would fail since long session timeout has been configured + final boolean cleanResult = tryCleanGlobal(false, null, null); + Assert.assertEquals(false, cleanResult); + + // Reset will success with --force, it will force delete active members on broker side + cleanGlobal(false, "--force", null); + assertThat("Group is not empty after cleanGlobal", isEmptyConsumerGroup(adminClient, appID)); + + assertInternalTopicsGotDeleted(null); + + // RE-RUN + streams.start(); + final List> resultRerun = IntegrationTestUtils.waitUntilMinKeyValueRecordsReceived(resultConsumerConfig, OUTPUT_TOPIC, 10); + streams.close(); + + assertThat(resultRerun, equalTo(result)); + cleanGlobal(false, "--force", null); + } + void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() throws Exception { appID = testId + "-from-scratch"; streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, appID); @@ -507,9 +543,9 @@ private Topology setupTopologyWithoutIntermediateUserTopic() { return builder.build(); } - private void cleanGlobal(final boolean withIntermediateTopics, - final String resetScenario, - final String resetScenarioArg) throws Exception { + private boolean tryCleanGlobal(final boolean withIntermediateTopics, + final String resetScenario, + final String resetScenarioArg) throws Exception { // leaving --zookeeper arg here to ensure tool works if users add it final List parameterList = new ArrayList<>( Arrays.asList("--application-id", appID, @@ -546,8 +582,14 @@ private void cleanGlobal(final boolean withIntermediateTopics, cleanUpConfig.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 100); cleanUpConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + CLEANUP_CONSUMER_TIMEOUT); - final int exitCode = new StreamsResetter().run(parameters, cleanUpConfig); - Assert.assertEquals(0, exitCode); + return new StreamsResetter().run(parameters, cleanUpConfig) == 0; + } + + private void cleanGlobal(final boolean withIntermediateTopics, + final String resetScenario, + final String resetScenarioArg) throws Exception { + final boolean cleanResult = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); + Assert.assertEquals(true, cleanResult); } private void assertInternalTopicsGotDeleted(final String intermediateUserTopic) throws Exception { diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java index ed04710561469..71fcc80fff4a3 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java @@ -72,6 +72,11 @@ public void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() super.testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic(); } + @Test + public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Exception { + super.testResetWhenLongSessionTimeoutConfiguredWithForceOption(); + } + @Test public void testReprocessingFromScratchAfterResetWithIntermediateUserTopic() throws Exception { super.testReprocessingFromScratchAfterResetWithIntermediateUserTopic(); diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java b/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java index eb7fa4c1cb38b..f12aa0d8259d5 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java @@ -858,16 +858,7 @@ private ConsumerGroupInactiveCondition(final Admin adminClient, @Override public boolean conditionMet() { - try { - final ConsumerGroupDescription groupDescription = - adminClient.describeConsumerGroups(Collections.singletonList(applicationId)) - .describedGroups() - .get(applicationId) - .get(); - return groupDescription.members().isEmpty(); - } catch (final ExecutionException | InterruptedException e) { - return false; - } + return isEmptyConsumerGroup(adminClient, applicationId); } } @@ -878,6 +869,20 @@ public static void waitForEmptyConsumerGroup(final Admin adminClient, "Test consumer group " + applicationId + " still active even after waiting " + timeoutMs + " ms."); } + public static boolean isEmptyConsumerGroup(final Admin adminClient, + final String applicationId) { + try { + final ConsumerGroupDescription groupDescription = + adminClient.describeConsumerGroups(Collections.singletonList(applicationId)) + .describedGroups() + .get(applicationId) + .get(); + return groupDescription.members().isEmpty(); + } catch (final ExecutionException | InterruptedException e) { + return false; + } + } + private static StateListener getStateListener(final KafkaStreams streams) { try { final Field field = streams.getClass().getDeclaredField("stateListener");